AI Engineering Academy · Lekcja

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.

Lekcja 4 z 413 kroki

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 newline

Asynchroniczne 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 in

astream_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 complete

Strumieniowanie 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.arguments

Anulowanie 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.

Bezpłatny start

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

  1. Architektura LangChain i podstawowe abstrakcje
  2. Tworzenie łańcuchów za pomocą LCEL
  3. Rozgałęzianie i równoległe łańcuchy
  4. Strumieniowanie danych wyjściowych w LangChain
← Powrót do AI Engineering Academy