Strumieniowanie danych wyjściowych w LangChain
Uczestnicy zaimplementują strumieniowanie tokenów przez łańcuchy LCEL, aby aplikacja wyświetlała każde słowo zaraz po jego nadejściu zamiast czekać na pełną odpowiedź, poprawiając odczuwalne opóźnienie.
Strumieniowanie danych wyjściowych w LangChain to bezpłatna lekcja AI Engineering Academy na CoddyKit. To lekcja 4 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AI Engineering Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AI Engineering Academy zawiera 4 lekcji w sumie.
Dlaczego strumieniowanie ma znaczenie
Bez strumieniowania użytkownicy patrzą na pusty ekran, czekając na zakończenie generowania przez LLM — w przypadku długich odpowiedzi może to potrwać od 5 do 30 sekund. Dzięki strumieniowaniu tokeny pojawiają się w miarę generowania, zapewniając natychmiastową informację zwrotną. Znacznie poprawia to odczuwaną responsywność. LCEL w LangChain automatycznie propaguje strumieniowanie przez cały łańcuch po wywołaniu .stream().
Podstawowe strumieniowanie za pomocą .stream()
Każdy łańcuch LCEL udostępnia metodę .stream(), która zwraca iterator fragmentów. W przypadku łańcucha zakończonego przez StrOutputParser każdy fragment jest częścią tekstu. Należy iterować po fragmentach i wyświetlać je lub zwracać w miarę ich nadejścia. Strumieniowanie odbywa się na poziomie HTTP — każdy token z API OpenAI jest przekazywany przez parser natychmiast po odebraniu.
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 newlineAsynchroniczne strumieniowanie za pomocą .astream()
.astream() to asynchroniczna wersja .stream(). Zwraca iterator asynchroniczny, który należy obsługiwać za pomocą async for. Jest to właściwe podejście w FastAPI, Starlette i innych asynchronicznych frameworkach internetowych, w których program obsługi żądania jest korutyną. Użycie synchronicznego strumieniowania w asynchronicznym programie obsługi blokowałoby pętlę zdarzeń.
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())Strumieniowanie w FastAPI za pomocą StreamingResponse
W FastAPI należy opakować generator asynchroniczny w StreamingResponse z użyciem media_type='text/plain', aby strumieniować tokeny tekstu do przeglądarki. W przypadku zdarzeń wysyłanych przez serwer (SSE) należy użyć media_type='text/event-stream' i formatować każdy fragment jako data: ...\n\n. Przeglądarka otrzymuje wtedy tokeny w miarę ich generowania, bez oczekiwania na kompletną odpowiedź.
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'Strumieniowanie przez kroki pośrednie
Łańcuchy LCEL propagują strumieniowanie przez każdy krok, który je obsługuje. StrOutputParser jest przystosowany do strumieniowania i natychmiast przekazuje fragmenty dalej. Niektóre parsery — na przykład JsonOutputParser — muszą jednak zbuforować cały wynik przed jego przeanalizowaniem, co przerywa strumieniowanie. LangChain jasno określa to zachowanie: jeśli krok nie obsługuje strumieniowania, gromadzi wynik przed przekazaniem go dalej.
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 — szczegółowa kontrola
.astream_events() udostępnia bardziej szczegółowy interfejs strumieniowania, który emituje zdarzenia dla każdego kroku łańcucha, a nie tylko dla końcowego wyniku. Każde zdarzenie ma pole kind (on_chain_start, on_llm_stream, on_chain_end) oraz dane w polu data. Pozwala to strumieniować wyniki wywołań narzędzi, rozumowanie pośrednie i końcowy wynik osobno do różnych części interfejsu użytkownika.
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"]}]')Buforowanie strumieniowanego wyniku
Czasami trzeba jednocześnie strumieniować tokeny do użytkownika i przechwycić kompletną odpowiedź na potrzeby rejestrowania lub dalszego przetwarzania. Proszę użyć .astream() wraz z akumulatorem w postaci listy. Po zakończeniu pętli należy połączyć fragmenty, aby uzyskać pełny tekst. Ten wzorzec pozwala wyświetlać wynik strumieniowo w czasie rzeczywistym, a jednocześnie przechowywać kompletną odpowiedź na potrzeby analityki, buforowania lub oceny.
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 completeStrumieniowanie z wywołaniami narzędzi
Gdy model generuje wywołanie narzędzia w strumieniowanej odpowiedzi, argumenty funkcji docierają jako fragmenty tokenów. Przed wykonaniem wywołania należy buforować ciąg argumentów JSON aż do zakończenia wywołania narzędzia. LangChain obsługuje to automatycznie w swoich wykonawcach agentów, ale podczas tworzenia własnej pętli strumieniowania trzeba sprawdzać finish_reason i gromadzić fragmenty 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.argumentsAnulowanie i limit czasu podczas strumieniowania
Długie odpowiedzi strumieniowane wymagają obsługi anulowania. W asynchronicznym Pythonie można anulować asyncio.Task opakowujący strumień. W FastAPI framework automatycznie obsługuje anulowanie po rozłączeniu klienta podczas korzystania z StreamingResponse. Limit czasu oczekiwania można ustawić za pomocą parametru timeout klienta OpenAI albo opakować strumień funkcją asyncio.wait_for(), aby przerwać go po maksymalnym czasie.
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 po stronie klienta w JavaScript
Po stronie frontendu przeglądarkowy natywny interfejs EventSource API odbiera zdarzenia wysyłane przez serwer. Gdy punkt końcowy FastAPI emituje fragmenty data: token\n\n, EventSource wywołuje zdarzenie message dla każdego z nich. Należy dołączać każdy token do DOM w miarę jego nadejścia, aby uzyskać efekt pisania na maszynie. Większą kontrolę zapewnia fetch() wraz z response.body.getReader(), który udostępnia pełny dostęp do strumienia.
// 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);
}Najlepsze praktyki dotyczące strumieniowania
Podczas implementowania strumieniowania należy stosować następujące dobre praktyki: podczas wypisywania do stdout zawsze używać flush=True, aby zapobiec buforowaniu. Ustawić stream_usage=True, jeśli podczas strumieniowania potrzebne są dokładne liczby tokenów. Na końcu strumieni SSE emitować znacznik data: [DONE]\n\n, aby klient wiedział, kiedy zamknąć połączenie. Punkty końcowe obsługujące strumieniowanie należy testować za pomocą curl --no-buffer, aby sprawdzić, czy tokeny docierają stopniowo.
# 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'}
)Szybki test
Proszę sprawdzić swoją wiedzę na temat strumieniowania danych wyjściowych w LangChain.
Podsumowanie lekcji
W tej lekcji poznali Państwo: stream() i astream() pozwalają iterować po fragmentach tokenów w miarę ich generowania, eliminując długie oczekiwanie na pełną odpowiedź; StreamingResponse w FastAPI w formacie SSE dostarcza tokeny do klientów przeglądarkowych w czasie rzeczywistym; a astream_events() udostępnia szczegółowe punkty obsługi zdarzeń dla każdego kroku łańcucha, w tym wywołań narzędzi i wyników pośrednich. W następnej części omówimy zarządzanie pamięcią w rozmowach wieloturowych.
Ucz się Python dzięki korepetycjom AI — za darmo
Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.
- Kursy
- 30
- Lekcje
- 120
Często zadawane pytania
Czy lekcja „Strumieniowanie danych wyjściowych w LangChain” jest bezpłatna?
Tak — pełny tekst „Strumieniowanie danych wyjściowych w LangChain” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AI Engineering Academy, przejdź na CoddyKit PRO. Kurs AI Engineering Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „Strumieniowanie danych wyjściowych w LangChain”?
Uczestnicy zaimplementują strumieniowanie tokenów przez łańcuchy LCEL, aby aplikacja wyświetlała każde słowo zaraz po jego nadejściu zamiast czekać na pełną odpowiedź, poprawiając odczuwalne opóźnien… Ćwiczysz AI Engineering Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć AI Engineering Academy?
Nie wymagamy żadnego doświadczenia. AI Engineering Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 4 z 4.
Ile czasu zajmuje lekcja „Strumieniowanie danych wyjściowych w LangChain”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji AI Engineering Academy?
Tak. Każda lekcja AI Engineering Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Architektura LangChain i podstawowe abstrakcje
- Tworzenie łańcuchów za pomocą LCEL
- Rozgałęzianie i równoległe łańcuchy
- Strumieniowanie danych wyjściowych w LangChain