스트리밍

스트리밍 (Streaming)

이 가이드에서는 DSPy 프로그램에서 스트리밍을 활성화하는 방법을 안내할게요. DSPy 스트리밍은 두 부분으로 구성돼요:

  • 출력 토큰 스트리밍 (Output Token Streaming): 완전한 응답을 기다리는 대신, 토큰이 생성될 때 개별적으로 스트리밍해요.
  • 중간 상태 스트리밍 (Intermediate Status Streaming): 프로그램의 실행 상태("웹 검색 호출 중...", "결과 처리 중...")에 대한 실시간 업데이트를 제공해요.

출처: 문서

본문

출력 토큰 스트리밍 (Output Token Streaming)

DSPy의 토큰 스트리밍 기능은 최종 출력뿐 아니라 파이프라인의 모든 모듈에서 동작해요. 유일한 요구사항은 스트리밍할 필드가 str 타입이어야 한다는 것뿐이에요. 토큰 스트리밍을 활성화하려면:

  1. 프로그램을 dspy.streamify로 감싸기
  2. 스트리밍할 필드를 지정하는 dspy.streaming.StreamListener 객체를 하나 이상 만들기

기본 예제를 살펴볼게요:

import os

import dspy

os.environ["OPENAI_API_KEY"] = "your_api_key"

dspy.configure(lm=dspy.LM("openai/gpt-4o-mini"))

predict = dspy.Predict("question->answer")

# Enable streaming for the 'answer' field
stream_predict = dspy.streamify(
    predict,
    stream_listeners=[dspy.streaming.StreamListener(signature_field_name="answer")],
)

스트리밍된 출력을 소비하려면:

import asyncio

async def read_output_stream():
    output_stream = stream_predict(question="Why did a chicken cross the kitchen?")

    async for chunk in output_stream:
        print(chunk)

asyncio.run(read_output_stream())

다음과 같은 출력이 생성돼요:

StreamResponse(predict_name='self', signature_field_name='answer', chunk='To')
StreamResponse(predict_name='self', signature_field_name='answer', chunk=' get')
StreamResponse(predict_name='self', signature_field_name='answer', chunk=' to')
StreamResponse(predict_name='self', signature_field_name='answer', chunk=' the')
StreamResponse(predict_name='self', signature_field_name='answer', chunk=' other')
StreamResponse(predict_name='self', signature_field_name='answer', chunk=' side of the frying pan!')
Prediction(
    answer='To get to the other side of the frying pan!'
)

참고: dspy.streamify는 async 제너레이터를 반환하므로 async 컨텍스트 안에서 사용해야 해요. Jupyter나 Google Colab처럼 이미 이벤트 루프(async 컨텍스트)가 있는 환경이라면 제너레이터를 직접 사용할 수 있어요.

위 스트리밍에 StreamResponse와 Prediction 두 가지 다른 엔티티가 포함되어 있다는 걸 눈치챘을 거예요. StreamResponse는 청취 중인 필드의 스트리밍 토큰을 감싸는 래퍼이고, 이 예제에서는 answer 필드예요. Prediction은 프로그램의 최종 출력입니다. DSPy에서 스트리밍은 sidecar 방식으로 구현돼요: LM에서 스트리밍을 활성화해 LM이 토큰의 스트림을 출력하게 해요. 이 토큰들을 측면 채널(side channel)로 보내고, 사용자가 정의한 리스너가 계속해서 읽습니다. 리스너는 스트림을 계속 해석하면서, 청취 중인 signature_field_name이 나타나기 시작했는지, 그리고 완결되었는지를 결정해요. 필드가 나타났다고 결정하면 리스너는 사용자가 읽을 수 있는 async 제너레이터에 토큰을 출력하기 시작합니다. 리스너의 내부 메커니즘은 뒤에서 사용하는 어댑터에 따라 바뀌는데, 보통 다음 필드를 볼 때까지 필드가 완결되었는지 결정할 수 없기 때문에 리스너는 최종 제너레이터로 보내기 전에 출력 토큰을 버퍼링해요. 그래서 마지막 StreamResponse 타입의 chunk에는 보통 토큰이 두 개 이상 들어 있답니다. 프로그램의 출력도 스트림에 기록되는데, 위 샘플 출력의 Prediction chunk가 바로 그거예요.

이 다양한 타입을 처리하고 커스텀 로직을 구현하려면:

import asyncio

