LangChain

Стриминг: поток вместо ожидания | Курс LangChain урок 7

Стриминг: поток вместо ожидания | Курс LangChain урок 7
Михаил Омельченко
Автор
Михаил Омельченко
Опубликовано 03.10.2026
5,0
Views 3

Цель урока: получать от модели и от агента ответ по частям вместо одного ответа в конце. Выбирать режим потока под задачу и собирать из частей целое сообщение. Отбирать из потока то, что нужно показать пользователю.

Необходимые знания:

1) урок 0: окружение собрано, ключ работает, переменные MODEL_NAME и MODEL_BASE_URL заполнены

2) урок 3: типы сообщений, свойство text, контент-блоки и content_blocks

3) урок 4: параметры модели и счёт токенов через usage_metadata

4) урок 5: системный промпт агента, middleware, что уходит в модель на каждом шаге

5) урок 6: response_format и разобранный ответ в ключе structured_response

6) Python на уровне джуниора: итераторы, генераторы, try/except, форматирование строк

Ключевые концепции:

1) stream() у модели отдаёт chunks, объекты AIMessageChunk, они складываются в сообщение оператором +

2) у агента поток идёт режимами: updates, messages, custom и ещё четыре для снимков состояния и отладки

3) форма chunk задаётся параметром version: в v1 она зависит от набора режимов, в v2 всегда одна

4) в режиме messages метаданные говорят, из какого узла графа пришёл chunk

5) свои события в поток пишутся из инструмента через get_stream_writer

6) аргументы вызова инструмента приходят обрывками JSON, полными они становятся только у собранного сообщения

7) метка nostream убирает служебную модель из потока, не отменяя её вызова


Зачем показывать ответ частями

Модель выдаёт ответ по токену, и длинный ответ собирается секунды. Агент с инструментами вдобавок обращается к модели несколько раз за запуск.

Метод invoke возвращает результат в самом конце, а до этого на экране у пользователя пусто. Со стримингом ответ приходит частями, и первые слова появляются на экране, пока остальной ответ ещё генерируется. Часть ответа, которая приходит до готовности целого, называется chunk.

У агента invoke тоже отдаёт ответ в конце прогона. Если по дороге было несколько вызовов инструментов, пользователю нужно видеть, что происходит, до завершения. Поток агента отдаёт и текст ответа, и эти промежуточные шаги.

Дальше в уроке два разных потока, и путать их не надо. Первый это поток самой модели: chunk с текстом. Второй это поток агента: шаги графа, токены, свои события. Начну с модели: там один вызов и один вид chunk.

Общая сборка модели

Модуль сборки модели тот же, что в уроках 4 и 5. Положите его рядом с примерами под именем course_model.py, и дальше каждый пример вызывает его одной строкой.

"""Общая сборка модели для примеров урока 7.

Тот же модуль, что 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)

Поток модели: chunk это не сообщение

Метод stream() возвращает итератор. На каждом обороте цикла приходит chunk, объект AIMessageChunk с обрывком ответа. Сложением chunks собираются в полное сообщение.

У отдельного chunk текст обрывочный, а вызов инструмента лежит в нём обрывком JSON. Поле tool_calls у chunk заполняется попыткой разобрать только этот обрывок. Аргументы там пустые или недописанные, а обрывок, который не разбирается, уходит в invalid_tool_calls, и tool_calls пуст. Цельный текст и полные аргументы вызова есть только у суммы chunks.

Расход токенов лежит в поле usage_metadata, и в потоке он приходит не всегда. Chat Completions от OpenAI присылает его в потоке только по запросу. ChatOpenAI делает такой запрос сам, когда вы идёте к OpenAI напрямую. Если адрес задан параметром base_url или переменной OPENAI_BASE_URL, запроса нет. Тогда расход придёт, только если провайдер присылает его сам. Явно запросить его можно параметром stream_usage=True при создании модели. Модель курса подключена к провайдеру через base_url, поэтому ChatOpenAI расход сам не запрашивает. Пример прогоняет оба случая: без параметра и с stream_usage=True.

Пример 01_model_stream.py

"""Пример 1 урока 7: поток у самой модели и куда девается расход токенов.

