AI Engineering Academy · Урок

Потоковый вывод в LangChain

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

Урок 4 из 413 шагов

«Потоковый вывод в LangChain» — бесплатный урок AI Engineering Academy на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Engineering Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Engineering Academy содержит 4 уроков всего.

Зачем нужна потоковая передача

Без потоковой передачи пользователи смотрят на пустой экран, пока LLM завершает генерацию, а для длинных ответов это может занимать 5–30 секунд. При потоковой передаче токены появляются по мере генерации, обеспечивая мгновенную обратную связь. Это значительно улучшает воспринимаемую отзывчивость приложения. LCEL в LangChain автоматически передаёт поток через всю цепочку при вызове .stream().

Базовая потоковая передача с помощью .stream()

Каждая цепочка LCEL предоставляет метод .stream(), возвращающий итератор фрагментов. Для цепочки, завершающейся StrOutputParser, каждый фрагмент представляет собой часть строки. Вы перебираете фрагменты и выводите или передаёте их по мере поступления. Потоковая передача происходит на уровне HTTP: каждый токен из API OpenAI сразу после получения передаётся через анализатор.

from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser

chain = (
    ChatPromptTemplate.from_template('Explain {topic} in detail.')
    | ChatOpenAI(model='gpt-4o-mini')
    | StrOutputParser()
)

# Stream tokens to stdout
for chunk in chain.stream({'topic': 'quantum entanglement'}):
    print(chunk, end='', flush=True)
print()  # final newline

Асинхронная потоковая передача с помощью .astream()

.astream() — асинхронная версия .stream(). Она возвращает асинхронный итератор, элементы которого обрабатываются с помощью async for. Это правильный подход в FastAPI, Starlette и других асинхронных веб-фреймворках, где обработчик запроса является сопрограммой. Использование синхронной потоковой передачи в асинхронном обработчике заблокировало бы цикл обработки событий.

import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate

chain = (
    ChatPromptTemplate.from_template('Write a poem about {subject}')
    | ChatOpenAI(model='gpt-4o-mini')
    | StrOutputParser()
)

async def stream_response():
    async for chunk in chain.astream({'subject': 'the ocean'}):
        print(chunk, end='', flush=True)

asyncio.run(stream_response())

Потоковая передача в FastAPI с помощью StreamingResponse

В FastAPI асинхронный генератор оборачивают в StreamingResponse с параметром media_type='text/plain', чтобы передавать текстовые токены в браузер. Для событий, отправляемых сервером (SSE), используйте media_type='text/event-stream' и форматируйте каждый фрагмент как data: ...\n\n. После этого браузер получает токены по мере их генерации, не дожидаясь полного ответа.

from fastapi import FastAPI
from fastapi.responses import StreamingResponse

app = FastAPI()

async def generate_stream(topic: str):
    async for chunk in chain.astream({'topic': topic}):
        yield chunk

@app.get('/stream')
async def stream_endpoint(topic: str):
    return StreamingResponse(
        generate_stream(topic),
        media_type='text/plain'
    )

# SSE format for frontend EventSource
async def sse_stream(topic: str):
    async for chunk in chain.astream({'topic': topic}):
        yield f'data: {chunk}\n\n'

Потоковая передача через промежуточные шаги

Цепочки LCEL передают поток через каждый поддерживающий его шаг. StrOutputParser поддерживает потоковую обработку и сразу передаёт фрагменты дальше. Однако некоторые анализаторы, например JsonOutputParser, должны накопить весь вывод перед его разбором, из-за чего потоковая передача прерывается. LangChain явно показывает такое поведение: если шаг несовместим с потоковой обработкой, он накапливает вывод перед передачей его дальше.

from langchain_core.output_parsers import JsonOutputParser

# This chain does NOT stream token by token
# JsonOutputParser must buffer the full response before parsing JSON
json_chain = (
    ChatPromptTemplate.from_template('Return JSON: {task}')
    | ChatOpenAI(model='gpt-4o-mini')
    | JsonOutputParser()  # buffers until complete
)

