LangChain

Поток событий: stream_events версии v3 | Курс LangChain урок 8

Поток событий: stream_events версии v3 | Курс LangChain урок 8
Михаил Омельченко
Автор
Михаил Омельченко
Опубликовано 04.10.2026
5,0
Views 1

Цель урока: читать прогон агента через 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 агентов, пишите, обсудим условия

Авторизуйтесь, чтобы оставить комментарий.

Комментариев: 0

Нет комментариев.

Тут может быть ваша реклама

Пишите info@aisferaic.ru

Похожие туториалы