Поток событий: stream_events версии v3 | Курс LangChain урок 8
Цель урока: читать прогон агента через stream_events(version="v3") и раскладывать его на проекции вместо ручного разбора chunks режима. Находить нужное в сыром потоке событий. Писать свой преобразователь потока, регистрируя его на вызове и в middleware.
Необходимые знания:
1) урок 0: окружение собрано, ключ работает, переменные MODEL_NAME и MODEL_BASE_URL заполнены
2) урок 2: init_chat_model, первый взгляд на create_agent
3) урок 4: middleware как место, куда встраивается своя логика
4) урок 7: режимы потока updates, messages и custom, метод stream и то, что показывать пользователю
5) Python на уровне джуниора: генераторы, TypedDict, наследование, Counter
Ключевые концепции:
1) поток событий, это слой над режимами stream
2) stream_events(version="v3") возвращает объект прогона с проекциями
3) проекции: messages, values, tool_calls, output, extensions
4) сырое событие протокола и его поля: seq, method, namespace, data
5) имя узла модели сменилось с agent на model, и старый отбор молча даёт ноль
6) проекция копит значения только с момента подписки, отсюда interleave
7) свой преобразователь потока: init, process, finalize, fail, required_stream_modes
8) регистрация преобразователя на вызове и в middleware
Зачем понадобилась вторая механика потока
В уроке 7 разобраны метод stream и его режимы. Там вы выбираете режим заранее и сами читаете chunks: этот от узла модели, в нём токен, а этот от узла инструментов, в нём правка состояния. Работает, но вся раскладка на вашей стороне.
Поток событий решает ту же задачу с другого конца. Вы получаете один объект прогона, а из него берёте отдельные проекции: сообщения отдельно, состояние отдельно, вызовы инструментов отдельно. Проекция это часть того же потока, в которой приходят объекты только одного вида, и вид этот известен заранее. Поэтому поля читаются сразу, без проверки, что пришло: message.text у сообщения, call.tool_name у вызова инструмента. Каждая проекция получает события независимо от других, но копит их только с момента подписки, об этом раздел "Две проекции сразу".
Для новых приложений авторы фреймворка рекомендуют поток событий. Он появился в LangChain 1.3.
Устроено это слоями.
1) нижний слой: движок Pregel выдаёт сырые события выполнения графа. Это и есть режимы из урока 7: updates, values, messages, custom, checkpoints, tasks, debug
2) верхний слой: поток событий приводит их к одному виду, прогоняет через преобразователи потока и раскладывает по проекциям
Между слоями стоит маршрутизатор событий: он получает нормализованное событие и передаёт его по очереди каждому зарегистрированному преобразователю. Встроенные дают stream.values, stream.messages, stream.lifecycle и stream.subgraphs, агент добавляет stream.tool_calls и stream.subagents. Итог stream.output собирает сам объект прогона, из последнего снимка состояния. Свои добавляют проекции в stream.extensions, и во второй половине урока вы такой напишете.
Хотя авторы рекомендуют поток событий, версия v3 пока экспериментальная. В коде пакета, в описании метода stream_events, сказано, что интерфейс version="v3" может измениться, а класс объекта прогона GraphRunStream помечен декоратором @beta. На страницах документации об этом не сказано. Поэтому держите версии пакетов закреплёнными, а после обновления первым делом проверяйте код с version="v3".
Об этом же напоминает предупреждение LangChainBetaWarning. В примерах с агентом оно выводится дважды, и обе строки указывают на langgraph/pregel/main.py. Это пометка о бета-версии, код после неё работает как обычно. У модели без агента, как в примере 1, предупреждения нет: там бета-метод вызывается изнутри самого LangChain, а для таких вызовов предупреждение не выводится.
Общая сборка модели
Модуль сборки модели тот же, что в уроках 4 и 5. Положите его рядом с примерами под именем course_model.py, и дальше каждый пример вызывает его одной строкой.
"""Общая сборка модели для примеров урока 8.
Тот же модуль, что course_model.py уроков 4 и 5, без изменений: у build_model
есть необязательный первый аргумент с именем модели, а рядом лежит
gateway_kwargs() для примеров, которые собирают модель сами.
Файл .env берётся тот же, что в уроке 0. Положите его рядом с этой папкой или
выше по дереву: load_dotenv() ищет файл начиная с папки этого модуля и поднимается
вверх.
"""
import os
from dotenv import load_dotenv
from langchain.chat_models import init_chat_model
load_dotenv()
def gateway_kwargs():
"""Возвращает аргументы доступа к провайдеру: имя, адрес, ключ.
Нужны примерам, которые вызывают init_chat_model сами. У настраиваемой модели
имени модели при создании нет, поэтому build_model ей не подходит.
"""
base_url = os.getenv("MODEL_BASE_URL")
if base_url:
return {
"model_provider": "openai",
"base_url": base_url,
"api_key": os.environ["OPENAI_API_KEY"],
}
return {}
def build_model(model_name=None, **kwargs):
"""Собирает модель курса.
model_name без значения означает модель из переменной MODEL_NAME. Явное имя
нужно примерам, где моделей в приложении больше одной.
Все именованные аргументы уходят в init_chat_model как есть: temperature,
max_tokens, timeout, max_retries, rate_limiter, profile и прочее из раздела
Parameters.
"""
model_name = model_name or os.environ["MODEL_NAME"]
access = gateway_kwargs()
if access:
# Путь для любого адреса, совместимого с OpenAI Chat Completions API.
return init_chat_model(model=model_name, **access, **kwargs)
# Путь напрямую к провайдеру.
return init_chat_model(model_name, **kwargs)
Первый прогон: поток событий у одной модели
Первый пример без агента: метод stream_events(version="v3") есть и у самой модели. На одной модели разница между объектом прогона и итератором из урока 7 видна в нескольких строках кода.
Пример 01_model_stream.py
"""Пример 1 урока 8: поток событий у одной модели, без агента.
Метод stream_events(version="v3") у модели возвращает не итератор, а объект
ChatModelStream с проекциями: .text, .reasoning, .tool_calls,
.output. Итоговое сообщение и расход токенов лежат в .output.
"""
from course_model import build_model
model = build_model(temperature=0, max_tokens=1024)
stream = model.stream_events(
"Назовите три версии протокола HTTP, по одной в строке, без пояснений.",
version="v3",
)
print("ТИП ОБЪЕКТА ПРОГОНА:", type(stream).__name__)
deltas = []
for delta in stream.text:
deltas.append(delta)
print("ДЕЛЬТ ПРИШЛО:", len(deltas))
print("ПЕРВЫЕ ТРИ ДЕЛЬТЫ:", deltas[:3])
final = stream.output
print("ТИП ИТОГОВОГО СООБЩЕНИЯ:", type(final).__name__)
print("ЗНАКОВ В ОТВЕТЕ:", len(final.text))
print("РАСХОД:", final.usage_metadata)
print("ОТВЕТ:")
print(final.text)
# Вывод:
# ТИП ОБЪЕКТА ПРОГОНА: ChatModelStream
# ДЕЛЬТ ПРИШЛО: 4
# ПЕРВЫЕ ТРИ ДЕЛЬТЫ: ['HTTP/1.0', 'HTTP/1', '.1 HTTP/2.']
# ТИП ИТОГОВОГО СООБЩЕНИЯ: AIMessage
# ЗНАКОВ В ОТВЕТЕ: 26
# РАСХОД: {'input_tokens': 26, 'output_tokens': 149, 'total_tokens': 175}
# ОТВЕТ:
# HTTP/1.0
# HTTP/1.1
# HTTP/2.0
Вызов вернул объект прогона, и до первого обхода запрос ещё не ушёл. Прогон двигает тот, кто читает: каждый шаг вашего цикла for забирает следующее событие, и в фоне ничего не выполняется.
stream.text это проекция текста ответа. При обходе она отдаёт дельты: дельта это текст, который добавил очередной chunk. Chunk без текста дельты не даёт, поэтому дельт меньше, чем chunks. str(stream.text) даёт готовый текст целиком. Так же устроены stream.reasoning для рассуждения модели и stream.tool_calls для chunks аргументов вызова. Только готовый результат у stream.tool_calls берётся вызовом .get(), это список собранных вызовов. У модели курса stream.reasoning останется пустым: рассуждение приходит отдельным полем reasoning_content, и ChatOpenAI его не извлекает (урок 7). Но токены рассуждения входят в расход, поэтому ответ из трёх строк занял 149 выходных токенов.
stream.output это собранное сообщение AIMessage, и расход токенов в Python берётся отсюда, из output.usage_metadata. В документации упомянута ещё проекция usage, но она есть только в версии для JavaScript.
Агент: одно сообщение на каждый вызов модели
За один запуск агент вызывает модель несколько раз. В примере ниже вызовов два: первый ответ модели содержит вызов инструмента, второй, уже с результатом инструмента, содержит ответ пользователю. Проекция stream.messages отдаёт на каждый вызов модели отдельное сообщение с теми же проекциями, что в примере 1: .text, .reasoning, .tool_calls, .output. Вдобавок у каждого сообщения есть поле node с именем узла графа, из которого оно пришло.
Пример 02_agent_messages.py
"""Пример 2 урока 8: проекция stream.messages у агента и имя узла.
У агента вызовов модели за один запуск несколько, и stream.messages отдаёт по
одному объекту на каждый вызов. У каждого объекта есть .node, имя узла графа,
из которого пришло сообщение. В текущей версии узел модели называется "model",
а не "agent", и строки "ОТБОР ПО ИМЕНИ" в выводе показывают разницу числом.
"""
from collections import Counter
from langchain.agents import create_agent
from course_model import build_model
def get_weather(city: str) -> str:
"""Возвращает погоду в указанном городе."""
return f"В городе {city} всегда солнечно!"
agent = create_agent(
model=build_model(temperature=0, max_tokens=1024),
tools=[get_weather],
system_prompt="Отвечайте по-русски и коротко.",
)
stream = agent.stream_events(
{"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]},
version="v3",
)
rows = []
for message in stream.messages:
# Дельты считаются обходом .text, полное сообщение берётся из .output.
deltas = sum(1 for _ in message.text)
output = message.output
calls = [call["name"] for call in output.tool_calls]
rows.append((message.node, deltas, len(output.text), calls))
for number, (node, deltas, length, calls) in enumerate(rows, start=1):
print(f"сообщение {number}: узел={node!r} дельт={deltas} знаков={length} вызовы={calls}")
print()
print("СООБЩЕНИЙ ПО УЗЛАМ:", dict(Counter(node for node, *_ in rows)))
print("ОТБОР ПО ИМЕНИ 'agent':", sum(1 for node, *_ in rows if node == "agent"))
print("ОТБОР ПО ИМЕНИ 'model':", sum(1 for node, *_ in rows if node == "model"))
final = stream.output
print("СООБЩЕНИЙ В ИТОГОВОМ СОСТОЯНИИ:", len(final["messages"]))
print("ОТВЕТ:", final["messages"][-1].text)
# Вывод:
# ...\site-packages\langgraph\pregel\main.py:3708: LangChainBetaWarning: The v3 streaming protocol on Pregel is experimental.
# return self._pregel_stream_v3(
# ...\site-packages\langgraph\pregel\main.py:3558: LangChainBetaWarning: The v3 streaming protocol on Pregel is experimental.
# return GraphRunStream(graph_iter, mux)
# сообщение 1: узел='model' дельт=0 знаков=0 вызовы=['get_weather']
# сообщение 2: узел='model' дельт=7 знаков=35 вызовы=[]
#
# СООБЩЕНИЙ ПО УЗЛАМ: {'model': 2}
# ОТБОР ПО ИМЕНИ 'agent': 0
# ОТБОР ПО ИМЕНИ 'model': 2
# СООБЩЕНИЙ В ИТОГОВОМ СОСТОЯНИИ: 4
# ОТВЕТ: В Сан-Франциско всегда солнечно! ☀️
По выводу видно:
1) у первого сообщения знаков=0 и вызовы=['get_weather'], у второго текст ответа и пустой список вызовов. Это цикл агента из урока 2, разложенный по шагам
2) message.output отдаёт собранный AIMessage, когда сообщение закончилось. Обращение к нему раньше само дочитывает сообщение до конца. В примере дельты к этому моменту уже прочитаны, и итог приходит сразу
3) stream.output после цикла отдаёт итоговое состояние агента, словарь с ключом messages. В нём четыре сообщения: вопрос, вызов инструмента, результат инструмента и ответ
Имя узла сменилось с agent на model
Посмотрите на две последние строки отбора в выводе примера 2.
До LangChain 1.0 узел модели назывался agent. В версии 1.0 его переименовали в model, чтобы имя точнее описывало назначение узла, а узел инструментов по-прежнему называется tools. Отбор по agent остался в старых проектах и статьях, и перенесённый в поток событий он выглядит так.
# Отбор из старого кода не сработает.
for message in stream.messages:
if message.node != "agent":
continue
print(message.text)
Такой код отработает без исключения и без предупреждения. Условие message.node != "agent" верно для каждого сообщения, continue пропустит их все, и экран останется пустым. Причина в одной строке "agent" в условии, и ни одна ошибка на неё не укажет.
Это переименование касается и режимов stream из урока 7. В режиме messages отбор metadata["langgraph_node"] == "agent" тоже даёт пустой вывод без ошибки. В режиме updates правки приходят под ключом model, поэтому обращение chunk["data"]["agent"] падает с KeyError.
Имена узлов в отборе не пишите по памяти: напечатайте message.node на своём агенте, как в примере 2, и отбирайте по тому, что увидели. С middleware это особенно важно: middleware с хуками добавляет в граф свои узлы с именами вида ИмяКласса.before_model.
Что можно читать из потока
У объекта прогона агента есть такие проекции и свойства.
| Что читать | Что даёт |
|---|---|
stream.messages |
сообщения модели, по одному на каждый вызов модели |
stream.values |
снимки состояния агента |
stream.tool_calls |
исполнение инструментов: вход, дельты вывода, итог, ошибка |
stream.output |
итоговое состояние агента |
stream.subgraphs |
вложенные графы и агенты |
stream.lifecycle |
статус прогона, вложенных графов и субагентов |
stream.subagents |
именованные субагенты, им отведён урок 19 (выйдет позже) |
stream.extensions |
все проекции по ключам, встроенные и свои |
обход самого stream |
сырые события протокола со всеми полями |
У каждого сообщения из stream.messages свои проекции, те же, что в примере 1.
| Что читать | Что даёт |
|---|---|
message.text |
дельты текста и итоговый текст |
message.reasoning |
дельты рассуждения у моделей, которые его отдают |
message.tool_calls |
chunks аргументов вызова, готовые вызовы через .get() |
message.output |
собранное сообщение после окончания вызова |
Для остановки на подтверждение человеком есть пара stream.interrupted и stream.interrupts. Обе доводят прогон до конца и отдают то же прерывание, которое в уроке 17 достаётся из результата invoke. Поэтому отдельного разбора в курсе у них нет.
Сырые события: поля, каналы, счётчик
Проекции покрывают частые случаи. Когда нужного вида нет, обходят сам объект прогона и получают события протокола как есть.
Каждое событие, это словарь. Главные поля в нём такие.
1) seq: строго растущий номер внутри прогона
2) method: имя канала
3) params.namespace: путь от корневого графа до места, откуда событие пришло
4) params.data: содержимое, форма которого зависит от канала
Время в событии тоже есть, поле params.timestamp, но порядок определяйте по seq: время берётся с системных часов и может расходиться с номером события.
Каналы такие: values, updates, messages, tools, lifecycle, checkpoints, input, tasks и custom, плюс custom:<имя> для своих преобразователей.
Даже короткий запуск даёт десятки событий. Поэтому пример печатает сводку: первые восемь событий и счёт по каналам.
Пример 03_raw_events.py
"""Пример 3 урока 8: сырые события протокола и счётчик по каналам.
Объект прогона можно обходить сам по себе, и тогда приходят события
протокола. У короткого запуска агента их десятки, поэтому скрипт печатает
сводку: первые восемь событий и счёт по каналам.
"""
from collections import Counter
from langchain.agents import create_agent
from course_model import build_model
def get_weather(city: str) -> str:
"""Возвращает погоду в указанном городе."""
return f"В городе {city} всегда солнечно!"
agent = create_agent(
model=build_model(temperature=0, max_tokens=1024),
tools=[get_weather],
system_prompt="Отвечайте по-русски и коротко.",
)
stream = agent.stream_events(
{"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]},
version="v3",
)
counts = Counter()
first_rows = []
for event in stream:
counts[event["method"]] += 1
if len(first_rows) < 8:
params = event["params"]
first_rows.append((event.get("seq"), event["method"], params["namespace"]))
print("ПЕРВЫЕ ВОСЕМЬ СОБЫТИЙ")
for seq, method, namespace in first_rows:
print(f" seq={seq} method={method!r} namespace={namespace}")
print()
print("ВСЕГО СОБЫТИЙ:", sum(counts.values()))
print("ПО КАНАЛАМ:")
for method, number in counts.most_common():
print(f" {method:<12} {number}")
final = stream.output
print()
print("КЛЮЧИ ИТОГОВОГО СОСТОЯНИЯ:", list(final))
print("СООБЩЕНИЙ В СОСТОЯНИИ:", len(final["messages"]))
# Вывод:
# [вывод подрезан: те же два предупреждения LangChainBetaWarning, что в примере 2]
# ПЕРВЫЕ ВОСЕМЬ СОБЫТИЙ
# seq=1 method='values' namespace=[]
# seq=2 method='messages' namespace=[]
# seq=3 method='messages' namespace=[]
# seq=4 method='messages' namespace=[]
# seq=5 method='messages' namespace=[]
# seq=6 method='messages' namespace=[]
# seq=7 method='messages' namespace=[]
# seq=8 method='messages' namespace=[]
#
# ВСЕГО СОБЫТИЙ: 28
# ПО КАНАЛАМ:
# messages 22
# values 4
# tools 2
#
# КЛЮЧИ ИТОГОВОГО СОСТОЯНИЯ: ['messages']
# СООБЩЕНИЙ В СОСТОЯНИИ: 4
Посмотрите на пустой список в колонке namespace у первых событий. Пустой путь означает корневой граф. У вложенного вызова там появились бы сегменты вида имя:идентификатор, по одному на уровень вложенности. Для вложенных графов отбор по ним уже сделан в проекции stream.subgraphs.
Канал messages устроен как поток контент-блоков с явными границами. Сначала идёт message-start, затем на каждый блок content-block-start, ноль или больше content-block-delta и content-block-finish, в конце message-finish. Из-за этих границ текст, рассуждение, вызов инструмента и картинка разбираются одинаково, без подстройки под провайдера. Проекция stream.messages собрана из таких событий, и обходить их руками стоит только ради точного порядка прибытия разного содержимого внутри одного сообщения.
Жизненный цикл инструмента
Вызов инструмента виден в потоке дважды, и это разные вещи.
1) message.tool_calls это вызов в том виде, в каком его пишет модель. Chunks аргументов приходят по мере генерации
2) stream.tool_calls это исполнение вызова. Отдельный объект на каждый запущенный инструмент, с именем, аргументами, потоком частичного вывода, итогом и ошибкой
Проекция stream.tool_calls собирается из событий канала tools: tool-started, tool-output-delta, tool-finished и tool-error. В каждом из них есть tool_call_id, тот же идентификатор, что у вызова в сообщении модели. По нему исполнение связывается с вызовом, который написала модель.
Пример 04_tool_calls.py
"""Пример 4 урока 8: жизненный цикл вызова инструмента через stream.tool_calls.
Проекция stream.tool_calls отдаёт по объекту на каждый запущенный инструмент.
В объекте лежат имя, аргументы, поток частичного вывода, итоговый результат и
строка ошибки. Подписка здесь одна, на эту проекцию, и прогон вперёд двигает
её обход.
"""
from langchain.agents import create_agent
from course_model import build_model
def get_weather(city: str) -> str:
"""Возвращает погоду в указанном городе."""
return f"В городе {city} всегда солнечно!"
def get_population(city: str) -> str:
"""Возвращает число жителей указанного города."""
numbers = {"Сан-Франциско": "808 437", "Бостон": "653 833"}
return numbers.get(city, "нет данных")
agent = create_agent(
model=build_model(temperature=0, max_tokens=1024),
tools=[get_weather, get_population],
system_prompt="Отвечайте по-русски и коротко.",
)
stream = agent.stream_events(
{
"messages": [
{
"role": "user",
"content": "Какая погода и сколько жителей в Сан-Франциско?",
}
]
},
version="v3",
)
number = 0
for call in stream.tool_calls:
number += 1
deltas = list(call.output_deltas)
print(f"вызов {number}: {call.tool_name}")
print(f" вход: {call.input}")
print(f" дельт вывода: {len(deltas)}")
print(f" результат: {call.output!r}")
print(f" ошибка: {call.error!r}")
print(f" завершён: {call.completed}")
print()
print("ВЫЗОВОВ ИНСТРУМЕНТОВ ВСЕГО:", number)
print("ОТВЕТ:", stream.output["messages"][-1].text)
# Вывод:
# [вывод подрезан: те же два предупреждения LangChainBetaWarning, что в примере 2]
# вызов 1: get_weather
# вход: {'city': 'Сан-Франциско'}
# дельт вывода: 0
# результат: ToolMessage(content='В городе Сан-Франциско всегда солнечно!', name='get_weather', id='d1d068dc-f129-4f0a-8bc2-7e1cad66f37f', tool_call_id='call_a019b93956b545c29b1a9f7b')
# ошибка: None
# завершён: True
# вызов 2: get_population
# вход: {'city': 'Сан-Франциско'}
# дельт вывода: 0
# результат: ToolMessage(content='808 437', name='get_population', tool_call_id='call_25c1d6d12e244de2a21cd059')
# ошибка: None
# завершён: True
#
# ВЫЗОВОВ ИНСТРУМЕНТОВ ВСЕГО: 2
# ОТВЕТ: Погода в Сан-Франциско: всегда солнечно. Население: 808 437 человек.
Модель вызывает инструмент не всегда, даже при нулевой температуре. Бывает, что она отвечает сама и выдумывает погоду с населением, и тогда проекция stream.tool_calls пуста, а в итоге стоит ВЫЗОВОВ ИНСТРУМЕНТОВ ВСЕГО: 0.
Дельт вывода у обычного инструмента не будет: функция возвращает результат разом. Событие tool-output-delta появляется, когда код инструмента сам отправляет часть вывода вызовом runtime.emit_output_delta(...). Параметр runtime с объектом ToolRuntime разобран в уроке 10. Поле error заполняется, когда инструмент упал. Что при падении происходит с прогоном и как его обработать, разобрано в уроке 9. Там ошибку читают по ToolMessage и его полю status.
Проекция stream.tool_calls у агента есть всегда: create_agent добавляет ToolCallTransformer в список преобразователей при сборке графа. В графе, собранном без create_agent, этой проекции нет. Её подключают так же, как свой преобразователь в примере 6: класс ToolCallTransformer из langgraph.prebuilt передают в stream_events аргументом transformers.
Две проекции сразу: подписка решает
Проекция копит значения только с того момента, как у неё появился подписчик. Поэтому в синхронном коде полный обход stream.messages доводит прогон до конца, и вторая проекция, к которой вы обратились следом, окажется пустой. Пустой потому, что во время прогона её никто не слушал.
Для чтения нескольких проекций сразу в синхронном коде есть stream.interleave(...). Метод принимает имена проекций и отдаёт пары "имя, значение" в порядке прибытия, подписываясь на все сразу. В асинхронном коде ту же задачу решает astream_events с несколькими потребителями через asyncio.gather.
Пример запускает одного и того же агента дважды. Первый раз он читает stream.messages и stream.tool_calls двумя циклами подряд, второй раз читает их же вместе с stream.values одним обходом через interleave.
Пример 05_interleave.py
"""Пример 5 урока 8: две проекции подряд против interleave.
Проекция накапливает значения только с того момента, как у неё появился
подписчик. Поэтому в синхронном коде два обхода подряд и один обход через
interleave дают разные числа. Скрипт прогоняет одного и того же агента дважды
и печатает обе сводки.
"""
from langchain.agents import create_agent
from course_model import build_model
QUESTION = {
"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]
}
def get_weather(city: str) -> str:
"""Возвращает погоду в указанном городе."""
return f"В городе {city} всегда солнечно!"
agent = create_agent(
model=build_model(temperature=0, max_tokens=1024),
tools=[get_weather],
system_prompt="Отвечайте по-русски и коротко.",
)
print("ДВА ОБХОДА ПОДРЯД")
stream = agent.stream_events(QUESTION, version="v3")
messages = [message.node for message in stream.messages]
calls = [call.tool_name for call in stream.tool_calls]
print(" сообщений:", len(messages), messages)
print(" вызовов инструментов:", len(calls), calls)
print()
print("ОДИН ОБХОД ЧЕРЕЗ INTERLEAVE")
stream = agent.stream_events(QUESTION, version="v3")
order = []
for name, item in stream.interleave("messages", "tool_calls", "values"):
if name == "messages":
order.append(f"messages(узел {item.node})")
elif name == "tool_calls":
order.append(f"tool_calls({item.tool_name})")
else:
order.append(f"values(сообщений {len(item['messages'])})")
for step, row in enumerate(order, start=1):
print(f" {step}. {row}")
print(" всего записей:", len(order))
# Вывод:
# [вывод подрезан: те же два предупреждения LangChainBetaWarning, что в примере 2]
# ДВА ОБХОДА ПОДРЯД
# сообщений: 2 ['model', 'model']
# вызовов инструментов: 0 []
#
# ОДИН ОБХОД ЧЕРЕЗ INTERLEAVE
# 1. values(сообщений 1)
# 2. messages(узел model)
# 3. values(сообщений 2)
# 4. tool_calls(get_weather)
# 5. values(сообщений 3)
# 6. messages(узел model)
# 7. values(сообщений 4)
# всего записей: 7
Сравните вызовы инструментов в двух сводках: в первой их ноль, во второй есть строка tool_calls(get_weather). Агент и вопрос одинаковые, отличается только способ чтения, и разница объясняется тем, кто был подписан в момент прогона.
Вторая половина вывода показывает порядок прибытия, и он же порядок работы агента. Снимок с вопросом, сообщение модели с вызовом, снимок, исполнение инструмента, снимок, ответ модели, итоговый снимок. Для показа хода работы пользователю это самая удобная форма.
У проекции объекта прогона один подписчик. Второй обход той же проекции падает с RuntimeError и подсказкой про tee(n). Если проекцию нужно прочитать дважды, вызовите, например, stream.messages.tee(2) до первого обхода: он вернёт два независимых итератора. Проекции одного сообщения, message.text, message.reasoning и message.tool_calls, под это правило не попадают: они хранят дельты и при каждом обходе отдают их с начала.
Свой преобразователь потока
Встроенные проекции stream.messages, stream.values и stream.tool_calls собирают преобразователи потока, подклассы StreamTransformer. Свой преобразователь пишется так же.
Преобразователь получает каждое событие прогона, копит в себе нужные данные и публикует их своей проекцией. Свой пишут тогда, когда встроенные проекции не дают нужного приложению вида: прогресс загрузки, счётчики, артефакты, события предметной области. Методы StreamTransformer такие.
1) init() возвращает словарь проекций. Его ключи становятся ключами stream.extensions
2) process(event) смотрит каждое событие. Возвращает True, чтобы оставить событие в общем потоке, и False, чтобы погасить его
3) finalize() вызывается один раз при нормальном окончании прогона, до закрытия каналов
4) fail(err) вызывается вместо finalize(), если прогон закончился ошибкой. В err приходит само исключение
Кроме методов у преобразователя есть атрибут required_stream_modes, список режимов потока, которые ему нужны. Перед прогоном движок собирает эти режимы со всех преобразователей и запрашивает у графа только их. Режим, который никто не запросил, граф не отдаёт, и ваш process таких событий не увидит.
Часть режимов уже запрошена: встроенные преобразователи берут values, messages и tasks, агент добавляет tools. Поэтому объявление ("tools",) в примере ниже ничего не меняет. А режим updates не запрашивает никто, и без своего объявления этих событий не будет.
Объявление только добавляет режимы. process всё равно получает все события прогона, поэтому свои отбирайте по event["method"].
Значения проекции публикуются через StreamChannel, канал именованный или безымянный.
1) StreamChannel() без имени: значения видны только через stream.extensions. Сюда можно класть любые объекты Python
2) StreamChannel("имя"): значения видны и через stream.extensions, и в общем потоке событий, каждое отдельным событием custom:имя. Поэтому они должны быть сериализуемыми, как любое событие протокола
Преобразователь этого урока собирает из канала tools ленту активности. Положите его рядом с примерами под именем tool_activity.py, его используют примеры 6 и 7.
"""Свой преобразователь потока для примеров 6 и 7 урока 8.
Преобразователь смотрит канал tools и собирает две проекции:
1) tool_activity, именованный канал. Каждая запись одновременно уходит в общий
поток событий отдельным событием custom:tool_activity
2) tool_totals, безымянный канал. Туда в finalize() кладётся один словарь со
счётчиком событий за прогон
Имя инструмента приходит только в событии tool-started, поэтому оно
запоминается по tool_call_id и подставляется в остальные три события.
"""
from collections import Counter
from typing import Any, TypedDict
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
class Activity(TypedDict):
"""Одна запись проекции tool_activity."""
tool: str
event: str
class ToolActivityTransformer(StreamTransformer):
"""Проекция "что происходит с инструментами" поверх канала tools."""
required_stream_modes = ("tools",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.activity: StreamChannel[Activity] = StreamChannel("tool_activity")
self.totals: StreamChannel[dict[str, int]] = StreamChannel()
self.names: dict[str, str] = {}
self.counts: Counter[str] = Counter()
def init(self) -> dict[str, Any]:
"""Ключи этого словаря станут ключами stream.extensions."""
return {"tool_activity": self.activity, "tool_totals": self.totals}
def process(self, event: ProtocolEvent) -> bool:
"""Смотрит каждое событие прогона и отбирает свои по имени канала."""
if event["method"] != "tools":
return True
data = event["params"]["data"]
if not isinstance(data, dict):
return True
kind = data.get("event")
call_id = data.get("tool_call_id")
if kind is None or call_id is None:
return True
if kind == "tool-started":
self.names[call_id] = data.get("tool_name", "")
self.counts[kind] += 1
self.activity.push({"tool": self.names.get(call_id, ""), "event": kind})
# True оставляет событие в общем потоке. False его гасит.
return True
def finalize(self) -> None:
"""Вызывается один раз в конце прогона, до закрытия каналов."""
self.totals.push(dict(self.counts))
Словарь names нужен потому, что имя инструмента приходит только в событии tool-started. В остальных трёх событиях есть лишь tool_call_id, и имя по нему берётся из словаря. Похожий пример в документации отбирает события по полю tool_name и поэтому записывает только начала вызовов. Его ветка для ошибок не сработает: в tool-error имени нет.
В примере 6 преобразователь регистрируется на вызове: класс передаётся в stream_events аргументом transformers, а проекции читаются из stream.extensions по ключам, которые вернул init().
Пример 06_transformer_call.py
"""Пример 6 урока 8: свой преобразователь, зарегистрированный на вызове.
Класс преобразователя передаётся в stream_events аргументом transformers, и его
проекции появляются в stream.extensions под теми же ключами, которые вернул
init(). Подписка на tool_totals оформляется до прогона: значение туда кладётся в
finalize(), и без подписчика оно потерялось бы.
"""
from collections import Counter
from langchain.agents import create_agent
from course_model import build_model
from tool_activity import ToolActivityTransformer
def get_weather(city: str) -> str:
"""Возвращает погоду в указанном городе."""
return f"В городе {city} всегда солнечно!"
def get_population(city: str) -> str:
"""Возвращает число жителей указанного города."""
numbers = {"Сан-Франциско": "808 437", "Бостон": "653 833"}
return numbers.get(city, "нет данных")
agent = create_agent(
model=build_model(temperature=0, max_tokens=1024),
tools=[get_weather, get_population],
system_prompt="Отвечайте по-русски и коротко.",
)
stream = agent.stream_events(
{
"messages": [
{
"role": "user",
"content": "Какая погода и сколько жителей в Сан-Франциско?",
}
]
},
version="v3",
transformers=[ToolActivityTransformer],
)
print("КЛЮЧИ EXTENSIONS:", sorted(stream.extensions))
# Подписка без обхода: значение из finalize() иначе некуда будет положить.
totals = iter(stream.extensions["tool_totals"])
print("ЗАПИСИ ПРОЕКЦИИ tool_activity")
seen = Counter()
for record in stream.extensions["tool_activity"]:
seen[record["event"]] += 1
print(f" {record['event']:<18} {record['tool']}")
print()
print("СЧЁТ НА СТОРОНЕ ЧИТАТЕЛЯ:", dict(seen))
print("СЧЁТ ИЗ finalize():", next(totals, None))
# Вывод:
# [вывод подрезан: те же два предупреждения LangChainBetaWarning, что в примере 2]
# КЛЮЧИ EXTENSIONS: ['lifecycle', 'messages', 'subagents', 'subgraphs', 'tool_activity', 'tool_calls', 'tool_totals', 'values']
# ЗАПИСИ ПРОЕКЦИИ tool_activity
# tool-started get_weather
# tool-finished get_weather
# tool-started get_population
# tool-finished get_population
#
# СЧЁТ НА СТОРОНЕ ЧИТАТЕЛЯ: {'tool-started': 2, 'tool-finished': 2}
# СЧЁТ ИЗ finalize(): {'tool-started': 2, 'tool-finished': 2}
Строка с iter(...) отвечает на то самое правило подписки, из-за которого в примере 5 разошлись числа. Значение в tool_totals кладётся в finalize(), то есть в самом конце прогона. Если подписаться на канал после обхода ленты, прогон к этому моменту уже закончится и класть будет некуда. Вызов iter() подписывает, ничего не читая, а next() потом забирает значение.
В списке ключей extensions вы увидите не только свои две проекции: встроенные лежат там же, и interleave берёт имена оттуда.
Преобразователь в middleware
Регистрация на вызове удобна для разовой отладки. Но тогда transformers=[...] нужно передавать везде, где вызывается агент. Пропустите один вызов, и проекции там не будет.
В middleware преобразователь объявляется один раз. Подкласс AgentMiddleware перечисляет фабрики в атрибуте transformers, и агент добавляет их при сборке графа. Фабрика это функция, которая получает scope, путь до графа, где работает преобразователь, и возвращает экземпляр преобразователя. Сам класс преобразователя тоже подходит как фабрика: его конструктор принимает scope.
scope это пустой кортеж для корневого графа и непустой для вложенного, и для каждого scope создаётся свой экземпляр. Корневой экземпляр видит и события вложенных графов, поэтому ToolCallTransformer сверяет event["params"]["namespace"] со своим scope. У агента без вложенных графов эта сверка не нужна, и в tool_activity.py её нет.
Пример 07_transformer_middleware.py
"""Пример 7 урока 8: тот же преобразователь, зарегистрированный через middleware.
Middleware объявляет фабрики преобразователей атрибутом transformers, и агент
добавляет их при сборке графа. Вызывающему коду про них знать не нужно:
stream_events вызывается без аргумента transformers, а проекция в extensions
уже есть.
"""
from langchain.agents import create_agent
from langchain.agents.middleware import AgentMiddleware
from course_model import build_model
from tool_activity import ToolActivityTransformer
class ToolActivityMiddleware(AgentMiddleware):
"""Middleware без единого хука: он нужен только ради проекции потока."""
transformers = (ToolActivityTransformer,)
def get_weather(city: str) -> str:
"""Возвращает погоду в указанном городе."""
return f"В городе {city} всегда солнечно!"
agent = create_agent(
model=build_model(temperature=0, max_tokens=1024),
tools=[get_weather],
system_prompt="Отвечайте по-русски и коротко.",
middleware=[ToolActivityMiddleware()],
)
stream = agent.stream_events(
{"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]},
version="v3",
)
print("КЛЮЧИ EXTENSIONS:", sorted(stream.extensions))
for record in stream.extensions["tool_activity"]:
print(f" {record['event']:<18} {record['tool']}")
print()
print("ОТВЕТ:", stream.output["messages"][-1].text)
# Второй прогон того же агента: видно ли записи именованного канала
# в общем потоке событий.
stream = agent.stream_events(
{"messages": [{"role": "user", "content": "Какая погода в Бостоне?"}]},
version="v3",
)
methods = [event["method"] for event in stream]
print("СОБЫТИЙ ВСЕГО:", len(methods))
print("ИЗ НИХ custom:tool_activity:", methods.count("custom:tool_activity"))
# Вывод:
# [вывод подрезан: те же два предупреждения LangChainBetaWarning, что в примере 2]
# КЛЮЧИ EXTENSIONS: ['lifecycle', 'messages', 'subagents', 'subgraphs', 'tool_activity', 'tool_calls', 'tool_totals', 'values']
# tool-started get_weather
# tool-finished get_weather
#
# ОТВЕТ: В Сан-Франциско всегда солнечно! 🌞
# СОБЫТИЙ ВСЕГО: 27
# ИЗ НИХ custom:tool_activity: 2
Последняя строка вывода показывает, что записи именованной проекции попадают и в общий поток событий, отдельным каналом custom:tool_activity. Записи безымянной проекции видны только через stream.extensions.
На каждом прогоне преобразователи регистрируются в таком порядке.
1) встроенные преобразователи состояния, сообщений, жизненного цикла и вложенных графов
2) ToolCallTransformer и преобразователь субагентов, их добавляет create_agent
3) фабрики из middleware, в порядке списка middleware
4) фабрики из аргумента transformers у create_agent
5) фабрики из аргумента transformers на вызове stream_events
Порядок важен, когда преобразователь правит содержимое события. Встроенный преобразователь сообщений сразу снимает текст в проекцию, и более поздняя правка туда не попадёт. Для этого у StreamTransformer есть флаг before_builtins: с before_builtins = True преобразователь встаёт впереди встроенных. Такой флаг стоит у преобразователя внутри готового PIIMiddleware, который прячет персональные данные, например адрес почты или номер карты. Он разобран в уроке 17. С apply_to_output=True он вычищает их из дельт текста, рассуждения и аргументов вызова, с apply_to_tool_results=True из вывода инструментов.
Middleware здесь нужен только для того, чтобы зарегистрировать преобразователь. Хуки, порядок middleware и переход прогона к другому шагу через jump_to разобраны в уроках 15 и 16.
Распространённые ошибки
Ошибка: отбор по имени узла agent
# Неправильно: узел так назывался до LangChain 1.0, отбор молча даёт ноль
for message in stream.messages:
if message.node == "agent":
print(message.text)
# Правильно
for message in stream.messages:
if message.node == "model":
print(message.text)
Почему: узел модели переименован в model. Исключения нет, пустой экран объясняется сам собой только после того, как вы напечатаете message.node.
Ошибка: забыли version="v3"
# Неправильно: без version вызов идёт по старой схеме v2 и отдаёт генератор
stream = agent.stream_events({"messages": [...]})
for message in stream.messages: # AttributeError
...
# Правильно
stream = agent.stream_events({"messages": [...]}, version="v3")
Почему: по умолчанию version равен "v2", а старая схема отдаёт итератор словарей, у которого никаких проекций нет.
Ошибка: stream_mode вместе с версией v3
# Неправильно: TypeError, аргумент stream_mode версия v3 не принимает
stream = agent.stream_events(
{"messages": [...]}, version="v3", stream_mode="updates"
)
# Правильно: режимы объявляют преобразователи
stream = agent.stream_events({"messages": [...]}, version="v3")
Почему: набор режимов версия v3 собирает сама, из required_stream_modes всех зарегистрированных преобразователей. Аргументы stream_mode и subgraphs она отвергает намеренно.
Ошибка: второй обход той же проекции
# Неправильно: RuntimeError про уже имеющегося подписчика
nodes = [message.node for message in stream.messages]
texts = [str(message.text) for message in stream.messages]
# Правильно: всё нужное собирается за один проход
rows = [(message.node, str(message.text)) for message in stream.messages]
Почему: у проекции объекта прогона ровно один подписчик, значения снимаются с очереди по мере чтения и не хранятся. Для двух проходов есть tee(n), его вызывают вместо первого обхода.
Практическое задание
Напишите скрипт run_watch.py, который показывает ход работы агента одной лентой и подводит итог прогона числами.
Требования:
1) соберите агента с двумя инструментами. Инструменты могут возвращать выдуманные строки, сеть им не нужна
2) напишите свой преобразователь потока, который считает события канала updates по имени узла и публикует счётчик безымянным каналом в finalize(). В таком событии event["params"]["data"] это словарь с именем узла в ключе, как в правках режима updates урока 7
3) объявите required_stream_modes так, чтобы преобразователь видел нужные события, и отберите свои события по event["method"] внутри process
4) зарегистрируйте преобразователь через middleware, а не аргументом вызова
5) ленту событий печатайте через stream.interleave("messages", "tool_calls", "values"): на сообщение модели строку с именем узла, на вызов инструмента строку с именем и аргументами, на снимок состояния строку с числом сообщений
6) в конце напечатайте счётчик из своей проекции и итоговый ответ агента
7) любой отказ провайдера печатайте строкой с типом исключения, скрипт при этом дорабатывает до конца
Как проверить результат:
1) в ленте есть строки всех трёх видов, и вызов инструмента стоит между двумя сообщениями модели
2) имена узлов в ленте взяты из вывода. Отбор по строке agent в скрипте не встречается ни разу
3) счётчик из finalize() напечатался словарём с узлами model и tools. Если вместо словаря напечаталось None, вы подписались на проекцию после прогона
4) если убрать required_stream_modes, вместо счётчика напечатается пустой словарь {}: канал updates не запрашивает ни один встроенный преобразователь. Верните объявление на место
Подсказка: подписку на итоговый канал оформляйте вызовом iter() до начала обхода ленты, как в примере 6.
Итоги урока
Метод stream_events(version="v3") возвращает объект прогона, и двигает его ваш цикл. Объект раскладывает поток на проекции: сообщения, снимки состояния, исполнение инструментов, итог и ваши собственные. Если проекции не хватает, обходите сам объект: придут сырые события с номером, каналом, путём и содержимым.
С LangChain 1.0 узел модели называется model, и старый отбор по agent без ошибки даёт пустой вывод. Имена узлов берите из message.node своего агента.
Проекция копит значения только с момента подписки. Поэтому несколько проекций сразу в синхронном коде читаются через interleave, а в асинхронном через astream_events и несколько потребителей в asyncio.gather.
Свой преобразователь, это подкласс StreamTransformer с методами init, process, finalize, fail и атрибутом required_stream_modes. Зарегистрированный на вызове, он работает только в этом прогоне. Объявленный в middleware, он работает в каждом прогоне агента. Встроенные преобразователи стоят первыми, а флаг before_builtins ставит ваш впереди них.
Инструменты в этом уроке были видны только снаружи: имя, аргументы, результат. У самого инструмента есть описание, по которому модель его выбирает, схема аргументов и обработка ошибок. Рядом с ними параллельные вызовы и возврат результата сразу пользователю.
В уроке 9, "Инструменты и цикл вызова", разберу декоратор @tool, docstring как часть промпта и схему аргументов на Pydantic. Там же покажу, что происходит при падении внутри инструмента и почему модель иногда вызывает три инструмента сразу.
Код урока
Примеры этого урока лежат в репозитории курса, папка lesson_08. Закреплённые версии, на которых получен вывод в тексте, лежат в requirements.txt в корне репозитория.
Предыдущий урок: Стриминг: поток вместо ожидания
Следующий урок: Инструменты и цикл вызова
Подписывайтесь на мой Telegram канал
Если вам нужен ментор и вы хотите научиться разрабатывать AI агентов, пишите, обсудим условия
Авторизуйтесь, чтобы оставить комментарий.
Нет комментариев.
Тут может быть ваша реклама
Пишите info@aisferaic.ru