Потоковый вывод в LangChain
Реализуйте потоковую передачу токенов через цепочки LCEL, чтобы приложение отображало каждое слово сразу после его получения, а не ожидало полного ответа, сокращая воспринимаемую задержку.
«Потоковый вывод в 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 inastream_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 — локальная установка не требуется.
Все уроки этого курса
- Архитектура LangChain и основные абстракции
- Создание цепочек с LCEL
- Разветвлённые и параллельные цепочки
- Потоковый вывод в LangChain