0Pricing
AI Engineering Academy · Lekcja

Obsługa strumieni za pomocą Python SDK

Użyj klienta asynchronicznego OpenAI wraz z async for, aby odbierać strumieniowane uzupełnienia, złożyć pełną odpowiedź i obsługiwać błędy występujące w trakcie strumienia bez utraty częściowego wyniku.

Obsługa strumieni za pomocą Python SDK to bezpłatna lekcja AI Engineering Academy na CoddyKit. To lekcja 2 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.

Klienci synchroniczni a asynchroniczni do strumieniowania

OpenAI Python SDK udostępnia zarówno synchronicznego klienta OpenAI, jak i asynchronicznego klienta AsyncOpenAI. W przypadku skryptów wiersza poleceń i prostych aplikacji klient synchroniczny jest łatwiejszy w użyciu. W przypadku serwerów WWW, interfejsów API i aplikacji obsługujących wiele współbieżnych żądań klient async jest niezbędny — nie blokuje pętli zdarzeń podczas oczekiwania na tokeny, dzięki czemu inne żądania mogą być obsługiwane współbieżnie.

# Synchronous client (simple scripts)
from openai import OpenAI
client = OpenAI()

# Asynchronous client (web servers, concurrent workloads)
from openai import AsyncOpenAI
async_client = AsyncOpenAI()

# The async client has the same API surface as the sync client
# but all methods are coroutines that must be awaited

Asynchroniczne strumieniowanie za pomocą AsyncOpenAI

W przypadku klienta AsyncOpenAI wywołanie strumieniowania staje się korutyną. Do iterowania po fragmentach służy async for, a nie zwykła pętla for. Pętla zdarzeń może planować wykonywanie innych korutyn między nadejściem kolejnych fragmentów, dzięki czemu serwer może obsługiwać inne żądania podczas oczekiwania na następny token z LLM — to kluczowa przewaga nad synchronicznym strumieniowaniem w środowisku WWW.

import asyncio
from openai import AsyncOpenAI

async_client = AsyncOpenAI()

async def async_stream_completion(prompt: str) -> str:
    stream = await async_client.chat.completions.create(
        model='gpt-4o-mini',
        messages=[{'role': 'user', 'content': prompt}],
        stream=True,
    )

    full_text = ''
    async for chunk in stream:
        delta = chunk.choices[0].delta.content
        if delta:
            print(delta, end='', flush=True)
            full_text += delta
    print()
    return full_text

# Run the coroutine
asyncio.run(async_stream_completion('Explain what async/await does in Python'))

Korzystanie z menedżera kontekstu strumienia

OpenAI SDK udostępnia również menedżera kontekstu strumienia za pośrednictwem client.chat.completions.stream(). To podejście automatycznie zamyka strumień po wyjściu z kontekstu i udostępnia pomocnicze metody, takie jak stream.text_stream, która zwraca wyłącznie niepuste delty tekstowe, oraz stream.get_final_completion(), która umożliwia uzyskanie statystyk użycia po zakończeniu strumienia bez konieczności ręcznego ich gromadzenia.

from openai import AsyncOpenAI
import asyncio

async def stream_with_context_manager(prompt: str):
    async with async_client.chat.completions.stream(
        model='gpt-4o-mini',
        messages=[{'role': 'user', 'content': prompt}],
    ) as stream:
        # text_stream filters None deltas automatically
        async for text in stream.text_stream:
            print(text, end='', flush=True)

    # Access final completion after stream ends
    completion = await stream.get_final_completion()
    print(f'\nUsage: {completion.usage}')
    return completion

asyncio.run(stream_with_context_manager('What are the benefits of async I/O?'))

Elegancka obsługa błędów występujących w trakcie strumieniowania

Błędy mogą wystąpić w dowolnym momencie działania strumienia: podczas nawiązywania połączenia, po pierwszym tokenie lub pod koniec długiej odpowiedzi. Iterację po strumieniu należy opakować w bloki try/except i osobno obsłużyć openai.APIConnectionError, openai.RateLimitError oraz openai.APIStatusError, ponieważ każdy z nich wymaga innej strategii odzyskiwania — ponowienia, narastającego opóźnienia lub powiadomienia użytkownika.

import openai

async def resilient_stream(prompt: str):
    try:
        stream = await async_client.chat.completions.create(
            model='gpt-4o-mini',
            messages=[{'role': 'user', 'content': prompt}],
            stream=True,
        )
        accumulated = ''
        async for chunk in stream:
            delta = chunk.choices[0].delta.content
            if delta:
                accumulated += delta
                yield delta  # async generator

    except openai.RateLimitError:
        yield '[Rate limit reached — please wait and retry]'
    except openai.APIConnectionError:
        yield '[Connection error — check your network]'
    except openai.APIStatusError as e:
        yield f'[API error {e.status_code}]'
    except Exception as e:
        yield f'[Unexpected error: {type(e).__name__}]'

Generator asynchroniczny do strumieniowania

