AI Engineering Academy · Les

Streams verwerken met de Python SDK

Gebruik de asynchrone client van OpenAI met async for om streaming-completions te verwerken, het volledige antwoord op te bouwen en fouten tijdens de stream af te handelen zonder gedeeltelijke uitvoer te verliezen.

Les 2 van 413 stappen

Streams verwerken met de Python SDK is een gratis AI Engineering Academy-les op CoddyKit. Dit is les 2 van 4. Je kunt de volledige les hieronder gratis lezen en daarna in de browser praktisch oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is. Deze les maakt deel uit van het leertraject AI Engineering Academy. Je voortgang wordt gesynchroniseerd op het web en in de CoddyKit-app. De cursus AI Engineering Academy bevat in totaal 4 lessen.

Synchrone versus asynchrone streamingclients

De OpenAI Python SDK biedt zowel een synchrone OpenAI-client als een asynchrone AsyncOpenAI-client. Voor opdrachtregelscripts en eenvoudige toepassingen is de synchrone client gemakkelijker te gebruiken. Voor webservers, API's en toepassingen die meerdere gelijktijdige verzoeken verwerken, is de asynchrone client essentieel — deze blokkeert de eventlus niet tijdens het wachten op tokens, zodat andere verzoeken gelijktijdig kunnen worden afgehandeld.

# 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

Asynchrone streaming met AsyncOpenAI

Met de AsyncOpenAI-client wordt de streamingaanroep een coroutine. Je gebruikt async for om over brokken te itereren in plaats van een gewone for-lus. De eventlus kan tussen de aankomst van elk brok andere coroutines inplannen, waardoor je server andere verzoeken kan verwerken terwijl deze wacht op het volgende token van de LLM — dit is het belangrijkste voordeel ten opzichte van synchrone streaming in een webcontext.

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'))

De streamcontextmanager gebruiken

De OpenAI SDK biedt ook een streamcontextmanager via client.chat.completions.stream(). Deze aanpak sluit de stream automatisch wanneer de context wordt verlaten en biedt handige methoden zoals stream.text_stream, die alleen tekstverschillen die niet None zijn oplevert, en stream.get_final_completion() voor gebruiksstatistieken na de stream, zonder dat je deze handmatig hoeft te verzamelen.

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?'))

Fouten halverwege een stream netjes afhandelen

Fouten kunnen op elk moment tijdens een stream optreden: tijdens de eerste verbinding, na het eerste token of tegen het einde van een lang antwoord. Plaats de iteratie over je stream in try/except-blokken en handel openai.APIConnectionError, openai.RateLimitError en openai.APIStatusError afzonderlijk af, omdat voor elk een andere herstelstrategie nodig is (opnieuw proberen, een wachttijd gebruiken of de gebruiker informeren).

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__}]'

Asynchrone generator voor streaming

Het duidelijkste asynchrone patroon voor streaming is een asynchrone generatorfunctie die tokens oplevert. Gebruikers itereren hierover met async for. Zo blijft de streaminglogica gescheiden van de manier waarop de uitvoer wordt gebruikt — een FastAPI-eindpunt, een WebSocket-handler en een test gebruiken allemaal dezelfde generator zonder iets van elkaar te weten.

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)])

Gelijktijdige streamingverzoeken

Een groot voordeel van asynchrone streaming is dat je meerdere streams gelijktijdig binnen één proces kunt uitvoeren. Met asyncio.gather kun je meerdere streamingverzoeken aan een LLM tegelijkertijd starten en hun tokens verwerken zodra ze binnenkomen. Dit is nuttig voor fan-outpatronen waarbij je meerdere varianten van een prompt wilt vergelijken of parallelle subtaken wilt uitvoeren.

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())

Time-outs en annulering

Langdurige streams moeten time-outs hebben om onbeperkt blokkeren te voorkomen. Gebruik asyncio.wait_for om een time-out op coroutineniveau toe te passen, of httpx.Timeout om verbindings- en leestime-outs in te stellen op het niveau van de HTTP-client. Beide aanpakken zorgen ervoor dat een vastgelopen stream een verzoek niet onbeperkt vasthoudt. Annuleer streams altijd expliciet wanneer de gebruiker de verbinding verbreekt, zodat er geen GPU-rekenkracht wordt verspild.

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)

Gedeeltelijke regels bufferen

Wanneer je streamt naar een client die volledige regels verwerkt (zoals een CLI die markdown weergeeft), wil je tokens mogelijk bufferen tot een nieuwe regel of zinseinde voordat je ze doorstuurt. Zo voorkom je flikkerende weergaven van onvolledige zinnen. Verzamel tokens in een buffer, stuur de buffer door naar de verbruiker zodra je leestekens aan het einde van een zin of een nieuweregelteken detecteert, en stuur de resterende buffer altijd door aan het einde van de 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 buffer

