Konsumere strømmer med Python SDK
Bruk OpenAI-asynkronklienten med async for til å konsumere strømmende fullføringer, samle hele svaret og håndtere feil underveis uten å miste delvis output.
Konsumere strømmer med Python SDK er en gratis leksjon i AI Engineering Academy på CoddyKit. Dette er leksjon 2 av 4. Du kan lese hele leksjonen gratis nedenfor – og deretter øve praktisk i nettleseren med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt. Den er en del av læringsløpet i AI Engineering Academy, og fremdriften din synkroniseres mellom nettet og CoddyKit-appen. Kurset i AI Engineering Academy inneholder totalt 4 leksjoner.
Synkrone og asynkrone strømmeklienter
OpenAI Python SDK tilbyr både en synkron OpenAI-klient og en asynkron AsyncOpenAI-klient. For kommandolinjeskript og enkle applikasjoner er den synkrone klienten enklere å bruke. For webservere, API-er og applikasjoner som håndterer flere samtidige forespørsler, er den asynkrone klienten avgjørende — den blokkerer ikke hendelsesløkken mens den venter på tokens, slik at andre forespørsler kan behandles samtidig.
# 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 awaitedAsynkron strømming med AsyncOpenAI
Med AsyncOpenAI-klienten blir strømmekallet en coroutine. De bruker async for til å iterere over chunks i stedet for en vanlig for-løkke. Hendelsesløkken kan planlegge andre coroutiner mellom hver chunk som kommer inn, slik at serveren kan håndtere andre forespørsler mens den venter på neste token fra LLM-en — dette er den viktigste fordelen fremfor synkron strømming i en webkontekst.
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'))Bruke stream-kontekstbehandleren
OpenAI SDK tilbyr også en stream-kontekstbehandler via client.chat.completions.stream(). Denne metoden lukker automatisk strømmen når konteksten avsluttes, og tilbyr praktiske metoder som stream.text_stream, som bare gir ikke-None-tekstendringer, og stream.get_final_completion() for bruksstatistikk etter strømmingen uten at De må samle dem manuelt.
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?'))Håndtere feil under strømming på en god måte
Feil kan oppstå når som helst under en strøm: under den første tilkoblingen, etter det første tokenet eller mot slutten av et langt svar. Pakk itereringen over strømmen inn i try/except-blokker, og håndter openai.APIConnectionError, openai.RateLimitError og openai.APIStatusError separat, siden hver av dem krever en ulik strategi for gjenoppretting (nytt forsøk, gradvis venting eller varsling av brukeren).
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__}]'Asynkron generator for strømming
Det ryddigste asynkrone mønsteret for strømming er en asynkron generatorfunksjon som gir fra seg tokens. Forbrukere itererer over den med async for. Dette holder strømmelogikken atskilt fra hvordan resultatet brukes — et FastAPI-endepunkt, en WebSocket-behandler og en test kan alle bruke den samme generatoren uten å kjenne til hverandre.
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)])Samtidige strømmingsforespørsler
En viktig fordel med asynkron strømming er muligheten til å kjøre flere strømmer samtidig i én enkelt prosess. Ved å bruke asyncio.gather kan De starte flere LLM-strømmeforespørsler samtidig og behandle tokenene etter hvert som de kommer inn. Dette er nyttig for fan-out-mønstre der De vil sammenligne flere promptvarianter eller kjøre parallelle deloppgaver.
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())Tidsavbrudd og kansellering
Langvarige strømmer bør ha tidsavbrudd for å hindre blokkering på ubestemt tid. Bruk asyncio.wait_for for å bruke et tidsavbrudd på coroutinenivå, eller httpx.Timeout for å angi tidsavbrudd for tilkobling og lesing på HTTP-klientnivå. Begge metodene sørger for at en strøm som har stoppet opp, ikke holder en forespørsel åpen på ubestemt tid. Kanseller alltid strømmer eksplisitt når brukeren kobler fra, slik at GPU-ressurser ikke sløses bort.
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)Bufre ufullstendige linjer
Når De strømmer til en klient som behandler komplette linjer (for eksempel en CLI som gjengir markdown), kan De ønske å bufre tokens frem til et linjeskift eller slutten på en setning før De videresender dem. Dette hindrer flimrende gjengivelse av ufullstendige setninger. Samle tokens i en buffer, tøm bufferen til forbrukeren når De oppdager tegnsetting som avslutter en setning, eller et linjeskift, og tøm alltid resten av bufferen når strømmen avsluttes.
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 bufferRegistrere strømmelatens i produksjon
I produksjon bør De instrumentere hver strøm for å registrere TTFT og total genereringstid til overvåking. Lagre disse målingene i en tidsseriedatabase, og varsle når TTFT overskrider SLA-terskelen (vanligvis 1–2 sekunder for interaktive applikasjoner). Sett TTFT-topper i sammenheng med promptlengde, modellbelastning og tidspunkt på dagen for å identifisere årsakene til forverret latenstid.
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, metricsTeste asynkron strømmekode
Testing av asynkron strømming krever særlig omtanke. Bruk pytest-asyncio til å kjøre asynkrone testfunksjoner, og mock OpenAI-klienten for å unngå ekte API-kall i enhetstester. Opprett en falsk strøm som gir fra seg forhåndsdefinerte chunks med konfigurerbare forsinkelser, slik at De kan teste både behandling av tokens i normaltilfellet og feilhåndtering uten å bruke API-kvoten.
# 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!'SDK-hjelpere: stream.text og stream.final_message
OpenAI Python SDKs stream-kontekstbehandler tilbyr hjelpeattributter som gjør manuell samling unødvendig. stream.text_stream er en asynkron itererbar som bare gir fra seg innholdsstrenger som ikke er None. Etter at strømmen er fullført, returnerer await stream.get_final_message() en fullstendig ChatCompletionMessage med hele teksten og bruksdata. Disse hjelperne reduserer mengden standardkode og håndterer automatisk kanttilfeller som tomme endringer.
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.contentKort sjekk
Test forståelsen Deres av asynkron strømming med OpenAI Python SDK fra denne leksjonen.
Oppsummering av leksjonen
I denne leksjonen har De lært: AsyncOpenAI muliggjør ikke-blokkerende strømming, slik at servere kan håndtere samtidige forespørsler, asynkrone generatorer er det ryddigste mønsteret for å gi fra seg strømmede tokens til nedstrømsforbrukere, og asyncio.wait_for og timeout-parametere hindrer strømmer som har stoppet opp på ubestemt tid, i å blokkere serveren. Stream-kontekstbehandleren tilbyr praktiske hjelpere som text_stream og get_final_completion. Deretter eksponerer vi LLM-strømming for nettleserklienter via FastAPI og Server-Sent Events.
Lær deg Python med en AI-veileder – gratis
Skriv og kjør ekte kode i nettleseren, få umiddelbar hjelp fra en AI-veileder som er tilgjengelig døgnet rundt, og fortsett der du slapp – på nettet eller i appen.
- Kurs
- 30
- Leksjoner
- 120
Ofte stilte spørsmål
Er leksjonen «Konsumere strømmer med Python SDK» gratis?
Ja – hele teksten i «Konsumere strømmer med Python SDK» er gratis å lese her på nettet. For å øve interaktivt med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt, og for å låse opp resten av AI Engineering Academy-kurset, kan du oppgradere til CoddyKit PRO. Kurset i AI Engineering Academy inneholder totalt 4 leksjoner.
Hva lærer jeg i «Konsumere strømmer med Python SDK»?
Bruk OpenAI-asynkronklienten med async for til å konsumere strømmende fullføringer, samle hele svaret og håndtere feil underveis uten å miste delvis output. Du øver på AI Engineering Academy med praktisk kode som du kjører direkte i nettleseren, mens en AI-veileder som er tilgjengelig døgnet rundt, svarer på spørsmålene dine mens du jobber deg gjennom leksjonen.
Trenger jeg erfaring for å begynne med AI Engineering Academy?
Ingen tidligere erfaring er nødvendig. AI Engineering Academy på CoddyKit er lagt opp for både nybegynnere og viderekomne, så De kan begynne her eller helt fra start og lære i Deres eget tempo. Dette er leksjon 2 av 4.
Hvor lang tid tar leksjonen «Konsumere strømmer med Python SDK»?
De fleste CoddyKit-leksjoner tar omtrent 5–10 minutter. Hver leksjon er kort og interaktiv, slik at De gjør jevne fremskritt og kan fortsette akkurat der De slapp – både på nettet og i appen.
Kan jeg skrive og kjøre kode i denne AI Engineering Academy-leksjonen?
Ja. Alle AI Engineering Academy-leksjoner har en innebygd kodeeditor, slik at De kan skrive og kjøre ekte kode direkte i nettleseren og få umiddelbar tilbakemelding fra AI – uten lokal konfigurering.
Alle leksjonene i dette kurset
- Forstå tokenstrømming
- Konsumere strømmer med Python SDK
- Strømming i FastAPI med Server-Sent Events
- Håndtere verktøykall i strømmende svar