Consumare gli stream con l'SDK Python
Utilizzi il client asincrono OpenAI con async for per consumare i completamenti in streaming, accumulare la risposta completa e gestire gli errori durante lo stream senza perdere l'output parziale.
Consumare gli stream con l'SDK Python è una lezione AI Engineering Academy gratuita su CoddyKit. Questa è la lezione 2 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento AI Engineering Academy, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso AI Engineering Academy include 4 lezioni in totale.
Client di streaming sincroni e asincroni
L'SDK Python di OpenAI fornisce sia un client sincrono OpenAI sia un client asincrono AsyncOpenAI. Per gli script da riga di comando e le applicazioni semplici, il client sincrono è più facile da usare. Per i server web, le API e le applicazioni che gestiscono più richieste simultanee, il client asincrono è essenziale: non blocca il loop degli eventi mentre attende i token, consentendo di gestire altre richieste contemporaneamente.
# 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 awaitedStreaming asincrono con AsyncOpenAI
Con il client AsyncOpenAI, la chiamata di streaming diventa una coroutine. Utilizzi async for per iterare sui chunk invece di un normale ciclo for. Il loop degli eventi può pianificare l'esecuzione di altre coroutine tra l'arrivo di un chunk e quello successivo, consentendo al server di gestire altre richieste mentre attende il token successivo dall'LLM: questo è il vantaggio principale rispetto allo streaming sincrono in un contesto web.
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'))Utilizzo del gestore di contesto per lo stream
L'SDK OpenAI fornisce anche un gestore di contesto per lo stream tramite client.chat.completions.stream(). Questo approccio chiude automaticamente lo stream all'uscita dal contesto e offre metodi pratici come stream.text_stream, che restituisce solo i delta di testo non-None, e stream.get_final_completion() per ottenere le statistiche di utilizzo al termine dello stream senza doverle accumulare manualmente.
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?'))Gestione corretta degli errori durante lo stream
Gli errori possono verificarsi in qualsiasi momento durante uno stream: durante la connessione iniziale, dopo il primo token o verso la fine di una risposta lunga. Inserisca l'iterazione sullo stream in blocchi try/except e gestisca separatamente openai.APIConnectionError, openai.RateLimitError e openai.APIStatusError, poiché ciascuno richiede una strategia di recupero diversa (nuovo tentativo, backoff o notifica all'utente).
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__}]'Generatore asincrono per lo streaming
Il pattern asincrono più chiaro per lo streaming consiste in una funzione generatore asincrona che restituisce i token con yield. I consumer vi iterano utilizzando async for. In questo modo la logica di streaming rimane separata dal modo in cui l'output viene utilizzato: un endpoint FastAPI, un gestore WebSocket e un test possono consumare lo stesso generatore senza conoscersi tra loro.
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)])Richieste di streaming simultanee
Un grande vantaggio dello streaming asincrono è la possibilità di eseguire più stream simultaneamente all'interno di un singolo processo. Utilizzando asyncio.gather, può avviare contemporaneamente diverse richieste di streaming all'LLM ed elaborarne i token man mano che arrivano. Questo è utile nei pattern fan-out, quando desidera confrontare più varianti di un prompt o eseguire sotto-attività in parallelo.
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())Timeout e annullamento
Gli stream di lunga durata dovrebbero avere dei timeout per evitare blocchi indefiniti. Utilizzi asyncio.wait_for per applicare un timeout a livello di coroutine oppure httpx.Timeout per impostare timeout di connessione e lettura a livello del client HTTP. Entrambi gli approcci assicurano che uno stream bloccato non mantenga una richiesta aperta indefinitamente. Annulli sempre esplicitamente gli stream quando l'utente si disconnette, per evitare di sprecare capacità di calcolo della 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)Bufferizzazione delle righe parziali
Quando esegue lo streaming verso un client che elabora righe complete, come una CLI che visualizza Markdown, potrebbe voler accumulare i token fino a una nuova riga o alla fine di una frase prima di inoltrarli. Questo evita lo sfarfallio causato dalla visualizzazione di frasi parziali. Accumuli i token in un buffer, lo invii al consumer quando rileva un segno di punteggiatura che conclude una frase o un carattere di nuova riga e svuoti sempre il buffer rimanente al termine dello stream.
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 bufferRegistrazione della latenza dello stream in produzione
In produzione, strumentalizzi ogni stream per registrare il TTFT e il tempo totale di generazione ai fini del monitoraggio. Memorizzi queste metriche in un database di serie temporali e configuri un avviso quando il TTFT supera la soglia SLA (in genere 1-2 secondi per le applicazioni interattive). Metta in correlazione i picchi di TTFT con la lunghezza del prompt, il carico del modello e l'ora del giorno per individuare le cause principali del peggioramento della latenza.
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, metricsTest del codice di streaming asincrono
Il test dello streaming asincrono richiede particolare attenzione. Utilizzi pytest-asyncio per eseguire funzioni di test asincrone e simuli il client OpenAI per evitare chiamate API reali nei test unitari. Crei uno stream fittizio che restituisca chunk predefiniti con ritardi configurabili, così da testare sia l'elaborazione corretta dei token sia i percorsi di gestione degli errori senza consumare il budget 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!'Helper dell'SDK: stream.text e stream.final_message
Il gestore di contesto per lo stream dell'SDK Python di OpenAI fornisce attributi di supporto che evitano l'accumulo manuale. stream.text_stream è un iterabile asincrono che restituisce solo stringhe di contenuto non-None. Al termine dello stream, await stream.get_final_message() restituisce un ChatCompletionMessage completo, con il testo completo e i dati di utilizzo. Questi helper riducono il codice ripetitivo e gestiscono automaticamente casi limite come i delta vuoti.
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.contentVerifica rapida
Verifichi la sua comprensione dello streaming asincrono con l'SDK Python di OpenAI trattato in questa lezione.
Riepilogo della lezione
In questa lezione ha imparato che AsyncOpenAI abilita uno streaming non bloccante che consente ai server di gestire richieste simultanee, che i generatori asincroni sono il pattern più chiaro per fornire token in streaming ai consumer a valle e che asyncio.wait_for e i parametri di timeout impediscono agli stream bloccati indefinitamente di arrestare il server. Il gestore di contesto dello stream fornisce helper pratici come text_stream e get_final_completion. Nel prossimo argomento esporremo lo streaming LLM ai client browser tramite FastAPI e Server-Sent Events.
Domande Frequenti
La lezione «Consumare gli stream con l'SDK Python» è gratuita?
Sì — il testo completo di «Consumare gli stream con l'SDK Python» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso AI Engineering Academy, passa a CoddyKit PRO. Il corso AI Engineering Academy include 4 lezioni in totale.
Cosa imparerò in «Consumare gli stream con l'SDK Python»?
Utilizzi il client asincrono OpenAI con async for per consumare i completamenti in streaming, accumulare la risposta completa e gestire gli errori durante lo stream senza perdere l'output parziale. Eserciti AI Engineering Academy con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.
Ho bisogno di esperienza per iniziare AI Engineering Academy?
Non è richiesta alcuna esperienza precedente. AI Engineering Academy su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 2 di 4.
Quanto tempo richiede la lezione «Consumare gli stream con l'SDK Python»?
La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.
Posso scrivere ed eseguire codice in questa lezione AI Engineering Academy?
Sì. Ogni lezione AI Engineering Academy include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.
Tutte le lezioni di questo corso
- Comprendere lo streaming dei token
- Consumare gli stream con l'SDK Python
- Streaming in FastAPI con Server-Sent Events
- Gestire le chiamate agli strumenti nelle risposte in streaming