Aller au contenu principal

Implémenter le streaming en Python

Mis à jour le 29 juillet 2026

Du prototype à la production

La leçon précédente a posé les concepts du streaming SSE. Reste le plus long : transformer une boucle de démonstration en client capable de tenir sous charge, de survivre à une erreur réseau et d’alimenter une interface web. C’est l’objet de cette leçon, qui suit le trajet réel d’un projet — le squelette, puis les erreurs, puis l’asynchrone, puis l’exposition HTTP, et enfin la mesure.

Le pattern de base

Tout commence par ce squelette d’une dizaine de lignes, que vous adapterez à chaque projet. Il fait deux choses en même temps : afficher au fil de l’eau et accumuler pour le retour de fonction.

from mistralai import Mistral
import os

client = Mistral(api_key=os.getenv("MISTRAL_API_KEY"))

def stream_chat(messages: list, model: str = "mistral-large-latest") -> str:
    """Envoie une requête en streaming et affiche les tokens en temps réel."""
    stream = client.chat.stream(
        model=model,
        messages=messages
    )

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

    print()  # Retour à la ligne
    return full_text

Cette version convient à un script local. Elle a un défaut rédhibitoire en production : elle imprime dans le terminal, donc elle ne sait rien faire d’autre que du terminal.

Gestion complète des erreurs

Les erreurs réseau, les timeouts et les limites de débit ne sont pas des hypothèses : ce sont des certitudes à l’échelle du trafic réel. Le client ci-dessous répond aux deux problèmes à la fois. Il retente automatiquement les erreurs 429 avec un délai croissant, et il n’imprime plus rien : il émet des événements typés via yield, à charge pour l’appelant de décider ce qu’il en fait — les afficher, les pousser dans une WebSocket, les écrire en base.

from mistralai import Mistral
from mistralai.models import SDKError
import os
import time

client = Mistral(api_key=os.getenv("MISTRAL_API_KEY"))

def stream_with_retry(
    messages: list,
    model: str = "mistral-large-latest",
    max_retries: int = 3,
    max_tokens: int = 2000
) -> dict:
    """Streaming avec retry automatique et métriques."""
    for attempt in range(max_retries):
        try:
            stream = client.chat.stream(
                model=model,
                messages=messages,
                max_tokens=max_tokens
            )

            full_text = ""
            finish_reason = None
            chunk_count = 0

            for chunk in stream:
                delta = chunk.data.choices[0].delta
                if delta.content:
                    full_text += delta.content
                    chunk_count += 1
                    yield {"type": "token", "content": delta.content}

                if chunk.data.choices[0].finish_reason:
                    finish_reason = chunk.data.choices[0].finish_reason

            yield {
                "type": "done",
                "full_text": full_text,
                "finish_reason": finish_reason,
                "chunks": chunk_count
            }
            return

        except SDKError as e:
            if e.status_code == 429 and attempt < max_retries - 1:
                wait_time = 2 ** attempt  # Backoff exponentiel
                yield {"type": "retry", "wait": wait_time, "attempt": attempt + 1}
                time.sleep(wait_time)
            else:
                yield {"type": "error", "message": str(e), "status": e.status_code}
                return

        except Exception as e:
            yield {"type": "error", "message": str(e), "status": None}
            return

Le backoff exponentiel 2 ** attempt fait attendre 1 seconde, puis 2, puis 4. Cette progression n’est pas décorative : si tous vos clients retentaient au même rythme fixe, ils frapperaient l’API en rafales synchronisées et prolongeraient eux-mêmes la saturation qu’ils subissent. Notez aussi que seule l’erreur 429 déclenche une reprise — retenter une 401 due à une clé invalide ne ferait que retarder le diagnostic.

Côté appelant, la consommation devient une simple lecture d’événements :

messages = [
    {"role": "system", "content": "Vous êtes un assistant technique."},
    {"role": "user", "content": "Expliquez les design patterns en Python."}
]

for event in stream_with_retry(messages):
    if event["type"] == "token":
        print(event["content"], end="", flush=True)
    elif event["type"] == "retry":
        print(f"\n[Retry {event['attempt']}, attente {event['wait']}s...]")
    elif event["type"] == "done":
        print(f"\n\n--- Terminé ({event['chunks']} chunks, {event['finish_reason']})")
    elif event["type"] == "error":
        print(f"\n[Erreur : {event['message']}]")