Метод stream() возвращает итератор chunks. Chunks складываются оператором + в одно
сообщение, и только у собранного сообщения есть цельный текст. Первые chunks
приходят без текста: модель курса сначала рассуждает, а ChatOpenAI текст
рассуждения не извлекает. Расход токенов в потоке приходит не всегда: у Chat
Completions его надо попросить отдельно.
"""

from course_model import build_model

QUESTION = "Назовите три причины использовать очередь задач, по одному предложению на причину."


def run(title: str, **extra) -> None:
    """Прогоняет один поток и печатает первые chunks с текстом, счётчики и сборку."""
    print(title)
    model = build_model(temperature=0, max_tokens=1024, **extra)

    full = None
    count = 0
    empty = 0
    shown = 0
    with_usage = 0

    for chunk in model.stream(QUESTION):
        count += 1
        if not chunk.text:
            empty += 1
        elif shown < 3:
            shown += 1
            print(f"  chunk {count}: {type(chunk).__name__}, text={chunk.text!r}")
        if chunk.usage_metadata:
            with_usage += 1
        full = chunk if full is None else full + chunk

    print(f"  всего chunks: {count}, из них без текста: {empty}")
    print(f"  chunks с расходом токенов: {with_usage}")

    if full is None:
        # Защитная ветка: провайдер закрыл поток, не прислав ни одного chunk.
        print("  провайдер не прислал ни одного chunk")
        return

    print(f"  тип собранного: {type(full).__name__}")
    print(f"  длина текста: {len(full.text)} знаков")
    print(f"  расход у собранного: {full.usage_metadata}")


run("ПОТОК ПО УМОЛЧАНИЮ")
print()
run("ПОТОК С stream_usage=True", stream_usage=True)

# Вывод:
# ПОТОК ПО УМОЛЧАНИЮ
#   chunk 60: AIMessageChunk, text='В'
#   chunk 61: AIMessageChunk, text='от три причины использовать очередь задач'
#   chunk 62: AIMessageChunk, text=':1. '
#   всего chunks: 104, из них без текста: 63
#   chunks с расходом токенов: 1
#   тип собранного: AIMessageChunk
#   длина текста: 684 знаков
#   расход у собранного: {'input_tokens': 23, 'output_tokens': 461, 'total_tokens': 484, 'input_token_details': {}, 'output_token_details': {}}
#
# ПОТОК С stream_usage=True
#   chunk 51: AIMessageChunk, text='1. **Раз'
#   chunk 52: AIMessageChunk, text='деление ответственности**:'
#   chunk 53: AIMessageChunk, text=' Очередь задач'
#   всего chunks: 90, из них без текста: 54
#   chunks с расходом токенов: 1
#   тип собранного: AIMessageChunk
#   длина текста: 617 знаков
#   расход у собранного: {'input_tokens': 23, 'output_tokens': 398, 'total_tokens': 421, 'input_token_details': {}, 'output_token_details': {}}

Насколько мелко резать ответ, решает провайдер. Ни число chunks, ни количество текста в каждом нигде не обещаны. Числа в выводе выше получены за один прогон, и два прогона подряд дают разное число chunks. Код разбора от размера chunk не зависит: он складывает то, что пришло.

Первые chunks пришли без текста. Модель курса сначала рассуждает и присылает рассуждение отдельным полем reasoning_content. ChatOpenAI это поле не извлекает (урок 3), и от chunk остаётся пустая оболочка. Для пользователя это пауза перед первым словом. Токены рассуждения при этом входят в output_tokens и расходуют тот же потолок max_tokens. Тесный потолок рассуждение съедает целиком, и текста нет вовсе (урок 4), поэтому в примерах урока потолок 1024.

В конце потока тоже приходят chunks без текста: с причиной остановки, с расходом токенов и последний, который дописывает сам LangChain как признак конца сообщения. О признаке подробнее рядом с примером 6. Ваш код должен пропускать пустые chunks.

Расход пришёл в обоих прогонах одним chunk, хотя в первом его никто не запрашивал: провайдер курса присылает его сам. Другой провайдер может его не прислать, и тогда usage_metadata у собранного сообщения будет None. Если приложение считает расход, передавайте stream_usage=True явно.

Сумма chunks остаётся объектом AIMessageChunk. Он наследует AIMessage, поэтому text, tool_calls и usage_metadata у суммы на месте. Если нужен именно AIMessage, его даёт функция message_chunk_to_message из langchain_core.messages.

Режимы потока агента

Агент это скомпилированный граф LangGraph. Его метод stream() принимает набор режимов: что именно вы хотите видеть.

Режимы updates, messages и custom нужны приложению постоянно. Остальные берут, когда нужен снимок состояния целиком или разбор работы графа по шагам.

Режим Что приходит Зачем
updates правки состояния после каждого шага прогресс агента: какой узел отработал и что вернул
values всё состояние после каждого шага когда нужно состояние целиком
messages пары (chunk сообщения, метаданные) токены модели для показа пользователю
custom то, что вы сами записали из узла прогресс инструмента, свои события
checkpoints события создания снимка, работают только с сохранением состояния (урок 12) отладка памяти
tasks старт и финиш задач с результатами и ошибками отладка
debug всё сразу, это сумма checkpoints и tasks отладка

Режимы не исключают друг друга: список в параметре stream_mode даёт их одновременно. Режимы values, checkpoints, tasks и debug разбираются в уроке 21 (выйдет позже).

Прогресс агента: режим updates

Начну с updates: в нём только итог каждого узла, без токенов. После каждого шага приходит словарь "имя узла, что он дописал в состояние". От агента с одним инструментом приходят три события: узел модели с запросом вызова, узел инструментов с результатом, узел модели с финальным ответом.

Имена узлов у агента create_agent: model и tools. Запомните их, они пригодятся в фильтрах.

Пример 02_agent_updates.py

"""Пример 2 урока 7: прогресс агента, режим потока updates.

