Diffuser la sortie avec LangChain
Implémentez la diffusion des tokens dans les chaînes LCEL afin que votre application affiche chaque mot dès son arrivée au lieu d’attendre la réponse complète, ce qui améliore la latence perçue.
Diffuser la sortie avec LangChain est une leçon AI Engineering Academy gratuite sur CoddyKit. Ceci est la leçon 4 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage AI Engineering Academy, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours AI Engineering Academy comprend 4 leçons au total.
Pourquoi la diffusion en continu est importante
Sans diffusion en continu, les utilisateurs regardent un écran vide en attendant que le LLM termine sa génération, ce qui peut prendre de 5 à 30 secondes pour les réponses longues. Avec la diffusion en continu, les jetons apparaissent au fur et à mesure de leur génération, ce qui fournit un retour immédiat. Cela améliore considérablement la réactivité perçue. Le LCEL de LangChain propage automatiquement la diffusion en continu dans toute la chaîne lorsque vous appelez .stream().
Diffusion en continu de base avec .stream()
Chaque chaîne LCEL expose une méthode .stream() qui renvoie un itérateur de fragments. Pour une chaîne se terminant par StrOutputParser, chaque fragment est une partie de chaîne. Vous parcourez les fragments et les affichez ou les produisez au fur et à mesure de leur arrivée. La diffusion en continu s’effectue au niveau HTTP : chaque jeton de l’API OpenAI est transmis par l’analyseur dès sa réception.
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
chain = (
ChatPromptTemplate.from_template('Explain {topic} in detail.')
| ChatOpenAI(model='gpt-4o-mini')
| StrOutputParser()
)
# Stream tokens to stdout
for chunk in chain.stream({'topic': 'quantum entanglement'}):
print(chunk, end='', flush=True)
print() # final newlineDiffusion en continu asynchrone avec .astream()
.astream() est la version asynchrone de .stream(). Elle renvoie un itérateur asynchrone que vous parcourez avec async for. Il s’agit de l’approche correcte dans FastAPI, Starlette et les autres infrastructures web asynchrones, où le gestionnaire de requête est une coroutine. Utiliser une diffusion synchrone dans un gestionnaire asynchrone bloquerait la boucle d’événements.
import asyncio
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate
chain = (
ChatPromptTemplate.from_template('Write a poem about {subject}')
| ChatOpenAI(model='gpt-4o-mini')
| StrOutputParser()
)
async def stream_response():
async for chunk in chain.astream({'subject': 'the ocean'}):
print(chunk, end='', flush=True)
asyncio.run(stream_response())Diffusion en continu dans FastAPI avec StreamingResponse
Dans FastAPI, vous encapsulez un générateur asynchrone dans StreamingResponse avec media_type='text/plain' afin de diffuser les jetons de texte vers le navigateur. Pour les événements envoyés par le serveur (SSE), utilisez media_type='text/event-stream' et formatez chaque fragment sous la forme data: ...\n\n. Le navigateur reçoit alors les jetons au fur et à mesure de leur génération, sans attendre la réponse complète.
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
async def generate_stream(topic: str):
async for chunk in chain.astream({'topic': topic}):
yield chunk
@app.get('/stream')
async def stream_endpoint(topic: str):
return StreamingResponse(
generate_stream(topic),
media_type='text/plain'
)
# SSE format for frontend EventSource
async def sse_stream(topic: str):
async for chunk in chain.astream({'topic': topic}):
yield f'data: {chunk}\n\n'Diffusion en continu à travers les étapes intermédiaires
Les chaînes LCEL propagent la diffusion en continu à travers chaque étape qui la prend en charge. Le StrOutputParser est compatible avec la diffusion en continu et transmet immédiatement les fragments. Toutefois, certains analyseurs, comme JsonOutputParser, doivent mettre en mémoire tampon la sortie complète avant de l’analyser, ce qui interrompt la diffusion en continu. LangChain rend ce comportement explicite : si une étape n’est pas compatible avec la diffusion en continu, elle accumule la sortie avant de la transmettre à l’étape suivante.
from langchain_core.output_parsers import JsonOutputParser
# This chain does NOT stream token by token
# JsonOutputParser must buffer the full response before parsing JSON
json_chain = (
ChatPromptTemplate.from_template('Return JSON: {task}')
| ChatOpenAI(model='gpt-4o-mini')
| JsonOutputParser() # buffers until complete
)
# But partial JSON streaming IS possible with streaming_json_parser
for partial in json_chain.stream({'task': 'list 3 colors'}):
print(partial) # prints partial dict as it fills inastream_events pour un contrôle précis
.astream_events() fournit une API de diffusion en continu plus détaillée, qui émet des événements pour chaque étape de la chaîne, et pas uniquement pour la sortie finale. Chaque événement possède un champ kind (on_chain_start, on_llm_stream, on_chain_end) et une charge utile data. Vous pouvez ainsi diffuser séparément les résultats des appels d’outils, le raisonnement intermédiaire et la sortie finale vers différentes parties d’une interface utilisateur.
async def stream_with_events(question: str):
async for event in chain.astream_events(
{'question': question},
version='v2'
):
kind = event['event']
if kind == 'on_llm_stream':
chunk = event['data']['chunk'].content
print(chunk, end='', flush=True)
elif kind == 'on_chain_end':
print('\n[Done]')
elif kind == 'on_tool_start':
print(f'\n[Tool: {event["name"]}]')Mise en mémoire tampon de la sortie diffusée
Vous devez parfois diffuser les jetons à l’utilisateur et récupérer la réponse complète pour la journalisation ou un traitement ultérieur. Utilisez .astream() avec un accumulateur sous forme de liste. Assemblez les fragments après la boucle pour obtenir le texte complet. Ce modèle vous permet d’afficher la sortie en continu en temps réel tout en stockant la réponse complète pour l’analyse, la mise en cache ou l’évaluation.
async def stream_and_capture(question: str) -> str:
full_response = []
async for chunk in chain.astream({'question': question}):
print(chunk, end='', flush=True) # stream to user
full_response.append(chunk) # also collect
print() # newline
complete = ''.join(full_response)
await log_response(question, complete) # log full text
return completeDiffusion en continu avec des appels d’outils
Lorsqu’un modèle génère un appel d’outil dans une réponse diffusée en continu, les arguments de la fonction arrivent sous forme de fragments de jetons. Vous devez mettre en mémoire tampon la chaîne d’arguments JSON jusqu’à la fin de l’appel d’outil avant de l’exécuter. LangChain gère automatiquement cette opération dans ses exécuteurs d’agents, mais si vous créez une boucle de diffusion personnalisée, vous devez vérifier finish_reason et accumuler les fragments de tool_call.function.arguments.
from openai import AsyncOpenAI
client = AsyncOpenAI()
async def stream_with_tools(prompt: str):
tool_call_buffer = {}
async with client.chat.completions.stream(
model='gpt-4o-mini',
messages=[{'role': 'user', 'content': prompt}],
tools=[weather_tool_schema]
) as stream:
async for chunk in stream:
delta = chunk.choices[0].delta
if delta.tool_calls:
for tc in delta.tool_calls:
idx = tc.index
if idx not in tool_call_buffer:
tool_call_buffer[idx] = ''
if tc.function.arguments:
tool_call_buffer[idx] += tc.function.argumentsAnnulation et délai d’expiration avec la diffusion en continu
Les réponses longues diffusées en continu doivent pouvoir être annulées. En Python asynchrone, vous pouvez annuler une asyncio.Task qui encapsule le flux. Dans FastAPI, l’infrastructure gère automatiquement l’annulation due à la déconnexion du client lors de l’utilisation de StreamingResponse. Définissez un délai d’expiration via le paramètre timeout du client OpenAI, ou encapsulez le flux avec asyncio.wait_for() afin de l’interrompre après une durée maximale.
import asyncio
async def stream_with_timeout(question: str, timeout: float = 30.0):
async def _stream():
async for chunk in chain.astream({'question': question}):
yield chunk
try:
async for chunk in asyncio.timeout(_stream(), timeout):
print(chunk, end='', flush=True)
except asyncio.TimeoutError:
print('\n[Stream timed out after 30 seconds]')
except asyncio.CancelledError:
print('\n[Stream cancelled by client disconnect]')SSE côté client avec JavaScript
Sur l’interface frontend, l’API EventSource native du navigateur consomme les événements envoyés par le serveur. Lorsque le point de terminaison FastAPI émet des fragments data: token\n\n, EventSource déclenche un événement message pour chacun d’eux. Ajoutez chaque jeton au DOM dès son arrivée afin de créer un effet de machine à écrire. Pour un contrôle accru, fetch() avec response.body.getReader() vous donne un accès complet à la diffusion en continu.
// Frontend JavaScript (not Python)
const source = new EventSource('/stream?topic=quantum+computing');
const outputDiv = document.getElementById('output');
source.onmessage = (event) => {
outputDiv.textContent += event.data;
};
source.onerror = () => {
source.close();
outputDiv.textContent += ' [done]';
};
// Alternative: fetch with ReadableStream
const response = await fetch('/stream?topic=ai');
const reader = response.body.getReader();
while (true) {
const {done, value} = await reader.read();
if (done) break;
outputDiv.textContent += new TextDecoder().decode(value);
}Bonnes pratiques pour la diffusion en continu
Suivez ces bonnes pratiques lors de la mise en œuvre de la diffusion en continu : utilisez toujours flush=True lors de l’affichage vers stdout afin d’éviter la mise en mémoire tampon. Définissez stream_usage=True si vous avez besoin d’un décompte précis des jetons pendant la diffusion. Émettez un marqueur data: [DONE]\n\n à la fin des flux SSE afin que le client sache quand fermer la connexion. Testez les points de terminaison de diffusion avec curl --no-buffer pour vérifier que les jetons arrivent progressivement.
# Complete SSE endpoint with DONE sentinel
async def sse_generator(question: str):
try:
async for chunk in chain.astream({'question': question}):
# Escape any newlines in the chunk
safe_chunk = chunk.replace('\n', ' ')
yield f'data: {safe_chunk}\n\n'
finally:
yield 'data: [DONE]\n\n'
@app.get('/chat/stream')
async def chat_stream(question: str):
return StreamingResponse(
sse_generator(question),
media_type='text/event-stream',
headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}
)Vérification rapide
Vérifiez votre compréhension de la sortie en continu dans LangChain.
Récapitulatif de la leçon
Dans cette leçon, vous avez appris que stream() et astream() permettent de parcourir les fragments de jetons au fur et à mesure de leur génération, ce qui évite la longue attente de la réponse complète ; que StreamingResponse dans FastAPI, au format SSE, transmet les jetons en temps réel aux navigateurs clients ; et que astream_events() fournit des points d’accroche événementiels précis pour chaque étape de la chaîne, y compris les appels d’outils et les sorties intermédiaires. Nous allons maintenant explorer la gestion de la mémoire pour les conversations à plusieurs tours.
Apprends Python avec un tuteur IA — gratuit
Écris et exécute du vrai code dans ton navigateur, obtiens de l'aide instantanée d'un tuteur IA disponible 24h/24, et reprends là où tu t'es arrêté sur le web ou dans l'app.
- Cours
- 30
- Leçons
- 120
Questions Fréquemment Posées
La leçon « Diffuser la sortie avec LangChain » est-elle gratuite ?
Oui — le texte complet de « Diffuser la sortie avec LangChain » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours AI Engineering Academy, passe à CoddyKit PRO. Le cours AI Engineering Academy comprend 4 leçons au total.
Qu'est-ce que j'apprendrai dans « Diffuser la sortie avec LangChain » ?
Implémentez la diffusion des tokens dans les chaînes LCEL afin que votre application affiche chaque mot dès son arrivée au lieu d’attendre la réponse complète, ce qui améliore la latence perçue. Tu pratiques AI Engineering Academy avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.
Dois-je avoir de l'expérience pour commencer AI Engineering Academy ?
Aucune expérience préalable n'est requise. AI Engineering Academy sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 4 sur 4.
Combien de temps prend la leçon « Diffuser la sortie avec LangChain » ?
La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.
Peux-tu écrire et exécuter du code dans cette leçon AI Engineering Academy ?
Oui. Chaque leçon AI Engineering Academy inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.
Toutes les leçons de ce cours
- Architecture de LangChain et abstractions fondamentales
- Construire des chaînes avec LCEL
- Chaînes conditionnelles et parallèles
- Diffuser la sortie avec LangChain