스트리밍
스트리밍 (Streaming)
이 가이드에서는 DSPy 프로그램에서 스트리밍을 활성화하는 방법을 안내할게요. DSPy 스트리밍은 두 부분으로 구성돼요:
- 출력 토큰 스트리밍 (Output Token Streaming): 완전한 응답을 기다리는 대신, 토큰이 생성될 때 개별적으로 스트리밍해요.
- 중간 상태 스트리밍 (Intermediate Status Streaming): 프로그램의 실행 상태("웹 검색 호출 중...", "결과 처리 중...")에 대한 실시간 업데이트를 제공해요.
출처: 문서
본문
출력 토큰 스트리밍 (Output Token Streaming)
DSPy의 토큰 스트리밍 기능은 최종 출력뿐 아니라 파이프라인의 모든 모듈에서 동작해요. 유일한 요구사항은 스트리밍할 필드가 str 타입이어야 한다는 것뿐이에요. 토큰 스트리밍을 활성화하려면:
- 프로그램을
dspy.streamify로 감싸기 - 스트리밍할 필드를 지정하는
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 파이프라인 같은 오래 실행되는 작업에서 사용자에게 프로그램의 진행 상황을 알려줘요. 상태 스트리밍을 구현하려면:
dspy.streaming.StatusMessageProvider를 서브클래싱해 커스텀 상태 메시지 제공자를 만들기- 원하는 훅 메서드를 오버라이드해 커스텀 상태 메시지 제공하기
- 제공자를
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}")