Режим updates отдаёт по событию после каждого шага агента: что вернул узел model
и что вернул узел tools. Токенов в этом режиме нет, видно только, из каких
шагов состоит прогон.
"""

from langchain.agents import create_agent

from course_model import build_model


def get_weather(city: str) -> str:
    """Возвращает погоду в городе."""
    return f"В городе {city} всегда солнечно!"


def describe(message) -> str:
    """Одна строка на сообщение: тип, вызовы инструментов, начало текста."""
    calls = getattr(message, "tool_calls", None)
    if calls:
        names = [call["name"] for call in calls]
        return f"{type(message).__name__} вызывает {names}"

    text = message.text.replace("
", " ")
    if len(text) > 60:
        text = text[:57] + "..."
    return f"{type(message).__name__} {text!r}"


agent = create_agent(
    model=build_model(temperature=0, max_tokens=1024),
    tools=[get_weather],
    system_prompt="Отвечайте по-русски, одним предложением.",
)

step = 0

for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]},
    stream_mode="updates",
    version="v2",
):
    if chunk["type"] != "updates":
        continue

    for node, update in chunk["data"].items():
        step += 1
        # У служебных ключей вроде __interrupt__ правка приходит не словарём.
        messages = update.get("messages", []) if isinstance(update, dict) else []
        if not messages:
            print(f"шаг {step}: узел {node}, правка без сообщений: {update!r}")
            continue
        for message in messages:
            print(f"шаг {step}: узел {node}, {describe(message)}")

print(f"всего шагов: {step}")

# Вывод:
# шаг 1: узел model, AIMessage вызывает ['get_weather']
# шаг 2: узел tools, ToolMessage 'В городе Сан-Франциско всегда солнечно!'
# шаг 3: узел model, AIMessage 'В Сан-Франциско всегда солнечно!'
# всего шагов: 3

Ветка про "правку без сообщений" нужна для служебных ключей. В правках режима updates приходит, например, __interrupt__ при остановке на подтверждение человеком (урок 17). Данные под ним лежат кортежем, метода get у кортежа нет, и без ветки цикл упал бы с AttributeError.

Две формы chunk: v1 и v2

Тот же цикл и тот же агент, а форма chunk зависит от параметра version.

Форма v1 действует по умолчанию, и в ней вид chunk зависит от параметров вызова. Один режим, и приходят сырые данные. Несколько режимов, и приходят кортежи "режим, данные". Включили подграфы, и впереди добавляется пространство имён: путь к узлу, где вызван подграф. Код разбора переписывается каждый раз, когда вы меняете набор режимов.

В форме v2 chunk всегда словарь с ключами type с именем режима, ns с пространством имён подграфа и data с полезной нагрузкой. У режима values к ним добавляется ключ interrupts. Включается параметром version="v2", требует LangGraph не ниже 1.1.

Пример 03_formats.py

"""Пример 3 урока 7: две формы chunk, v1 по умолчанию и v2 по требованию.

Один и тот же агент, один и тот же набор режимов, разная форма chunk. В v1 форма
зависит от того, сколько режимов вы запросили и нужны ли подграфы. В v2 форма
всегда одна: словарь с ключами type, ns и data.
"""

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="Отвечайте по-русски, одним предложением.",
)

PAYLOAD = {"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]}
MODES = ["updates", "messages"]

print("ФОРМА v1, ПО УМОЛЧАНИЮ")
count = 0
for chunk in agent.stream(PAYLOAD, stream_mode=MODES):
    count += 1
    if count <= 3:
        head = chunk[0] if isinstance(chunk, tuple) else None
        print(f"  chunk {count}: {type(chunk).__name__}, первый элемент {head!r}")
print(f"  всего chunks: {count}")

print()
print("ФОРМА v2")
count = 0
for chunk in agent.stream(PAYLOAD, stream_mode=MODES, version="v2"):
    count += 1
    if count <= 3:
        print(
            f"  chunk {count}: {type(chunk).__name__}, ключи {sorted(chunk)}, "
            f"type={chunk['type']!r}, ns={chunk['ns']!r}"
        )
print(f"  всего chunks: {count}")

# Вывод:
# ФОРМА v1, ПО УМОЛЧАНИЮ
#   chunk 1: tuple, первый элемент 'messages'
#   chunk 2: tuple, первый элемент 'messages'
#   chunk 3: tuple, первый элемент 'messages'
#   всего chunks: 58
#
# ФОРМА v2
#   chunk 1: dict, ключи ['data', 'ns', 'type'], type='messages', ns=()
#   chunk 2: dict, ключи ['data', 'ns', 'type'], type='messages', ns=()
#   chunk 3: dict, ключи ['data', 'ns', 'type'], type='messages', ns=()
#   всего chunks: 55

Разное число chunks в двух прогонах, это разные ответы модели, на число chunks форма не влияет.

Что выбирать. Пишете новый код, берите v2: разбор по chunk["type"] не сломается от добавления режима. Редактор при этом подскажет тип нагрузки: каждому режиму соответствует свой TypedDict из langgraph.types. Читаете чужой код без version, помните, что там v1, и форма зависит от параметров вызова.

Про имя параметра. Параметр version есть и у stream(), и у stream_events, но значит разное. У stream() значения "v1" и "v2" задают форму chunk. У stream_events значения "v1", "v2" и "v3", по умолчанию "v2". Первые два дают общий для объектов LangChain поток событий, где каждое событие, это словарь с полем event. Значение "v3" даёт интерфейс урока 8.

Токены и метаданные: режим messages

Режим messages отдаёт пары: chunk сообщения и метаданные. Метаданные это словарь с описанием источника chunk, главное поле в нём langgraph_node, имя узла графа.

Без отбора этот режим показывать пользователю нельзя.

1) Кроме текста, в потоке приходят вызовы инструментов. Их аргументы идут блоками tool_call_chunk, где args это обрывок JSON вида {" или city.

2) Кроме chunks модели, в поток попадают готовые ToolMessage, поэтому перед разбором стоит проверка isinstance(token, AIMessageChunk).

Пример 4 считает chunks по узлам и по типам сообщений, а обрывки аргументов складывает в список: видно, как JSON собирается по частям.

Пример 04_messages_tokens.py

"""Пример 4 урока 7: режим messages, токены модели и метаданные chunk.

