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 awaitedAsynchroniczne 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 bufferRejestrowanie 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, metricsTestowanie 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.contentSzybki 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
- Zrozumienie strumieniowania tokenów
- Obsługa strumieni za pomocą Python SDK
- Streaming w FastAPI z Server-Sent Events
- Obsługa wywołań narzędzi w strumieniowanych odpowiedziach