async def read_output_stream():
  output_stream = stream_predict(question="Why did a chicken cross the kitchen?")

  return_value = None
  async for chunk in output_stream:
    if isinstance(chunk, dspy.streaming.StreamResponse):
      print(f"Output token of field {chunk.signature_field_name}: {chunk.chunk}")
    elif isinstance(chunk, dspy.Prediction):
      return_value = chunk
  return return_value


program_output = asyncio.run(read_output_stream())
print("Final output: ", program_output)

StreamResponse 이해하기 (Understand StreamResponse)

StreamResponse(dspy.streaming.StreamResponse)는 스트리밍 토큰의 래퍼 클래스예요. 3개의 필드가 있습니다:

  • predict_name: signature_field_name을 보유한 predict의 이름이에요. your_program.named_predictors()를 실행할 때의 키 이름과 같아요. 위 코드에서 answer는 predict 자체에서 왔기 때문에 predict_name이 self로 나타나는데, 이는 predict.named_predictors()를 실행했을 때의 유일한 키입니다.
  • signature_field_name: 이 토큰들이 매핑되는 출력 필드예요. predict_name과 signature_field_name이 함께 필드의 고유 식별자를 형성합니다. 여러 필드 스트리밍과 중복 필드 이름을 처리하는 방법은 이 가이드 뒷부분에서 보여줄게요.
  • chunk: 스트림 chunk의 값이에요.

캐시와 함께하는 스트리밍 (Streaming with Cache)

캐시된 결과를 찾으면 스트림은 개별 토큰을 건너뛰고 최종 Prediction만 생성해요. 예를 들어:

Prediction(
    answer='To get to the other side of the dinner plate!'
)

여러 필드 스트리밍 (Streaming Multiple Fields)

각 필드에 대해 StreamListener를 만들어 여러 필드를 모니터링할 수 있어요. 다중 모듈 프로그램 예제를 살펴볼게요:

import asyncio

import dspy

lm = dspy.LM("openai/gpt-4o-mini", cache=False)
dspy.configure(lm=lm)


class MyModule(dspy.Module):
    def __init__(self):
        super().__init__()

        self.predict1 = dspy.Predict("question->answer")
        self.predict2 = dspy.Predict("answer->simplified_answer")

    def forward(self, question: str, **kwargs):
        answer = self.predict1(question=question)
        simplified_answer = self.predict2(answer=answer)
        return simplified_answer


predict = MyModule()
stream_listeners = [
    dspy.streaming.StreamListener(signature_field_name="answer"),
    dspy.streaming.StreamListener(signature_field_name="simplified_answer"),
]
stream_predict = dspy.streamify(
    predict,
    stream_listeners=stream_listeners,
)

async def read_output_stream():
    output = stream_predict(question="why did a chicken cross the kitchen?")

    return_value = None
    async for chunk in output:
        if isinstance(chunk, dspy.streaming.StreamResponse):
            print(chunk)
        elif isinstance(chunk, dspy.Prediction):
            return_value = chunk
    return return_value

program_output = asyncio.run(read_output_stream())
print("Final output: ", program_output)

출력은 다음과 같을 거예요:

StreamResponse(predict_name='predict1', signature_field_name='answer', chunk='To')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' get')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' to')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' the')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' other side of the recipe!')
StreamResponse(predict_name='predict2', signature_field_name='simplified_answer', chunk='To')
StreamResponse(predict_name='predict2', signature_field_name='simplified_answer', chunk=' reach')
StreamResponse(predict_name='predict2', signature_field_name='simplified_answer', chunk=' the')
StreamResponse(predict_name='predict2', signature_field_name='simplified_answer', chunk=' other side of the recipe!')
Final output:  Prediction(
    simplified_answer='To reach the other side of the recipe!'
)

같은 필드를 여러 번 스트리밍하기 (dspy.ReAct처럼) (Streaming the Same Field Multiple Times)

기본적으로 StreamListener는 단일 스트리밍 세션을 완료하면 자동으로 스스로 닫아요. 이 설계는 성능 문제를 방지하는데, 모든 토큰이 구성된 모든 스트림 리스너에 브로드캐스트되므로 활성 리스너가 너무 많으면 상당한 오버헤드가 발생할 수 있기 때문이에요.

하지만 dspy.ReAct처럼 DSPy 모듈이 루프에서 반복적으로 사용되는 시나리오에서는, 매번 사용될 때마다 각 예측에서 같은 필드를 스트리밍하고 싶을 수 있어요. 이 동작을 활성화하려면 StreamListener를 만들 때 allow_reuse=True를 설정하세요. 아래 예제를 봐요:

import asyncio

import dspy