Режим messages отдаёт пары (chunk сообщения, метаданные). Метаданные говорят, из
какого узла графа пришёл chunk. Аргументы инструмента приходят обрывками текста и
собираются в JSON только к концу вызова.
"""

from collections import Counter

from langchain.agents import create_agent
from langchain.messages import AIMessageChunk

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="Отвечайте по-русски, одним предложением.",
)

by_node = Counter()
by_type = Counter()
arg_pieces = []
metadata_keys = []
shown = 0

for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]},
    stream_mode="messages",
    version="v2",
):
    if chunk["type"] != "messages":
        continue

    token, metadata = chunk["data"]
    node = metadata.get("langgraph_node", "?")
    by_node[node] += 1
    by_type[type(token).__name__] += 1
    if not metadata_keys:
        metadata_keys = sorted(metadata)

    blocks = [block["type"] for block in token.content_blocks]

    if isinstance(token, AIMessageChunk):
        for piece in token.tool_call_chunks:
            arg_pieces.append(piece["args"])

    # Пустые chunks рассуждения пропускаются, см. пример 1.
    if blocks and shown < 6:
        shown += 1
        print(f"  {shown}. узел {node}, {type(token).__name__}, блоки {blocks}, text={token.text!r}")

print()
print("CHUNKS ПО УЗЛАМ:", dict(by_node))
print("ТИПЫ СООБЩЕНИЙ В ПОТОКЕ:", dict(by_type))
print("КЛЮЧИ МЕТАДАННЫХ:", metadata_keys)

if arg_pieces:
    print("АРГУМЕНТЫ ИНСТРУМЕНТА ПО CHUNKS:", arg_pieces)
    print("ОНИ ЖЕ, СКЛЕЕННЫЕ:", "".join(piece or "" for piece in arg_pieces))
else:
    # Защитная ветка: модель обошлась без инструмента или отдала вызов целиком.
    print("АРГУМЕНТЫ ИНСТРУМЕНТА ПО CHUNKS: chunks не было")

# Вывод:
#   1. узел model, AIMessageChunk, блоки ['tool_call_chunk'], text=''
#   2. узел model, AIMessageChunk, блоки ['tool_call_chunk'], text=''
#   3. узел model, AIMessageChunk, блоки ['tool_call_chunk'], text=''
#   4. узел model, AIMessageChunk, блоки ['tool_call_chunk'], text=''
#   5. узел model, AIMessageChunk, блоки ['tool_call_chunk'], text=''
#   6. узел model, AIMessageChunk, блоки ['tool_call_chunk'], text=''
#
# CHUNKS ПО УЗЛАМ: {'model': 49, 'tools': 1}
# ТИПЫ СООБЩЕНИЙ В ПОТОКЕ: {'AIMessageChunk': 49, 'ToolMessage': 1}
# КЛЮЧИ МЕТАДАННЫХ: ['checkpoint_ns', 'langgraph_checkpoint_ns', 'langgraph_node', 'langgraph_path', 'langgraph_step', 'langgraph_triggers', 'lc_versions', 'ls_integration', 'ls_max_tokens', 'ls_model_name', 'ls_model_type', 'ls_provider', 'ls_temperature']
# АРГУМЕНТЫ ИНСТРУМЕНТА ПО CHUNKS: [None, '{', '"city": "Сан', '-Ф', 'ранци', 'ско"}']
# ОНИ ЖЕ, СКЛЕЕННЫЕ: {"city": "Сан-Франциско"}

Ключи метаданных смотрите в выводе примера: набор зависит от версии пакета.

Свой канал в поток: режим custom

Токены модели закрывают только часть ожидания. Вторая часть, это время внутри инструмента: запрос к базе, выгрузка файла, обход внешнего API. Модель в это время ничего не присылает, и в режиме messages пусто.

Для этого есть функция get_stream_writer из langgraph.config. Внутри инструмента она возвращает писателя, обычную функцию с одним аргументом. Что вы передадите писателю, то и придёт в поток режимом custom: строка, словарь, объект. Если режим custom не запрошен, вызов писателя ничего не делает.

Писатель берётся из настроек текущего запуска графа. Поэтому инструмент с писателем нельзя вызвать как обычную функцию вне агента, разбор этой ошибки стоит в конце урока.

Пример 5 показывает и поток из двух режимов сразу.

Пример 05_custom_progress.py

"""Пример 5 урока 7: свои сообщения о прогрессе из инструмента.