Streamlatentie in productie registreren

Instrumenteer in productie elke stream om TTFT en de totale generatietijd vast te leggen voor monitoring. Sla deze metingen op in een tijdreeksdatabase en geef een waarschuwing wanneer TTFT je SLA-drempel overschrijdt (doorgaans 1-2 seconden voor interactieve toepassingen). Breng pieken in TTFT in verband met promptlengte, modelbelasting en het tijdstip van de dag om de grondoorzaken van verslechterde latentie te identificeren.

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

Asynchrone streamingcode testen

Voor het testen van asynchrone streaming is extra zorg nodig. Gebruik pytest-asyncio om asynchrone testfuncties uit te voeren en simuleer de OpenAI-client om echte API-aanroepen in unittests te vermijden. Maak een nepstream die vooraf gedefinieerde brokken met instelbare vertragingen oplevert, zodat je zowel de verwerking van tokens bij een succesvol verloop als foutpaden kunt testen zonder API-budget te besteden.

# 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-hulpmiddelen: stream.text en stream.final_message

De streamcontextmanager van de OpenAI Python SDK biedt hulpkenmerken waarmee handmatig verzamelen niet nodig is. stream.text_stream is een asynchrone iterable die alleen inhoudstekenreeksen oplevert die niet None zijn. Nadat de stream is voltooid, retourneert await stream.get_final_message() een volledige ChatCompletionMessage met de volledige tekst en gebruiksgegevens. Deze hulpmiddelen verminderen herhalende code en handelen randgevallen zoals lege verschillen automatisch af.

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

Korte controle

Test in deze les je begrip van asynchrone streaming met de OpenAI Python SDK.

Samenvatting van de les

In deze les heb je geleerd: AsyncOpenAI maakt niet-blokkerende streaming mogelijk, zodat servers gelijktijdige verzoeken kunnen verwerken, asynchrone generators het duidelijkste patroon zijn om streamingtokens op te leveren aan downstreamgebruikers, en asyncio.wait_for en time-outparameters voorkomen dat onbeperkt vastgelopen streams je server blokkeren. De streamcontextmanager biedt handige hulpmiddelen zoals text_stream en get_final_completion. Hierna maken we LLM-streaming beschikbaar voor browserclients via FastAPI en Server-Sent Events.

Gratis beginnen

Leer Python met een AI-tutor — gratis

Schrijf echte code en voer die uit in je browser, krijg direct hulp van een AI-tutor die 24/7 beschikbaar is en ga verder waar je gebleven bent op het web of in de app.

Cursussen
30
Lessen
120

Veelgestelde vragen

Is de les “Streams verwerken met de Python SDK” gratis?

Ja — de volledige tekst van “Streams verwerken met de Python SDK” kun je hier gratis op het web lezen. Als je interactief wilt oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is, en de rest van de cursus AI Engineering Academy wilt ontgrendelen, kun je upgraden naar CoddyKit PRO. De cursus AI Engineering Academy bevat in totaal 4 lessen.

Wat leer ik in “Streams verwerken met de Python SDK”?

Gebruik de asynchrone client van OpenAI met async for om streaming-completions te verwerken, het volledige antwoord op te bouwen en fouten tijdens de stream af te handelen zonder gedeeltelijke uitvoe… Je oefent met AI Engineering Academy door code rechtstreeks in de browser uit te voeren. Een AI-begeleider die 24/7 beschikbaar is beantwoordt je vragen terwijl je de les doorwerkt.

Heb ik ervaring nodig om met AI Engineering Academy te beginnen?

Ervaring vooraf is niet nodig. AI Engineering Academy op CoddyKit is opgebouwd voor beginners tot gevorderden, zodat je hier of bij het begin kunt starten en in je eigen tempo kunt leren. Dit is les 2 van 4.

Hoe lang duurt de les “Streams verwerken met de Python SDK”?

De meeste lessen van CoddyKit duren ongeveer 5–10 minuten. Elke les is kort en interactief, zodat je gestaag vooruitgaat en op het web en in de app precies verdergaat waar je was gebleven.

Kan ik code schrijven en uitvoeren in deze les over AI Engineering Academy?

Ja. Elke les over AI Engineering Academy bevat een ingebouwde code-editor, zodat je rechtstreeks in je browser echte code kunt schrijven en uitvoeren en direct feedback van AI krijgt — lokale installatie is niet nodig.

Alle lessen in deze cursus

  1. Tokenstreaming begrijpen
  2. Streams verwerken met de Python SDK
  3. Streaming in FastAPI met Server-Sent Events
  4. Toolaanroepen in gestreamde antwoorden verwerken
← Terug naar AI Engineering Academy