lm = dspy.LM("openai/gpt-4o-mini", cache=False)
dspy.configure(lm=lm)


def fetch_user_info(user_name: str):
    """Get user information like name, birthday, etc."""
    return {
        "name": user_name,
        "birthday": "2009-05-16",
    }


def get_sports_news(year: int):
    """Get sports news for a given year."""
    if year == 2009:
        return "Usane Bolt broke the world record in the 100m race."
    return None


react = dspy.ReAct("question->answer", tools=[fetch_user_info, get_sports_news])

stream_listeners = [
    # dspy.ReAct has a built-in output field called "next_thought".
    dspy.streaming.StreamListener(signature_field_name="next_thought", allow_reuse=True),
]
stream_react = dspy.streamify(react, stream_listeners=stream_listeners)


async def read_output_stream():
    output = stream_react(question="What sports news happened in the year Adam was born?")
    return_value = None
    async for chunk in output:
        if isinstance(chunk, dspy.streaming.StreamResponse):
            print(chunk)
        elif isinstance(chunk, dspy.Prediction):
            return_value = chunk
    return return_value


print(asyncio.run(read_output_stream()))

이 예제에서 StreamListener에 allow_reuse=True를 설정하면 "next_thought"에 대한 스트리밍이 첫 번째 반복뿐 아니라 모든 반복에서 사용 가능해져요. 이 코드를 실행하면 필드가 생성될 때마다 next_thought의 스트리밍 토큰이 출력되는 것을 볼 수 있어요.

중복 필드 이름 처리 (Handling Duplicate Field Names)

서로 다른 모듈에서 같은 이름의 필드를 스트리밍할 때는 StreamListener에서 predict와 predict_name을 모두 지정하세요:

import asyncio

import dspy

lm = dspy.LM("openai/gpt-4o-mini", cache=False)
dspy.configure(lm=lm)


class MyModule(dspy.Module):
    def __init__(self):
        super().__init__()

        self.predict1 = dspy.Predict("question->answer")
        self.predict2 = dspy.Predict("question, answer->answer, score")

    def forward(self, question: str, **kwargs):
        answer = self.predict1(question=question)
        simplified_answer = self.predict2(answer=answer)
        return simplified_answer


predict = MyModule()
stream_listeners = [
    dspy.streaming.StreamListener(
        signature_field_name="answer",
        predict=predict.predict1,
        predict_name="predict1"
    ),
    dspy.streaming.StreamListener(
        signature_field_name="answer",
        predict=predict.predict2,
        predict_name="predict2"
    ),
]
stream_predict = dspy.streamify(
    predict,
    stream_listeners=stream_listeners,
)


async def read_output_stream():
    output = stream_predict(question="why did a chicken cross the kitchen?")

    return_value = None
    async for chunk in output:
        if isinstance(chunk, dspy.streaming.StreamResponse):
            print(chunk)
        elif isinstance(chunk, dspy.Prediction):
            return_value = chunk
    return return_value


program_output = asyncio.run(read_output_stream())
print("Final output: ", program_output)

출력은 다음과 같을 거예요:

StreamResponse(predict_name='predict1', signature_field_name='answer', chunk='To')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' get')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' to')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' the')
StreamResponse(predict_name='predict1', signature_field_name='answer', chunk=' other side of the recipe!')
StreamResponse(predict_name='predict2', signature_field_name='answer', chunk="I'm")
StreamResponse(predict_name='predict2', signature_field_name='answer', chunk=' ready')
StreamResponse(predict_name='predict2', signature_field_name='answer', chunk=' to')
StreamResponse(predict_name='predict2', signature_field_name='answer', chunk=' assist')
StreamResponse(predict_name='predict2', signature_field_name='answer', chunk=' you')
StreamResponse(predict_name='predict2', signature_field_name='answer', chunk='! Please provide a question.')
Final output:  Prediction(
    answer="I'm ready to assist you! Please provide a question.",
    score='N/A'
)

중간 상태 스트리밍 (Intermediate Status Streaming)

상태 스트리밍은 특히 도구 호출이나 복잡한 AI 파이프라인 같은 오래 실행되는 작업에서 사용자에게 프로그램의 진행 상황을 알려줘요. 상태 스트리밍을 구현하려면:

  1. dspy.streaming.StatusMessageProvider를 서브클래싱해 커스텀 상태 메시지 제공자를 만들기
  2. 원하는 훅 메서드를 오버라이드해 커스텀 상태 메시지 제공하기
  3. 제공자를 dspy.streamify에 전달하기

예시:

class MyStatusMessageProvider(dspy.streaming.StatusMessageProvider):
    def lm_start_status_message(self, instance, inputs):
        return f"Calling LM with inputs {inputs}..."

    def lm_end_status_message(self, outputs):
        return f"Tool finished with output: {outputs}!"

사용 가능한 훅:

  • lm_start_status_message: dspy.LM 호출 시작 시 상태 메시지.
  • lm_end_status_message: dspy.LM 호출 종료 시 상태 메시지.
  • module_start_status_message: dspy.Module 호출 시작 시 상태 메시지.
  • module_end_status_message: dspy.Module 호출 종료 시 상태 메시지.
  • tool_start_status_message: dspy.Tool 호출 시작 시 상태 메시지.
  • tool_end_status_message: dspy.Tool 호출 종료 시 상태 메시지.

각 훅은 상태 메시지를 포함한 문자열을 반환해야 해요.

메시지 제공자를 만든 뒤 dspy.streamify에 전달하면 상태 메시지 스트리밍과 출력 토큰 스트리밍을 모두 활성화할 수 있어요. 아래 예제를 봐요. 중간 상태 메시지는 dspy.streaming.StatusMessage 클래스로 표현되므로, 이를 캡처하려면 또 다른 조건 검사가 필요합니다.

import asyncio

import dspy

lm = dspy.LM("openai/gpt-4o-mini", cache=False)
dspy.configure(lm=lm)


class MyModule(dspy.Module):
    def __init__(self):
        super().__init__()

        self.tool = dspy.Tool(lambda x: 2 * x, name="double_the_number")
        self.predict = dspy.ChainOfThought("num1, num2->sum")

    def forward(self, num, **kwargs):
        num2 = self.tool(x=num)
        return self.predict(num1=num, num2=num2)


class MyStatusMessageProvider(dspy.streaming.StatusMessageProvider):
    def tool_start_status_message(self, instance, inputs):
        return f"Calling Tool {instance.name} with inputs {inputs}..."

    def tool_end_status_message(self, outputs):
        return f"Tool finished with output: {outputs}!"


predict = MyModule()
stream_listeners = [
    # dspy.ChainOfThought has a built-in output field called "reasoning".
    dspy.streaming.StreamListener(signature_field_name="reasoning"),
]
stream_predict = dspy.streamify(
    predict,
    stream_listeners=stream_listeners,
    status_message_provider=MyStatusMessageProvider(),
)


async def read_output_stream():
    output = stream_predict(num=3)

    return_value = None
    async for chunk in output:
        if isinstance(chunk, dspy.streaming.StreamResponse):
            print(chunk)
        elif isinstance(chunk, dspy.Prediction):
            return_value = chunk
        elif isinstance(chunk, dspy.streaming.StatusMessage):
            print(chunk)
    return return_value


program_output = asyncio.run(read_output_stream())
print("Final output: ", program_output)

샘플 출력:

StatusMessage(message='Calling tool double_the_number...')
StatusMessage(message='Tool calling finished! Querying the LLM with tool calling results...')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk='To')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' find')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' the')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' sum')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' of')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' the')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' two')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' numbers')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=',')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' we')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' simply')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' add')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' them')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' together')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk='.')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' Here')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=',')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' ')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk='3')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' plus')
StreamResponse(predict_name='predict.predict', signature_field_name='reasoning', chunk=' 6 equals 9.')
Final output:  Prediction(
    reasoning='To find the sum of the two numbers, we simply add them together. Here, 3 plus 6 equals 9.',
    sum='9'
)

동기 스트리밍 (Synchronous Streaming)

기본적으로 streamify된 DSPy 프로그램을 호출하면 async 제너레이터가 생성돼요. sync 제너레이터를 얻으려면 async_streaming=False 플래그를 설정하면 됩니다:

import os

import dspy

os.environ["OPENAI_API_KEY"] = "your_api_key"

dspy.configure(lm=dspy.LM("openai/gpt-4o-mini"))

predict = dspy.Predict("question->answer")

# Enable streaming for the 'answer' field
stream_predict = dspy.streamify(
    predict,
    stream_listeners=[dspy.streaming.StreamListener(signature_field_name="answer")],
    async_streaming=False,
)

output = stream_predict(question="why did a chicken cross the kitchen?")

program_output = None
for chunk in output:
    if isinstance(chunk, dspy.streaming.StreamResponse):
        print(chunk)
    elif isinstance(chunk, dspy.Prediction):
        program_output = chunk
print(f"Program output: {program_output}")

더 알아보기 (Learn more)