Najbardziej przejrzystym asynchronicznym wzorcem strumieniowania jest asynchroniczna funkcja generatora, która zwraca tokeny za pomocą yield. Odbiorcy iterują po niej przy użyciu async for. Dzięki temu logika strumieniowania jest oddzielona od sposobu wykorzystania danych wyjściowych — endpoint FastAPI, obsługa WebSocketu i test mogą korzystać z tego samego generatora, nie wiedząc o swoim istnieniu.

from typing import AsyncGenerator

async def token_stream(
    messages: list[dict],
    model: str = 'gpt-4o-mini',
) -> AsyncGenerator[str, None]:
    stream = await async_client.chat.completions.create(
        model=model,
        messages=messages,
        stream=True,
    )
    async for chunk in stream:
        delta = chunk.choices[0].delta.content
        if delta:
            yield delta

# Consumer 1: print to terminal
async def print_stream(messages):
    async for token in token_stream(messages):
        print(token, end='', flush=True)

# Consumer 2: collect to string
async def collect_stream(messages) -> str:
    return ''.join([t async for t in token_stream(messages)])

Współbieżne żądania strumieniowania

Jedną z głównych zalet asynchronicznego strumieniowania jest możliwość uruchamiania wielu strumieni współbieżnie w ramach jednego procesu. Za pomocą asyncio.gather można jednocześnie uruchomić kilka żądań strumieniowania do LLM i przetwarzać ich tokeny w miarę nadejścia. Jest to przydatne we wzorcach fan-out, gdy chcą Państwo porównać kilka wariantów promptu lub uruchomić równoległe podzadania.

import asyncio

async def run_parallel_streams(queries: list[str]) -> list[str]:
    async def collect(query):
        messages = [{'role': 'user', 'content': query}]
        return ''.join([t async for t in token_stream(messages)])

    results = await asyncio.gather(*[collect(q) for q in queries])
    return results

queries = [
    'What is RAG?',
    'What is a vector database?',
    'What is BM25?',
]

async def main():
    answers = await run_parallel_streams(queries)
    for q, a in zip(queries, answers):
        print(f'Q: {q}\nA: {a[:100]}\n')

asyncio.run(main())

Limity czasu i anulowanie

Długotrwałe strumienie powinny mieć ustawione limity czasu, aby zapobiec blokowaniu bez określonego końca. Użyj asyncio.wait_for, aby zastosować limit czasu na poziomie korutyny, albo httpx.Timeout, aby ustawić limity czasu połączenia i odczytu na poziomie klienta HTTP. Oba podejścia gwarantują, że zatrzymany strumień nie będzie bezterminowo zajmował żądania. Po rozłączeniu użytkownika zawsze należy jawnie anulować strumienie, aby uniknąć marnowania mocy obliczeniowej GPU.

import asyncio

async def stream_with_timeout(messages, timeout_seconds: float = 30.0):
    try:
        async with asyncio.timeout(timeout_seconds):
            stream = await async_client.chat.completions.create(
                model='gpt-4o-mini',
                messages=messages,
                stream=True,
                timeout=timeout_seconds,  # HTTP-level timeout
            )
            async for chunk in stream:
                delta = chunk.choices[0].delta.content
                if delta:
                    yield delta
    except asyncio.TimeoutError:
        yield '\n[Stream timed out after {:.0f}s]'.format(timeout_seconds)

Buforowanie niepełnych wierszy

Podczas strumieniowania do klienta, który przetwarza kompletne wiersze, na przykład interfejsu CLI renderującego Markdown, można buforować tokeny do napotkania znaku nowego wiersza lub granicy zdania, zanim zostaną przekazane dalej. Zapobiega to migotaniu renderowanych niepełnych zdań. Należy gromadzić tokeny w buforze, opróżniać go po wykryciu znaku kończącego zdanie lub znaku nowego wiersza, a na końcu strumienia zawsze opróżnić pozostałą zawartość bufora.

async def buffered_line_stream(messages):
    buffer = ''
    flush_on = {'.', '!', '?', '\n'}

    async for token in token_stream(messages):
        buffer += token
        if any(c in buffer for c in flush_on):
            # Find the last sentence-ending position
            for i, c in enumerate(reversed(buffer)):
                if c in flush_on:
                    split_pos = len(buffer) - i
                    yield buffer[:split_pos]
                    buffer = buffer[split_pos:]
                    break

    if buffer:  # flush remainder
        yield buffer

Rejestrowanie opóźnienia strumienia w środowisku produkcyjnym

W środowisku produkcyjnym należy instrumentować każdy strumień, aby rejestrować TTFT i całkowity czas generowania na potrzeby monitorowania. Metryki te należy przechowywać w bazie danych szeregów czasowych i wysyłać alerty, gdy TTFT przekroczy próg SLA, zwykle 1–2 sekundy w aplikacjach interaktywnych. Warto korelować skoki TTFT z długością promptu, obciążeniem modelu i porą dnia, aby identyfikować główne przyczyny pogorszenia opóźnień.

import time
from dataclasses import dataclass

@dataclass
class StreamMetrics:
    prompt_chars: int
    ttft_ms: float
    total_ms: float
    token_count: int