# But partial JSON streaming IS possible with streaming_json_parser
for partial in json_chain.stream({'task': 'list 3 colors'}):
    print(partial)  # prints partial dict as it fills in

astream_events для детального управления

.astream_events() предоставляет более детализированный API потоковой передачи, который создаёт события для каждого шага цепочки, а не только для итогового вывода. Каждое событие содержит поле kind (on_chain_start, on_llm_stream, on_chain_end) и данные в поле data. Это позволяет отдельно передавать результаты вызовов инструментов, промежуточные рассуждения и итоговый вывод в разные части пользовательского интерфейса.

async def stream_with_events(question: str):
    async for event in chain.astream_events(
        {'question': question},
        version='v2'
    ):
        kind = event['event']
        if kind == 'on_llm_stream':
            chunk = event['data']['chunk'].content
            print(chunk, end='', flush=True)
        elif kind == 'on_chain_end':
            print('\n[Done]')
        elif kind == 'on_tool_start':
            print(f'\n[Tool: {event["name"]}]')

Буферизация передаваемого потока

Иногда нужно одновременно передавать токены пользователю и сохранять полный ответ для ведения журнала или дальнейшей обработки. Используйте .astream() со списком-накопителем. После завершения цикла объедините фрагменты, чтобы получить полный текст. Этот подход позволяет отображать потоковый вывод в реальном времени и одновременно сохранять полный ответ для аналитики, кэширования или оценки.

async def stream_and_capture(question: str) -> str:
    full_response = []
    async for chunk in chain.astream({'question': question}):
        print(chunk, end='', flush=True)  # stream to user
        full_response.append(chunk)        # also collect
    print()  # newline
    complete = ''.join(full_response)
    await log_response(question, complete)  # log full text
    return complete

Потоковая передача с вызовами инструментов

Когда модель создаёт вызов инструмента в передаваемом потоке ответа, аргументы функции поступают в виде фрагментов токенов. Необходимо буферизовать строку аргументов JSON до завершения вызова инструмента, прежде чем выполнять его. LangChain автоматически обрабатывает это в своих исполнителях агентов, но при создании собственного цикла потоковой передачи нужно проверить finish_reason и накапливать фрагменты tool_call.function.arguments.

from openai import AsyncOpenAI

client = AsyncOpenAI()

async def stream_with_tools(prompt: str):
    tool_call_buffer = {}
    async with client.chat.completions.stream(
        model='gpt-4o-mini',
        messages=[{'role': 'user', 'content': prompt}],
        tools=[weather_tool_schema]
    ) as stream:
        async for chunk in stream:
            delta = chunk.choices[0].delta
            if delta.tool_calls:
                for tc in delta.tool_calls:
                    idx = tc.index
                    if idx not in tool_call_buffer:
                        tool_call_buffer[idx] = ''
                    if tc.function.arguments:
                        tool_call_buffer[idx] += tc.function.arguments

Отмена и ограничение времени при потоковой передаче

Для длительных потоковых ответов необходима поддержка отмены. В асинхронном Python можно отменить asyncio.Task, в котором выполняется поток. В FastAPI фреймворк автоматически обрабатывает отмену при отключении клиента, если используется StreamingResponse. Задайте ограничение времени через параметр timeout клиента OpenAI или оберните поток в asyncio.wait_for(), чтобы прервать его по истечении максимальной продолжительности.

import asyncio

async def stream_with_timeout(question: str, timeout: float = 30.0):
    async def _stream():
        async for chunk in chain.astream({'question': question}):
            yield chunk

    try:
        async for chunk in asyncio.timeout(_stream(), timeout):
            print(chunk, end='', flush=True)
    except asyncio.TimeoutError:
        print('\n[Stream timed out after 30 seconds]')
    except asyncio.CancelledError:
        print('\n[Stream cancelled by client disconnect]')

