Стриминг с stream_events, асинхронный стриминг, обнаружение прерываний в потоке.
Стриминг позволяет получать промежуточные результаты выполнения агента в реальном времени, без ожидания завершения всего цикла.
Метод invoke возвращает результат только после завершения. Если агент вызывает несколько инструментов, пользователь не видит промежуточных шагов. Стриминг решает эту проблему, выдавая данные по мере поступления.
Основной способ стриминга — stream_events с version="v3". Возвращает объект с несколькими проекциями:
stream.messages — сообщения модели по токенамstream.values — полное состояние после каждого шагаstream.interrupted — прерван ли поток прерываниемstream.interrupts — данные прерыванийstream.output — финальное состояниеstream = agent.stream_events(
{"messages": [{"role": "user", "content": "Найди новости об ИИ"}]},
version="v3",
)
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)Через stream.values можно видеть состояние после каждого шага — какие инструменты вызваны, какие ответы получены:
stream = agent.stream_events(input, version="v3")
for snapshot in stream.values:
msg = snapshot["messages"][-1]
if msg.tool_calls:
print(f"Вызов инструмента: {[tc['name'] for tc in msg.tool_calls]}")
elif msg.content:
print(f"Ответ: {msg.content}")Для асинхронного использования — astream_events:
import asyncio
async def stream_agent():
stream = agent.astream_events(input, version="v3")
async for message in stream.messages:
async for token in message.text:
print(token, end="", flush=True)
asyncio.run(stream_agent())При использовании прерываний стриминг позволяет обнаружить паузу и получить данные прерывания. Проверяйте stream.interrupted после завершения — если True, граф ждёт внешнего ввода.
Для интерактивных агентов с прерываниями используйте цикл со стримингом:
stream_input = initial_input
while True:
stream = agent.stream_events(stream_input, config=config, version="v3")
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)
if not stream.interrupted:
break
user_response = input(stream.interrupts[0].value)
stream_input = Command(resume=user_response)Цикл продолжается, пока граф не завершится без прерываний. После каждого прерывания запрашивается ввод пользователя и граф возобновляется.
Далее: Память и сохранение состояния