Streaming asynchrone avec asyncio

Le client synchrone bloque son thread pendant toute la génération. Sur un serveur web qui doit tenir cinquante conversations simultanées, cela revient à immobiliser cinquante threads pour attendre. Le client asynchrone lève cette contrainte :

from mistralai import Mistral
import asyncio
import os

client = Mistral(api_key=os.getenv("MISTRAL_API_KEY"))

async def async_stream_chat(messages: list) -> str:
    """Streaming asynchrone pour les applications concurrentes."""
    stream = await client.chat.stream_async(
        model="mistral-large-latest",
        messages=messages
    )

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

    print()
    return full_text

# Exécution
async def main():
    messages = [{"role": "user", "content": "Bonjour, présentez-vous."}]
    result = await async_stream_chat(messages)
    print(f"\nTotal : {len(result)} caractères")

asyncio.run(main())

La logique est identique à la version synchrone ; seuls changent stream_async, le await sur l’appel et le async for sur la boucle.

Intégration avec FastAPI

Reste à faire sortir ce flux de votre processus Python pour l’amener jusqu’au navigateur. FastAPI s’en charge avec StreamingResponse, à condition de reformater chaque fragment au format SSE vu à la leçon précédente :

from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from mistralai import Mistral
import os
import json

app = FastAPI()
client = Mistral(api_key=os.getenv("MISTRAL_API_KEY"))

async def generate_stream(prompt: str):
    """Générateur SSE pour FastAPI."""
    stream = await client.chat.stream_async(
        model="mistral-large-latest",
        messages=[{"role": "user", "content": prompt}]
    )

    async for chunk in stream:
        content = chunk.data.choices[0].delta.content
        if content:
            data = json.dumps({"content": content})
            yield f"data: {data}\n\n"

    yield "data: [DONE]\n\n"

@app.get("/chat/stream")
async def chat_stream(prompt: str):
    return StreamingResponse(
        generate_stream(prompt),
        media_type="text/event-stream"
    )

Deux détails conditionnent le bon fonctionnement de ce endpoint. Le double retour à la ligne \n\n qui termine chaque yield est ce qui sépare les événements dans le protocole SSE : sans lui, le navigateur attend indéfiniment un événement qu’il considère inachevé. Et le media_type="text/event-stream" est ce qui autorise le client JavaScript à consommer la réponse avec un EventSource plutôt qu’avec un fetch classique.

Mesurer la performance

Une fois le flux en place, vous voudrez savoir ce qu’il vaut. Deux métriques suffisent à piloter l’optimisation :

import time

start = time.perf_counter()
first_token_time = None
token_count = 0

stream = client.chat.stream(
    model="mistral-large-latest",
    messages=[{"role": "user", "content": "Écrivez un paragraphe sur Paris."}]
)

for chunk in stream:
    if chunk.data.choices[0].delta.content:
        if first_token_time is None:
            first_token_time = time.perf_counter() - start
        token_count += 1

total_time = time.perf_counter() - start

print(f"Time to First Token (TTFT) : {first_token_time:.3f}s")
print(f"Tokens générés             : {token_count}")
print(f"Temps total                : {total_time:.3f}s")
print(f"Tokens/seconde             : {token_count / total_time:.1f}")

Le TTFT (Time to First Token), c’est-à-dire le temps entre l’envoi de la requête et l’affichage du premier caractère, est la métrique la plus importante pour l’expérience utilisateur. Elle prime sur le temps total : une génération de quinze secondes dont le premier mot apparaît en trois dixièmes de seconde sera perçue comme rapide, tandis qu’une génération de cinq secondes livrée d’un bloc sera perçue comme lente. Mesurez-la sur vos prompts réels, système compris — un message système de deux pages allonge le TTFT avant même que le premier token ne soit généré.

Points clés à retenir

  • Utilisez un pattern générateur (yield) pour propager les tokens en streaming
  • Implémentez un backoff exponentiel pour les erreurs 429 (rate limit)
  • Préférez stream_async pour les applications web concurrentes
  • Mesurez le TTFT et les tokens/seconde pour optimiser la performance
  • Avec FastAPI, StreamingResponse + text/event-stream propagent le SSE au navigateur