SSE на стороне клиента с JavaScript

Во внешнем интерфейсе встроенный в браузер API EventSource получает события, отправляемые сервером. Когда конечная точка FastAPI выдаёт фрагменты data: token\n\n, EventSource создаёт событие message для каждого из них. Добавляйте каждый токен в DOM по мере его поступления, чтобы создать эффект печати. Для большего контроля используйте fetch() с response.body.getReader(): это обеспечивает полный доступ к потоковой передаче.

// Frontend JavaScript (not Python)
const source = new EventSource('/stream?topic=quantum+computing');
const outputDiv = document.getElementById('output');

source.onmessage = (event) => {
    outputDiv.textContent += event.data;
};

source.onerror = () => {
    source.close();
    outputDiv.textContent += ' [done]';
};

// Alternative: fetch with ReadableStream
const response = await fetch('/stream?topic=ai');
const reader = response.body.getReader();
while (true) {
    const {done, value} = await reader.read();
    if (done) break;
    outputDiv.textContent += new TextDecoder().decode(value);
}

Рекомендации по потоковой передаче

При реализации потоковой передачи соблюдайте следующие рекомендации: всегда используйте flush=True при выводе в stdout, чтобы предотвратить буферизацию. Установите stream_usage=True, если Вам нужны точные подсчёты токенов во время потоковой передачи. В конце потоков SSE отправляйте маркер data: [DONE]\n\n, чтобы клиент понимал, когда закрыть соединение. Проверяйте конечные точки потоковой передачи с помощью curl --no-buffer, чтобы убедиться, что токены поступают постепенно.

# Complete SSE endpoint with DONE sentinel
async def sse_generator(question: str):
    try:
        async for chunk in chain.astream({'question': question}):
            # Escape any newlines in the chunk
            safe_chunk = chunk.replace('\n', ' ')
            yield f'data: {safe_chunk}\n\n'
    finally:
        yield 'data: [DONE]\n\n'

@app.get('/chat/stream')
async def chat_stream(question: str):
    return StreamingResponse(
        sse_generator(question),
        media_type='text/event-stream',
        headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}
    )

Быстрая проверка

Проверьте, насколько хорошо Вы понимаете потоковый вывод в LangChain.

Итоги урока

В этом уроке Вы узнали, что stream() и astream() позволяют перебирать фрагменты токенов по мере их генерации, устраняя длительное ожидание полного ответа; StreamingResponse в FastAPI в формате SSE доставляет токены клиентам браузера в реальном времени; а astream_events() предоставляет детализированные обработчики событий для каждого шага цепочки, включая вызовы инструментов и промежуточный вывод. Далее мы рассмотрим управление памятью в многоходовых диалогах.

Можно начать бесплатно

Изучай Python с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
30
Уроки
120

Часто задаваемые вопросы

Урок «Потоковый вывод в LangChain» бесплатный?

Да — полный текст урока «Потоковый вывод в LangChain» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Engineering Academy, подпишись на CoddyKit PRO. Курс AI Engineering Academy содержит 4 уроков всего.

Чему я научусь в уроке «Потоковый вывод в LangChain»?

Реализуйте потоковую передачу токенов через цепочки LCEL, чтобы приложение отображало каждое слово сразу после его получения, а не ожидало полного ответа, сокращая воспринимаемую задержку. Ты практикуешь AI Engineering Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать AI Engineering Academy?

Предыдущий опыт не требуется. AI Engineering Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.

Сколько времени занимает урок «Потоковый вывод в LangChain»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке AI Engineering Academy?

Да. Каждый урок AI Engineering Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Архитектура LangChain и основные абстракции
  2. Создание цепочек с LCEL
  3. Разветвлённые и параллельные цепочки
  4. Потоковый вывод в LangChain
← Назад к AI Engineering Academy