async def instrumented_stream(messages) -> tuple[str, StreamMetrics]:
    t_start = time.perf_counter()
    t_first = None
    token_count = 0
    full_text = ''

    stream = await async_client.chat.completions.create(
        model='gpt-4o-mini', messages=messages, stream=True
    )
    async for chunk in stream:
        delta = chunk.choices[0].delta.content
        if delta:
            if t_first is None:
                t_first = time.perf_counter()
            token_count += 1
            full_text += delta

    t_end = time.perf_counter()
    prompt_len = sum(len(m.get('content', '')) for m in messages)
    metrics = StreamMetrics(
        prompt_chars=prompt_len,
        ttft_ms=(t_first - t_start) * 1000 if t_first else 0,
        total_ms=(t_end - t_start) * 1000,
        token_count=token_count,
    )
    return full_text, metrics

Testowanie kodu asynchronicznego strumieniowania

Testowanie asynchronicznego strumieniowania wymaga szczególnej uwagi. Należy użyć pytest-asyncio do uruchamiania asynchronicznych funkcji testowych oraz zamockować klienta OpenAI, aby w testach jednostkowych uniknąć rzeczywistych wywołań API. Należy utworzyć sztuczny strumień zwracający zdefiniowane wcześniej fragmenty z konfigurowalnymi opóźnieniami, aby testować zarówno przetwarzanie tokenów w poprawnym scenariuszu, jak i ścieżki obsługi błędów bez zużywania limitu API.

# pip install pytest pytest-asyncio
import pytest
from unittest.mock import AsyncMock, MagicMock

async def fake_stream(tokens: list[str]):
    for token in tokens:
        chunk = MagicMock()
        chunk.choices[0].delta.content = token
        yield chunk

@pytest.mark.asyncio
async def test_stream_accumulates_correctly(monkeypatch):
    mock_create = AsyncMock(return_value=fake_stream(['Hello', ', ', 'world', '!']))
    monkeypatch.setattr(async_client.chat.completions, 'create', mock_create)

    result = await collect_stream([{'role': 'user', 'content': 'Hi'}])
    assert result == 'Hello, world!'

Pomocnicze elementy SDK: stream.text i stream.final_message

Menedżer kontekstu strumienia w OpenAI Python SDK udostępnia pomocnicze atrybuty, które eliminują konieczność ręcznego gromadzenia danych. stream.text_stream to asynchroniczny iterowalny obiekt zwracający wyłącznie niepuste ciągi treści. Po zakończeniu strumienia await stream.get_final_message() zwraca pełny obiekt ChatCompletionMessage z kompletnym tekstem i danymi o użyciu. Te pomocnicze elementy ograniczają ilość kodu dodatkowego i automatycznie obsługują przypadki brzegowe, takie jak puste delty.

async def clean_streaming_example(prompt: str):
    async with async_client.chat.completions.stream(
        model='gpt-4o-mini',
        messages=[{'role': 'user', 'content': prompt}],
    ) as stream:
        # Iterate only over text tokens, None deltas filtered automatically
        async for text in stream.text_stream:
            print(text, end='', flush=True)

    # After context exit, get accumulated result
    final = await stream.get_final_completion()
    return final.choices[0].message.content

Szybki test

Sprawdź swoją wiedzę na temat asynchronicznego strumieniowania za pomocą OpenAI Python SDK omówionego w tej lekcji.

Podsumowanie lekcji

W tej lekcji nauczyli się Państwo, że: AsyncOpenAI umożliwia nieblokujące strumieniowanie, dzięki któremu serwery mogą obsługiwać współbieżne żądania; asynchroniczne generatory to najbardziej przejrzysty wzorzec zwracania tokenów strumienia odbiorcom; a asyncio.wait_for i parametry timeout zapobiegają blokowaniu serwera przez strumienie, które zatrzymały się na czas nieokreślony. Menedżer kontekstu strumienia udostępnia pomocnicze elementy, takie jak text_stream i get_final_completion. Następnie udostępnimy strumieniowanie LLM klientom przeglądarkowym za pomocą FastAPI i Server-Sent Events.

Często zadawane pytania

Czy lekcja „Obsługa strumieni za pomocą Python SDK” jest bezpłatna?

Tak — pełny tekst „Obsługa strumieni za pomocą Python SDK” 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 „Obsługa strumieni za pomocą Python SDK”?

Użyj klienta asynchronicznego OpenAI wraz z async for, aby odbierać strumieniowane uzupełnienia, złożyć pełną odpowiedź i obsługiwać błędy występujące w trakcie strumienia bez utraty częściowego wyni… Ć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 2 z 4.

Ile czasu zajmuje lekcja „Obsługa strumieni za pomocą Python SDK”?

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. Zrozumienie strumieniowania tokenów
  2. Obsługa strumieni za pomocą Python SDK
  3. Streaming w FastAPI z Server-Sent Events
  4. Obsługa wywołań narzędzi w strumieniowanych odpowiedziach
← Powrót do AI Engineering Academy