Пока работает инструмент, модель ничего не присылает, и для пользователя это
выглядит как зависание. Режим custom даёт инструменту собственный канал в поток: что угодно,
что вы передали в writer, приходит в цикл потока.
"""

from langchain.agents import create_agent
from langgraph.config import get_stream_writer

from course_model import build_model

PAGES = 3


def load_orders(day: str) -> str:
    """Загружает заказы за указанный день и возвращает сводку."""
    writer = get_stream_writer()

    loaded = 0
    for page in range(1, PAGES + 1):
        loaded += 40
        writer({"stage": "loading", "page": page, "of": PAGES, "rows": loaded})

    writer({"stage": "done", "rows": loaded})
    return f"за {day} заказов {loaded}, из них оплачено 97"


agent = create_agent(
    model=build_model(temperature=0, max_tokens=1024),
    tools=[load_orders],
    system_prompt="Отвечайте по-русски, одним предложением.",
)

for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "Сколько заказов за вчера?"}]},
    stream_mode=["updates", "custom"],
    version="v2",
):
    if chunk["type"] == "custom":
        event = chunk["data"]
        if event["stage"] == "loading":
            print(f"[прогресс] страница {event['page']} из {event['of']}, строк {event['rows']}")
        else:
            print(f"[прогресс] загрузка закончена, строк {event['rows']}")

    elif chunk["type"] == "updates":
        for node, update in chunk["data"].items():
            messages = update.get("messages", []) if isinstance(update, dict) else []
            for message in messages:
                calls = getattr(message, "tool_calls", None)
                if calls:
                    print(f"[шаг] {node}: вызов {[call['name'] for call in calls]}")
                else:
                    text = message.text.replace("
", " ")
                    print(f"[шаг] {node}: {text[:70]!r}")

# Вывод:
# [шаг] model: вызов ['load_orders']
# [прогресс] страница 1 из 3, строк 40
# [прогресс] страница 2 из 3, строк 80
# [прогресс] страница 3 из 3, строк 120
# [прогресс] загрузка закончена, строк 120
# [шаг] tools: 'за вчера заказов 120, из них оплачено 97'
# [шаг] model: 'За вчера было 120 заказов.'

В примере писатель получает словарь. Цикл потока берёт из него этап, номер страницы и число строк по ключам, текст разбирать не нужно. Формат вывода задаётся в цикле, и для его смены инструмент не правится.

Где стриминг ломает логику обработки

Код, который работал на invoke, в потоке ломается по одной причине: chunk это не сообщение.

1) Разобрать вызов инструмента по chunk нельзя. В chunk лежит tool_call_chunk с обрывком JSON, а tool_calls у chunk собран из этого обрывка, с пустыми или недописанными аргументами. Полные аргументы есть только у собранного сообщения.

2) Проверка ответа опаздывает. Логика, которой нужен ответ целиком, работает только после последнего chunk: фильтр запрещённых слов, разбор схемы, подсчёт длины. А пользователь к этому моменту текст уже прочитал. Отзывать показанное поздно. Либо проверяйте до показа и теряйте выигрыш в скорости, либо показывайте и ставьте ограждения на входе, это урок 17.

3) Структурированного ответа в потоке нет. Схема из урока 6 собирается из тех же обрывков JSON, и до конца вызова модели объекта нет. Готовый объект приходит в правке узла model режима updates, под ключом structured_response.

В примере 6 целое сообщение получается так:

1) chunks складываются в цикле, а конец сообщения ловится по полю chunk_position: у последнего chunk оно равно "last"

2) параллельно идёт режим updates, и готовые сообщения берутся из правок состояния. Сборка здесь не нужна, но способ работает только там, где сообщение попадает в состояние

Пример 06_assemble.py

"""Пример 6 урока 7: chunk это не сообщение.

Разобрать вызов инструмента по одному chunk нельзя: в chunk лежит обрывок JSON.
Целое собирается двумя способами, и оба показаны здесь: сложением chunks в цикле
и чтением готового сообщения из режима updates.
"""

from langchain.agents import create_agent
from langchain.messages import AIMessage, AIMessageChunk, ToolMessage

from course_model import build_model


def get_weather(city: str) -> str:
    """Возвращает погоду в городе."""
    return f"В городе {city} всегда солнечно!"


def show_assembled(message: AIMessageChunk, reason: str) -> None:
    """Печатает то, что видно только у собранного сообщения."""
    text = message.text.replace("
", " ")
    print(f"СОБРАННОЕ СООБЩЕНИЕ ({reason})")
    print(f"  вызовы инструментов: {[call['name'] for call in message.tool_calls]}")
    print(f"  аргументы: {[call['args'] for call in message.tool_calls]}")
    print(f"  текст: {text[:70]!r}")
    print(f"  расход: {message.usage_metadata}")


agent = create_agent(
    model=build_model(temperature=0, max_tokens=1024),
    tools=[get_weather],
    system_prompt="Отвечайте по-русски, одним предложением.",
)

full = None
partial_calls = 0

for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "Какая погода в Сан-Франциско?"}]},
    stream_mode=["messages", "updates"],
    version="v2",
):
    if chunk["type"] == "messages":
        token, _metadata = chunk["data"]
        if not isinstance(token, AIMessageChunk):
            continue

        if token.tool_call_chunks:
            partial_calls += 1
            piece = token.tool_call_chunks[0]
            print(f"chunk вызова: name={piece['name']!r}, args={piece['args']!r}")

        full = token if full is None else full + token

        if token.chunk_position == "last":
            show_assembled(full, "по признаку chunk_position")
            full = None

    elif chunk["type"] == "updates":
        for node, update in chunk["data"].items():
            messages = update.get("messages", []) if isinstance(update, dict) else []
            for message in messages:
                if isinstance(message, AIMessage) and message.tool_calls:
                    print(f"из режима updates, узел {node}: {message.tool_calls}")
                elif isinstance(message, ToolMessage):
                    print(f"из режима updates, узел {node}: ответ {message.text!r}")

if full is not None:
    # Защитная ветка: класс модели обошёл базовый класс LangChain,
    # и последний chunk пришёл без признака chunk_position.
    show_assembled(full, "признака конца не было")

print(f"chunks с обрывками вызова: {partial_calls}")

# Вывод:
# chunk вызова: name='get_weather', args=None
# chunk вызова: name=None, args='{'
# chunk вызова: name=None, args='"city": "С'
# chunk вызова: name=None, args='ан-'
# chunk вызова: name=None, args='Фран'
# chunk вызова: name=None, args='циско"}'
# СОБРАННОЕ СООБЩЕНИЕ (по признаку chunk_position)
#   вызовы инструментов: ['get_weather']
#   аргументы: [{'city': 'Сан-Франциско'}]
#   текст: ''
#   расход: {'input_tokens': 300, 'output_tokens': 82, 'total_tokens': 382, 'input_token_details': {}, 'output_token_details': {}}
# из режима updates, узел model: [{'name': 'get_weather', 'args': {'city': 'Сан-Франциско'}, 'id': 'call_ea222d072a154d1885f9b0d7', 'type': 'tool_call'}]
# из режима updates, узел tools: ответ 'В городе Сан-Франциско всегда солнечно!'
# СОБРАННОЕ СООБЩЕНИЕ (по признаку chunk_position)
#   вызовы инструментов: []
#   аргументы: []
#   текст: 'В Сан-Франциско всегда солнечно!'
#   расход: {'input_tokens': 375, 'output_tokens': 35, 'total_tokens': 410, 'input_token_details': {}, 'output_token_details': {}}
# chunks с обрывками вызова: 6

Признак chunk_position ставит сам LangChain. Если интеграция не пометила последний chunk, базовый класс модели дописывает в конец потока пустой chunk с chunk_position="last". Ветка после цикла нужна для класса модели, который обходит этот механизм.

Что в поток пускать не надо

Служебные вызовы модели в потоке мешают. Классификатор внутри инструмента, проверка ответа ограждением, суммаризатор истории: всё это вызовы модели, и все они попадают в режим messages наравне с ответом пользователю.

Убираются они меткой. Модель с тегом nostream пропадает из режима messages, при этом вызывается и работает как работала.

Рядом стоит другой выключатель с другим действием. Параметр streaming=False при создании модели отключает выдачу по токенам. Если у класса модели такого параметра нет, то же делает disable_streaming=True из базового класса. Ответ такой модели всё равно приходит в режим messages, одним объектом AIMessage в конце вызова, и фильтр isinstance(token, AIMessageChunk) его отбросит. Метка nostream убирает вызов из потока, streaming=False меняет способ получения ответа от провайдера.

Пример 7 прогоняет одного и того же агента дважды, с меткой и без, и считает chunks по узлам.

Пример 07_nostream.py

"""Пример 7 урока 7: служебный вызов модели в потоке не нужен.

