Стриминг с stream_events, асинхронный стриминг, обнаружение прерываний в потоке.
Стриминг даёт промежуточные результаты по ходу выполнения агента, вместо того чтобы ждать конца всего цикла.
invoke возвращает результат только после завершения. Если агент сходил в три инструмента и дважды подумал, пользователь двадцать секунд смотрит на пустой экран. Хуже того, он не понимает, работает система или зависла.
Стриминг решает не только вопрос скорости отклика. Он даёт видимость процесса: какой инструмент вызван, что он вернул, о чём модель рассуждает сейчас. Для отладки агентов это ценнее, чем финальный ответ.
Основной способ для новых приложений — это stream_events с version="v3". Возвращается объект потока, у которого несколько независимых проекций одних и тех же событий.
stream.messages отдаёт сообщения модели по мере генерации, stream.values показывает состояние после каждого шага, stream.output содержит финальное состояние, stream.interrupted и stream.interrupts сообщают о приостановке графа.
Самый частый сценарий — это печать ответа по токенам:
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)Внешний цикл идёт по сообщениям (модель за запуск может ответить несколько раз), внутренний по кускам текста внутри сообщения. У рассуждающих моделей рядом с text доступны блоки рассуждений, поэтому ход мысли можно показывать отдельно от ответа.
Когда нужно показывать и текст, и вызовы инструментов в правильном порядке, проекции объединяют через interleave. Метод отдаёт пары «вид события и само событие»:
stream = agent.stream_events(
{"messages": [{"role": "user", "content": "Какая погода в Казани?"}]},
config=config,
version="v3",
)
for kind, item in stream.interleave("messages", "tool_calls"):
if kind == "messages":
for token in item.text:
print(token, end="", flush=True)
else:
print(f"\n[вызов инструмента: {item}]")Без interleave пришлось бы читать проекции по очереди, теряя их взаимный порядок. Именно так строят интерфейсы, где под ответом появляются карточки «поиск выполнен», «страница загружена».
Проекция values даёт полный снимок состояния после каждого супершага. Это способ увидеть, как граф движется, не влезая внутрь узлов:
stream = agent.stream_events(user_input, version="v3")
for snapshot in stream.values:
msg = snapshot["messages"][-1]
if getattr(msg, "tool_calls", None):
print(f"Вызов инструмента: {[tc['name'] for tc in msg.tool_calls]}")
elif msg.text:
print(f"Ответ: {msg.text}")Финальное состояние доступно как stream.output после того, как поток исчерпан. Это избавляет от повторного invoke ради результата.
У графов LangGraph есть более низкоуровневый метод stream, где режим выбирается явно. Начиная с формата v2 каждый элемент потока — это словарь с полями type, ns и data, поэтому разбор одинаков для всех режимов.
Режим values отдаёт полное состояние после каждого шага, updates только изменения от конкретного узла, messages токены модели вместе с метаданными, custom произвольные данные, отправленные из узла, debug подробную служебную информацию.
for chunk in graph.stream(
{"topic": "мороженое"},
stream_mode="updates",
version="v2",
):
if chunk["type"] == "updates":
for node_name, state in chunk["data"].items():
print(f"Узел {node_name} обновил состояние: {state}")Режимы можно запрашивать списком: stream_mode=["updates", "custom"]. Тогда в потоке смешаны события разных типов, и поле type подсказывает, как читать data.
Разница между values и updates практическая: values удобен, когда интерфейс перерисовывает всё состояние целиком, updates дешевле и показывает именно то, что изменилось.
Долгий узел может сообщать о прогрессе, не дожидаясь конца работы. Для этого берут писатель потока:
from langgraph.config import get_stream_writer
def fetch_pages(state: State):
writer = get_stream_writer()
for i, url in enumerate(state["urls"], start=1):
writer({"progress": f"Загружаю {i} из {len(state['urls'])}"})
download(url)
return {"status": "done"}Эти события приходят в режиме custom. Ограничение стоит помнить заранее: функция с get_stream_writer внутри перестаёт вызываться вне контекста LangGraph, поэтому обычный юнит-тест на такой инструмент придётся запускать через граф либо выносить логику в отдельную функцию.
По умолчанию события подграфов не попадают в поток родителя. Флаг subgraphs=True включает их, а поле ns показывает, из какого именно подграфа пришло событие:
for chunk in graph.stream(
{"foo": "foo"},
subgraphs=True,
stream_mode="updates",
version="v2",
):
print(chunk["ns"]) # путь к подграфу
print(chunk["data"]) # его обновленияВ мультиагентных системах без этого флага прогресс субагентов выглядит как пауза: родительский граф ничего не отдаёт, пока субагент работает.
В асинхронном приложении используют astream_events и astream. Логика та же, меняется только форма цикла:
import asyncio
async def stream_agent(user_input):
stream = agent.astream_events(user_input, version="v3")
async for message in stream.messages:
async for token in message.text:
print(token, end="", flush=True)
asyncio.run(stream_agent({"messages": [{"role": "user", "content": "Привет"}]}))Для веб-приложения это основной вариант: асинхронный поток напрямую ложится на ответ вида server-sent events или на веб-сокет.
Когда граф останавливается на interrupt, поток просто заканчивается, а признак остановки остаётся на объекте потока. Рабочий цикл выглядит так: стримим, проверяем interrupted, спрашиваем человека, возобновляем.
from langgraph.types import Command
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(f"\n{stream.interrupts[0].value}\n> ")
stream_input = Command(resume=user_response)Обратите внимание на два момента. Проверять interrupted нужно после того, как поток дочитан до конца: до этого признак ещё не установлен. И для возобновления на вход подаётся Command(resume=...) вместо обычного словаря с сообщениями.
Печать message.content вместо message.text. У части провайдеров содержимое приходит списком блоков, и на экране появляются словари.
Проверка stream.interrupted до окончания чтения потока. Признак появляется, когда события закончились.
Смешивание синхронного и асинхронного вызовов. stream_events внутри async def заблокирует цикл событий, для асинхронного кода нужен astream_events.
Забытый subgraphs=True при отладке мультиагентной системы. Кажется, что субагент завис, хотя он просто работает молча.
Далее: Память и сохранение состояния