Перейти к основному контенту
Tech Path Finder
КурсыИнтервьюКод-ревьюБлог
Tech Path Finder

Персонализированный путеводитель в IT. Квизы, мок-интервью, код ревью и аналитика прогресса.

@potapov_me

Платформа

  • Курсы
  • Прогресс
  • Мок-интервью
  • Код ревью
  • Живое ревью с ИИ
  • Тренажёр переговоров
  • Закладки

Контент

  • Блог
  • Главная
  • Обратная связь

Компания

  • О проекте
  • Тарифы
  • Условия использования
  • Конфиденциальность
  • Согласие на обработку данных
  • Cookie
  • Реквизиты

Аккаунт

  • Войти
  • Зарегистрироваться
  • Профиль

© 2026 Tech Path Finder. Все права защищены.

·ИП Потапов К.С.·Политика конфиденциальности·
Сделано с ❤️ в России
  1. Потоковая передача данных
streaming

Потоковая передача данных

Стриминг с stream_events, асинхронный стриминг, обнаружение прерываний в потоке.

Потоковая передача данных

Стриминг даёт промежуточные результаты по ходу выполнения агента, вместо того чтобы ждать конца всего цикла.

#Почему без стриминга интерфейс кажется сломанным

invoke возвращает результат только после завершения. Если агент сходил в три инструмента и дважды подумал, пользователь двадцать секунд смотрит на пустой экран. Хуже того, он не понимает, работает система или зависла.

Стриминг решает не только вопрос скорости отклика. Он даёт видимость процесса: какой инструмент вызван, что он вернул, о чём модель рассуждает сейчас. Для отладки агентов это ценнее, чем финальный ответ.

#Событийный стриминг: stream_events

Основной способ для новых приложений — это 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 при отладке мультиагентной системы. Кажется, что субагент завис, хотя он просто работает молча.

Устранение неисправностей

Далее: Память и сохранение состояния