Внутри инструмента работает вторая модель, служебная. Её ответ пользователю не
нужен, а в поток режима messages она попадает наравне с основной. Метка nostream
убирает её из потока, не отменяя самого вызова.
"""

from collections import Counter

from langchain.agents import create_agent

from course_model import build_model

INTERNAL = {"model": None}


def classify(text: str) -> str:
    """Определяет тему обращения одним словом."""
    answer = INTERNAL["model"].invoke(
        [{"role": "user", "content": f"Одним словом назовите тему обращения: {text}"}]
    )
    return answer.text.strip()


def run(title: str, tagged: bool) -> None:
    """Прогоняет агента и считает chunks потока по узлам графа."""
    print(title)

    service = build_model(temperature=0, max_tokens=1024)
    INTERNAL["model"] = service.with_config({"tags": ["nostream"]}) if tagged else service

    agent = create_agent(
        model=build_model(temperature=0, max_tokens=1024),
        tools=[classify],
        system_prompt="Отвечайте по-русски, одним предложением. Тему определяйте инструментом.",
    )

    by_node = Counter()
    answer = ""

    for chunk in agent.stream(
        {"messages": [{"role": "user", "content": "У меня дважды списали деньги за заказ 4412."}]},
        stream_mode="messages",
        version="v2",
    ):
        if chunk["type"] != "messages":
            continue
        token, metadata = chunk["data"]
        by_node[metadata.get("langgraph_node", "?")] += 1
        if metadata.get("langgraph_node") == "model":
            answer += token.text

    print(f"  chunks по узлам: {dict(by_node)}")
    print(f"  ответ пользователю: {answer.strip()[:70]!r}")


run("СЛУЖЕБНАЯ МОДЕЛЬ БЕЗ МЕТКИ", tagged=False)
print()
run("СЛУЖЕБНАЯ МОДЕЛЬ С МЕТКОЙ nostream", tagged=True)

# Вывод:
# СЛУЖЕБНАЯ МОДЕЛЬ БЕЗ МЕТКИ
#   chunks по узлам: {'model': 82, 'tools': 138}
#   ответ пользователю: 'Ваш вопрос относится к переплате — ошибка двойного списания за заказ 4'
#
# СЛУЖЕБНАЯ МОДЕЛЬ С МЕТКОЙ nostream
#   chunks по узлам: {'model': 84, 'tools': 1}
#   ответ пользователю: 'По вашему заказу 4412 произошло двойное списание, мы проверим и вернем'

Сравнивайте счётчики по узлу tools. Chunks оттуда приходят двух видов: ответ служебной модели и готовый ToolMessage с результатом инструмента. Служебная модель тоже рассуждает, поэтому её chunks много и почти все без текста. Метка убирает их все, ToolMessage остаётся: метка действует на вызов модели, узел она не трогает.

Что показывать пользователю

Соберу всё в одно. Финальный пример печатает в консоль то, что имеет смысл видеть человеку.

1) текст ответа, по мере готовности

2) строку статуса на вызове инструмента, из режима updates, где вызов уже разобран

3) прогресс загрузки, из режима custom

4) отказ провайдера отдельной строкой, программа при этом не падает

И не печатает обрывки JSON с аргументами, готовые ToolMessage и пустые chunks рассуждения. Заодно пример засекает время до первого показанного знака и до конца ответа.

Пример 08_console_chat.py

"""Пример 8 урока 7: что из потока показывают пользователю.

Три режима сразу: шаги агента, токены модели и прогресс инструмента. Пользователь
видит текст ответа по мере готовности, строку статуса на вызове инструмента и
прогресс загрузки. Обрывки JSON с аргументами в окно не попадают.

Заодно печатается то, ради чего поток и заводят: сколько секунд прошло до первого
показанного знака и сколько до конца ответа.
"""

import time

from langchain.agents import create_agent
from langchain.messages import AIMessage, AIMessageChunk
from langgraph.config import get_stream_writer

from course_model import build_model

PAGES = 3


def load_orders(day: str) -> str:
    """Загружает заказы за указанный день и возвращает сводку."""
    writer = get_stream_writer()

    loaded = 0
    for page in range(1, PAGES + 1):
        loaded += 40
        writer({"stage": "loading", "page": page, "of": PAGES, "rows": loaded})

    return f"за {day} заказов {loaded}, из них оплачено 97, отменено 6"


agent = create_agent(
    model=build_model(temperature=0, max_tokens=1024),
    tools=[load_orders],
    system_prompt="Вы помощник аналитика. Отвечайте по-русски, до двух предложений.",
)

started = time.perf_counter()
first_text_at = None

try:
    for chunk in agent.stream(
        {"messages": [{"role": "user", "content": "Что со вчерашними заказами?"}]},
        stream_mode=["updates", "messages", "custom"],
        version="v2",
    ):
        if chunk["type"] == "messages":
            token, metadata = chunk["data"]
            # В поток режима messages попадают и готовые ToolMessage, и обрывки
            # вызовов. Пользователю показывается только текст ответа модели.
            if not isinstance(token, AIMessageChunk) or not token.text:
                continue
            if metadata.get("langgraph_node") != "model":
                continue
            if first_text_at is None:
                first_text_at = time.perf_counter() - started
            print(token.text, end="", flush=True)

        elif chunk["type"] == "custom":
            event = chunk["data"]
            print(f"
[загрузка] страница {event['page']} из {event['of']}, строк {event['rows']}")

        elif chunk["type"] == "updates":
            for _node, update in chunk["data"].items():
                messages = update.get("messages", []) if isinstance(update, dict) else []
                for message in messages:
                    if isinstance(message, AIMessage) and message.tool_calls:
                        for call in message.tool_calls:
                            print(f"
[инструмент] {call['name']}({call['args']})")

except Exception as error:
    # Отказ провайдера посреди потока: часть ответа уже показана пользователю.
    print(f"
[сбой] {type(error).__name__}: {error}")

print()
if first_text_at is None:
    # Защитная ветка: текста в потоке не было вовсе, показывать было нечего.
    print("первый знак ответа: текста в потоке не было")
else:
    print(f"первый знак ответа через: {first_text_at:.2f} с")
print(f"весь ответ через: {time.perf_counter() - started:.2f} с")

# Вывод:
#
# [инструмент] load_orders({'day': 'вчера'})
#
# [загрузка] страница 1 из 3, строк 40
#
# [загрузка] страница 2 из 3, строк 80
#
# [загрузка] страница 3 из 3, строк 120
# Вчера было 120 заказов: 97 оплачено, 6 отменено.
# первый знак ответа через: 6.12 с
# весь ответ через: 6.26 с

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

В примере перехвачен Exception целиком. Поток рвётся отказом модели, сетью и таймаутом, и в любом из этих случаев пользователю нужна строка с причиной. В приложении на этом месте стоит запись в журнал и разбор типов исключений из урока 23 (выйдет позже).

Чего в этом уроке нет

1) Метод stream_events с version="v3", это урок 8. Вместо разбора chunks по chunk["type"] он даёт отдельные типизированные потоки под сообщения, вызовы инструментов и состояние. В закреплённой версии LangGraph этот интерфейс помечен как экспериментальный.

2) Поток из подграфов и из субагентов, параметр subgraphs=True и поле lc_agent_name в метаданных. Разбирается в уроке 19 (выйдет позже) вместе с самими субагентами.

3) Поток при остановке на подтверждение человеком, ключ __interrupt__ в правках режима updates. Урок 17.

Распространённые ошибки

Разборы ниже даны фрагментами, агент и модель в них уже собраны.

Ошибка: разбор chunk как готового сообщения

# Неправильно: tool_calls у chunk собран из обрывка JSON
for chunk in agent.stream(payload, stream_mode="messages", version="v2"):
    token, metadata = chunk["data"]
    for call in token.tool_calls:      # {} у первых chunks, дальше пустой список
        run_side_effect(call["args"])  # сработает с неполными аргументами

# Правильно: собрать сообщение и разобрать собранное
full = None
for chunk in agent.stream(payload, stream_mode="messages", version="v2"):
    token, metadata = chunk["data"]
    if not isinstance(token, AIMessageChunk):
        continue
    full = token if full is None else full + token
    if token.chunk_position == "last":
        for call in full.tool_calls:
            run_side_effect(call["args"])
        full = None  # следующее сообщение модели собирается с нуля

Почему так: tool_calls у chunk собирается из того обрывка JSON, который успел прийти. Побочное действие сработает с пустыми и недописанными аргументами. Следом в потоке придёт ToolMessage, поля tool_calls у него нет, и цикл упадёт с AttributeError. В правильном варианте full сбрасывается после каждого сообщения. Без сброса ответ модели после инструмента прибавится к первому сообщению, и вызов выполнится повторно.

Ошибка: код разбора написан под один режим, а режимов стало два

# Было и работало: один режим, v1, приходят сырые данные
for chunk in agent.stream(payload, stream_mode="updates"):
    for node, update in chunk.items():
        ...

# Добавили второй режим, и тот же код падает:
# AttributeError: 'tuple' object has no attribute 'items'
for chunk in agent.stream(payload, stream_mode=["updates", "messages"]):
    for node, update in chunk.items():
        ...

# Правильно: форма v2 не зависит от числа режимов
for chunk in agent.stream(payload, stream_mode=["updates", "messages"], version="v2"):
    if chunk["type"] == "updates":
        for node, update in chunk["data"].items():
            ...

Почему так: в v1 форма chunk зависит от набора режимов и от подграфов, в v2 она всегда одна и та же.

Ошибка: писатель потока вызывается вне графа

# Неправильно: тот же инструмент вызывают в тесте напрямую
def load_orders(day: str) -> str:
    writer = get_stream_writer()
    writer({"stage": "loading"})
    return "готово"

load_orders("вчера")  # RuntimeError: Called get_config outside of a runnable context

# Правильно: в тестах вызывать инструмент через агента,
# а получение данных вынести в функцию без писателя
def fetch_orders(day: str) -> str:
    return "готово"

def load_orders(day: str) -> str:
    writer = get_stream_writer()
    writer({"stage": "loading"})
    return fetch_orders(day)

Почему так: вне запуска графа его настроек нет, поэтому текст ошибки говорит про get_config и писателя не называет. Разделение на "инструмент с писателем" и "чистую функцию с логикой" возвращает вам обычный юнит-тест, это урок 22 (выйдет позже).

Практическое задание

Напишите скрипт stream_desk.py, который показывает работу агента в консоли так, как её увидел бы пользователь.

Требования:

1) агент с двумя инструментами: один возвращает выдуманный курс валюты, второй считает сумму заказа. Второй шлёт в поток три события прогресса через get_stream_writer

2) поток в трёх режимах сразу, форма v2

3) текст ответа печатайте по мере поступления, без переводов строки между chunks

4) на вызове инструмента печатайте отдельную строку с именем и разобранными аргументами, взяв их из режима updates

5) обрывки tool_call_chunk в консоль не выводите

6) считайте два числа: сколько chunks пришло всего и сколько из них были текстом для пользователя. Напечатайте их в конце

7) любой отказ провайдера печатайте строкой с типом исключения, программа при этом должна завершиться штатно

Как проверить результат:

1) в выводе есть строка статуса инструмента, и аргументы в ней напечатаны словарём

2) текстовых chunks меньше, чем всех chunks: разница это события режимов updates и custom, обрывки вызовов, пустые chunks рассуждения и сообщения инструментов

3) если в цикле нет отбора по langgraph_node и вы уберёте проверку isinstance(token, AIMessageChunk), в выводе появятся ответы инструментов, а число текстовых chunks вырастет

4) если убрать version="v2", скрипт упадёт на разборе chunk["type"], и это нормально

Подсказка: имя узла модели, это model, имя узла инструментов, это tools, оба видны в метаданных chunk под ключом langgraph_node.

Итоги урока

У модели поток даёт метод stream(), chunks приходят объектами AIMessageChunk и складываются в сообщение оператором +. Цельный текст и полные аргументы вызовов инструментов есть только у собранного сообщения. Расход токенов в потоке приходит, если его шлёт провайдер или если передан stream_usage=True. При работе с OpenAI напрямую ChatOpenAI включает этот параметр сам.

У агента поток разложен по режимам. Три из них для приложения: updates даёт шаги, messages даёт токены с метаданными, custom даёт ваши собственные события из инструмента. Ещё четыре, values, checkpoints, tasks и debug, разбираются в уроке 21 про отладку (он выйдет позже). Режимы включаются списком, а форма chunk задаётся параметром version: берите v2, там у chunk всегда есть ключи type, ns и data.

Chunk это не сообщение, и отсюда всё остальное. Вызов инструмента по chunk не разобрать, проверку ответа целиком в потоке не сделать, схему из урока 6 из обрывков не собрать. Целое достаётся сборкой по chunk_position или чтением готовых сообщений из режима updates. Служебные вызовы модели пользователю не нужны, они убираются меткой nostream.

Чего в этом уроке не хватило. Весь разбор потока держится на ветвлении по chunk["type"] и на ручных проверках типа сообщения. Код получается похожим на разбор сырого протокола, и один и тот же фильтр пишется в каждом цикле. В LangChain 1.3 для этого появился отдельный интерфейс.

В уроке 8, "Поток событий", разберу stream_events версии v3 и совмещение его потоков. Там же свои преобразователи потока в middleware и старый код с отбором по узлу agent: узел модели теперь называется model, и такой отбор молча даёт пустой вывод.

Код урока

Примеры этого урока лежат в репозитории курса, папка lesson_07. Закреплённые версии, на которых получен вывод в тексте, лежат в requirements.txt в корне репозитория.


Предыдущий урок: Структурированный вывод

Следующий урок: Поток событий ---

Подписывайтесь на мой Telegram канал

Если вам нужен ментор и вы хотите научиться разрабатывать AI агентов, пишите, обсудим условия

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

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

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

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

Пишите info@aisferaic.ru

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