GuideAvancé

RAG Temps Reel : Architectures WebSocket pour des Reponses Instantanees

7 août 2026
21 min de lecture
Equipe Ailog

Guide complet des architectures RAG temps reel : WebSocket vs SSE vs HTTP streaming. Pipeline event-driven, mises a jour live, optimisation de latence avec FastAPI.

TL;DR

Le RAG temps reel combine WebSocket pour la communication bidirectionnelle, SSE pour le streaming unidirectionnel et un pipeline event-driven pour les mises a jour de documents instantanees. Resultat : des reponses en p50 < 800ms (TTFT) contre 2-5s avec une architecture classique. Ce guide couvre les 3 patterns architecturaux, un exemple complet FastAPI + WebSocket, et les benchmarks de latence par approche.

Pourquoi le RAG classique n'est pas assez rapide

Le probleme de la latence en RAG

Un pipeline RAG classique (HTTP request-response) a une latence inherente :

Utilisateur -> [HTTP Request]
  -> Embedding de la query (50-100ms)
  -> Recherche vectorielle (20-50ms)
  -> Reranking (100-200ms)
  -> Generation LLM (2000-5000ms)
  -> [HTTP Response complete]
Utilisateur recoit TOUT d'un coup apres 2-5 secondes

L'utilisateur attend la totalite de la reponse. Avec le streaming, il voit les premiers mots en < 1 seconde.

Impact sur l'experience utilisateur

ApprocheTTFT (premier token)Temps completPerception utilisateur
HTTP classique2-5s2-5s"C'est lent"
HTTP + streaming0,5-1,5s3-6s"C'est rapide" (premiers mots visibles)
WebSocket + streaming0,3-0,8s2-5s"C'est instantane"
WebSocket + cache0,05-0,2s0,5-2s"Wow"

La perception est tout. Meme si le temps total est similaire, le streaming change radicalement l'experience.

Les 3 patterns architecturaux

Pattern 1 : HTTP Polling (a eviter)

DEVELOPERpython
# ❌ Anti-pattern : polling HTTP # Le client interroge le serveur regulierement # Client (JavaScript) """ setInterval(async () => { const response = await fetch('/api/chat/status/' + taskId); if (response.data.status === 'complete') { displayAnswer(response.data.answer); } }, 500); // Poll toutes les 500ms """ # Problemes : # - Gaspillage de bande passante (requetes inutiles) # - Latence = intervalle de polling (500ms minimum) # - Charge serveur elevee # - Pas de streaming possible

Pattern 2 : Server-Sent Events (SSE)

DEVELOPERpython
# ✅ Bon pour le streaming unidirectionnel # Le serveur pousse des evenements vers le client from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio app = FastAPI() async def rag_stream(query: str): """Pipeline RAG avec streaming SSE""" # Phase 1 : Retrieval (envoyer un signal de progression) yield f"data: {json.dumps({'type': 'status', 'message': 'Recherche...'})}\n\n" chunks = await vector_search(query, top_k=10) reranked = await rerank(query, chunks, top_n=5) yield f"data: {json.dumps({'type': 'status', 'message': 'Generation...'})}\n\n" # Phase 2 : Generation avec streaming context = build_context(reranked) async for token in llm_stream(query, context): yield f"data: {json.dumps({'type': 'token', 'content': token})}\n\n" # Phase 3 : Sources sources = [{"title": c.title, "url": c.url} for c in reranked] yield f"data: {json.dumps({'type': 'sources', 'data': sources})}\n\n" yield f"data: {json.dumps({'type': 'done'})}\n\n" @app.get("/api/chat/stream") async def chat_stream(query: str): return StreamingResponse( rag_stream(query), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", } )

Avantages SSE : Simple, compatible HTTP/2, reconnexion automatique. Limites SSE : Unidirectionnel (serveur -> client uniquement), pas ideal pour le chat.

Pattern 3 : WebSocket (recommande pour le chat)

DEVELOPERpython
# ✅ Ideal pour le chat RAG bidirectionnel from fastapi import FastAPI, WebSocket, WebSocketDisconnect import json import asyncio app = FastAPI() class RAGWebSocketHandler: def __init__(self): self.retriever = VectorRetriever() self.reranker = CohereReranker() self.llm = StreamingLLM() async def handle_message(self, websocket: WebSocket, data: dict): query = data.get("message", "") conversation_id = data.get("conversation_id") # Signal : debut du traitement await websocket.send_json({ "type": "thinking", "message": "Analyse de votre question..." }) # Etape 1 : Retrieval parallele retrieval_start = time.time() chunks = await self.retriever.search(query, top_k=20) reranked = await self.reranker.rerank(query, chunks, top_n=5) await websocket.send_json({ "type": "retrieval_done", "latency_ms": int((time.time() - retrieval_start) * 1000), "chunks_found": len(reranked) }) # Etape 2 : Generation en streaming context = self.build_context(reranked) full_response = "" async for token in self.llm.stream(query, context): full_response += token await websocket.send_json({ "type": "token", "content": token }) # Etape 3 : Envoi des metadonnees await websocket.send_json({ "type": "complete", "sources": [ {"title": c.title, "score": c.score, "url": c.url} for c in reranked ], "usage": { "input_tokens": count_tokens(context), "output_tokens": count_tokens(full_response), } }) handler = RAGWebSocketHandler() @app.websocket("/ws/chat") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() try: while True: data = await websocket.receive_json() if data.get("type") == "message": await handler.handle_message(websocket, data) elif data.get("type") == "ping": await websocket.send_json({"type": "pong"}) except WebSocketDisconnect: print("Client disconnected")

Comparaison detaillee des 3 approches

Tableau comparatif

CritereHTTP PollingSSEWebSocket
DirectionClient -> ServeurServeur -> ClientBidirectionnel
TTFT500ms+ (intervalle)300-800ms200-500ms
StreamingNonOuiOui
ReconnexionManuelleAutomatiqueManuelle
OverheadEleve (headers HTTP)FaibleTres faible
CompatibiliteUniverselleExcellenteBonne
Proxies/CDNAucun problemeParfois bloqueParfois bloque
Multi-messagesN/ANonOui
Typing indicatorsImpossiblePossibleNaturel
Upload fichiersRequete separeeImpossiblePossible
ScalabiliteHauteMoyenneA gerer

Quand utiliser quoi ?

DEVELOPERpython
# Decision framework def choose_architecture(requirements): if requirements.get("bidirectional"): return "WebSocket" # Chat, collaboration if requirements.get("streaming") and not requirements.get("bidirectional"): return "SSE" # Notifications, dashboards if requirements.get("simple") and not requirements.get("real_time"): return "HTTP" # API REST classique # Cas d'usage typiques use_cases = { "WebSocket": [ "Chatbot RAG interactif", "Chat d'equipe avec IA", "Collaboration en temps reel", "Gaming / applications interactives", ], "SSE": [ "Streaming de reponses LLM (comme ChatGPT)", "Notifications de mise a jour de documents", "Dashboards en temps reel", "Progression de taches longues", ], "HTTP": [ "API batch processing", "Webhooks", "Integrations tierces", "Requetes ponctuelles", ], }

Architecture event-driven pour les mises a jour live

Le probleme des mises a jour de documents

Quand un document est modifie, comment le rendre disponible instantanement ?

Document modifie
  -> Webhook recu (50ms)
  -> Re-chunking (200-500ms)
  -> Re-embedding (100-300ms par chunk)
  -> Upsert dans Qdrant (50ms)
  -> Disponible pour les prochaines requetes
  Total : 500ms - 2s

Architecture event-driven complete

DEVELOPERpython
# Pipeline de mise a jour event-driven import asyncio from datetime import datetime class EventDrivenRAGPipeline: def __init__(self): self.event_bus = AsyncEventBus() self.chunker = SmartChunker() self.embedder = BatchEmbedder() self.vector_db = QdrantClient() async def on_document_updated(self, event: DocumentEvent): """Webhook recu : document modifie""" doc = event.document # Etape 1 : Re-chunking chunks = await self.chunker.chunk(doc.content, doc.metadata) # Etape 2 : Supprimer les anciens chunks await self.vector_db.delete( filter={"document_id": doc.id} ) # Etape 3 : Embeddings en batch embeddings = await self.embedder.embed_batch( [c.text for c in chunks] ) # Etape 4 : Upsert dans Qdrant points = [ { "id": chunk.id, "vector": embedding, "payload": { "text": chunk.text, "document_id": doc.id, "updated_at": datetime.utcnow().isoformat(), **chunk.metadata, } } for chunk, embedding in zip(chunks, embeddings) ] await self.vector_db.upsert(points) # Etape 5 : Notifier les clients connectes await self.event_bus.publish("document_updated", { "document_id": doc.id, "chunks_updated": len(chunks), }) # Webhook endpoint @app.post("/webhooks/document") async def document_webhook(self, payload: dict): event = DocumentEvent.from_webhook(payload) # Traitement asynchrone (ne bloque pas le webhook) asyncio.create_task(self.on_document_updated(event)) return {"status": "accepted"}

Notification des clients en temps reel

DEVELOPERpython
# Notifier les utilisateurs connectes que la base a ete mise a jour class ConnectionManager: def __init__(self): self.active_connections: dict[str, WebSocket] = {} async def broadcast_update(self, document_id: str): """Informer les clients que de nouvelles donnees sont disponibles""" message = { "type": "knowledge_updated", "document_id": document_id, "message": "La base de connaissances a ete mise a jour.", "timestamp": datetime.utcnow().isoformat() } for ws in self.active_connections.values(): try: await ws.send_json(message) except Exception: pass # Client deconnecte

Optimisation de la latence : retrieval parallele

Pipeline sequentiel vs parallele

DEVELOPERpython
# ❌ Pipeline sequentiel (lent) async def sequential_rag(query: str): # Total : 50 + 50 + 150 + 100 = 350ms avant la generation embedding = await embed_query(query) # 50ms chunks = await vector_search(embedding) # 50ms reranked = await rerank(query, chunks) # 150ms history = await get_conversation_history() # 100ms return await generate(query, reranked, history) # ✅ Pipeline parallele (rapide) async def parallel_rag(query: str): # Total : max(50+50+150, 100) = 250ms avant la generation # Soit 100ms de gagne (29% plus rapide) # Lancer en parallele les taches independantes embedding_task = asyncio.create_task(embed_query(query)) history_task = asyncio.create_task(get_conversation_history()) # Attendre l'embedding (necessaire pour la recherche) embedding = await embedding_task # Lancer la recherche vectorielle chunks = await vector_search(embedding) # Reranking et historique en parallele rerank_task = asyncio.create_task(rerank(query, chunks)) history = await history_task # Deja termine probablement reranked = await rerank_task return await generate(query, reranked, history)

Optimisations avancees

DEVELOPERpython
# Speculative retrieval : commencer la generation avant le reranking async def speculative_rag(query: str): embedding = await embed_query(query) chunks = await vector_search(embedding, top_k=20) # Commencer la generation avec le top-3 brut (sans reranking) # pendant que le reranking tourne quick_context = build_context(chunks[:3]) gen_task = asyncio.create_task( generate_stream(query, quick_context) ) # Reranking en parallele reranked = await rerank(query, chunks, top_n=5) # Si les resultats du reranking different significativement if reranked_differs_significantly(chunks[:3], reranked[:3]): gen_task.cancel() # Relancer avec le bon contexte return generate_stream(query, build_context(reranked)) # Sinon, continuer avec la generation en cours return gen_task

Exemple complet : FastAPI + WebSocket + RAG streaming

Backend complet

DEVELOPERpython
# app/main.py from fastapi import FastAPI, WebSocket, WebSocketDisconnect from fastapi.middleware.cors import CORSMiddleware import json import asyncio import time app = FastAPI(title="Real-time RAG API") app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"], ) class RealtimeRAG: """Pipeline RAG temps reel avec WebSocket""" def __init__(self): self.embedder = OpenAIEmbedder() self.qdrant = QdrantClient(url="http://localhost:6333") self.reranker = CohereReranker() self.llm = AsyncOpenAI() async def process_query(self, ws: WebSocket, query: str, conv_id: str): start = time.perf_counter() metrics = {} try: # Phase 1 : Retrieval await ws.send_json({"type": "phase", "phase": "retrieval"}) t0 = time.perf_counter() embedding = await self.embedder.embed(query) results = await self.qdrant.search( collection_name="knowledge_base", query_vector=embedding, limit=20, ) metrics["retrieval_ms"] = int((time.perf_counter() - t0) * 1000) # Phase 2 : Reranking await ws.send_json({"type": "phase", "phase": "reranking"}) t0 = time.perf_counter() texts = [r.payload["text"] for r in results] reranked = await self.reranker.rerank(query, texts, top_n=5) metrics["reranking_ms"] = int((time.perf_counter() - t0) * 1000) # Phase 3 : Generation streaming await ws.send_json({"type": "phase", "phase": "generation"}) context = "\n\n".join([r["text"] for r in reranked]) prompt = f"Contexte:\n{context}\n\nQuestion: {query}\n\nReponse:" full_response = "" t0 = time.perf_counter() first_token = True stream = await self.llm.chat.completions.create( model="gpt-4o", messages=[ {"role": "system", "content": "Reponds en te basant sur le contexte fourni."}, {"role": "user", "content": prompt} ], stream=True, ) async for chunk in stream: if chunk.choices[0].delta.content: token = chunk.choices[0].delta.content full_response += token if first_token: metrics["ttft_ms"] = int((time.perf_counter() - t0) * 1000) first_token = False await ws.send_json({ "type": "token", "content": token }) metrics["generation_ms"] = int((time.perf_counter() - t0) * 1000) metrics["total_ms"] = int((time.perf_counter() - start) * 1000) # Phase 4 : Completion await ws.send_json({ "type": "complete", "sources": [ {"text": r["text"][:100], "score": r["score"]} for r in reranked ], "metrics": metrics, }) except Exception as e: await ws.send_json({ "type": "error", "message": str(e) }) rag = RealtimeRAG() @app.websocket("/ws/chat/{conversation_id}") async def websocket_chat(websocket: WebSocket, conversation_id: str): await websocket.accept() await websocket.send_json({ "type": "connected", "conversation_id": conversation_id }) try: while True: data = await websocket.receive_json() if data["type"] == "message": await rag.process_query( websocket, data["content"], conversation_id ) elif data["type"] == "ping": await websocket.send_json({"type": "pong"}) except WebSocketDisconnect: pass

Client JavaScript

DEVELOPERjavascript
// client.js - Client WebSocket pour RAG temps reel class RAGWebSocketClient { constructor(conversationId) { this.ws = new WebSocket(`wss://api.example.com/ws/chat/${conversationId}`); this.responseContainer = document.getElementById('response'); this.setupHandlers(); } setupHandlers() { this.ws.onmessage = (event) => { const data = JSON.parse(event.data); switch (data.type) { case 'phase': this.showPhase(data.phase); break; case 'token': this.appendToken(data.content); break; case 'complete': this.showSources(data.sources); this.showMetrics(data.metrics); break; case 'error': this.showError(data.message); break; } }; } sendMessage(text) { this.responseContainer.innerHTML = ''; this.ws.send(JSON.stringify({ type: 'message', content: text })); } appendToken(token) { this.responseContainer.textContent += token; } }

Benchmarks de latence par approche

Protocole de test

DEVELOPERpython
# Benchmark sur 1000 requetes, base de 10K documents benchmark = { "queries": 1000, "documents": 10_000, "vector_db": "Qdrant", "llm": "GPT-4o", "embedding": "text-embedding-3-small", "server": "4 vCPU, 8GB RAM", }

Resultats p50/p95

MetriqueHTTP classiqueHTTP + SSEWebSocketWS + Cache
TTFT p502 100ms850ms620ms120ms
TTFT p954 200ms1 800ms1 200ms350ms
Temps total p503 500ms3 200ms2 800ms1 500ms
Temps total p956 800ms5 500ms4 800ms2 800ms
Connexions/serveurIllimite~5 000~10 000~10 000
Bande passanteHauteMoyenneFaibleFaible

Decomposition de la latence (WebSocket, p50)

Total : 620ms TTFT
  ├── Embedding query    :  45ms (7%)
  ├── Recherche Qdrant   :  35ms (6%)
  ├── Reranking Cohere   : 140ms (23%)
  ├── Overhead reseau    :  20ms (3%)
  └── LLM TTFT           : 380ms (61%)

Le goulot d'etranglement est le LLM. C'est pourquoi le caching est si efficace : il elimine completement le temps LLM pour les requetes repetees.

Gestion de la scalabilite WebSocket

Le defi des connexions persistantes

DEVELOPERpython
# Architecture scalable avec Redis pub/sub # aioredis est fusionne dans redis-py depuis la 4.2 : on importe via redis.asyncio from redis import asyncio as aioredis class ScalableWebSocketManager: """Gere les WebSocket sur plusieurs instances de serveur""" def __init__(self): self.redis = aioredis.from_url("redis://localhost:6379") self.local_connections: dict[str, WebSocket] = {} async def register(self, user_id: str, ws: WebSocket): self.local_connections[user_id] = ws # S'abonner au channel Redis de l'utilisateur pubsub = self.redis.pubsub() await pubsub.subscribe(f"user:{user_id}") asyncio.create_task(self._listen_redis(pubsub, ws)) async def _listen_redis(self, pubsub, ws: WebSocket): """Relayer les messages Redis vers le WebSocket""" async for message in pubsub.listen(): if message["type"] == "message": await ws.send_text(message["data"]) async def broadcast_to_user(self, user_id: str, data: dict): """Envoyer a un utilisateur (meme sur un autre serveur)""" await self.redis.publish( f"user:{user_id}", json.dumps(data) )

Load balancing avec sticky sessions

DEVELOPERnginx
# nginx.conf pour WebSocket avec sticky sessions upstream rag_backend { ip_hash; # Sticky sessions basees sur l'IP server backend1:8000; server backend2:8000; server backend3:8000; } server { location /ws/ { proxy_pass http://rag_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_set_header Host $host; proxy_read_timeout 86400; # 24h pour les connexions longues } }

FAQ

WebSocket ou SSE pour un chatbot RAG ?

WebSocket est recommande pour un chatbot car la communication est bidirectionnelle : l'utilisateur envoie des messages, le serveur stream les reponses. SSE convient si vous avez besoin uniquement de streaming serveur -> client (ex: notifications). En pratique, la plupart des chatbots RAG modernes utilisent WebSocket.

Comment gerer la reconnexion WebSocket ?

Implementez un backoff exponentiel cote client avec sauvegarde de l'etat de la conversation. Lors de la reconnexion, envoyez le dernier message_id recu pour que le serveur puisse reprendre la ou il s'est arrete. C'est exactement ce que fait le widget Ailog nativement.

Quel est l'impact sur la consommation serveur ?

Chaque connexion WebSocket consomme ~50KB de RAM. Un serveur avec 8GB peut gerer ~100 000 connexions inactives. En pratique, avec le traitement RAG, comptez ~5 000-10 000 connexions simultanees par serveur. Utilisez Redis pub/sub pour scaler horizontalement.

Les CDN supportent-ils les WebSocket ?

Oui, la plupart des CDN modernes (Cloudflare, AWS CloudFront, Fastly) supportent les WebSocket. Cependant, certains proxies d'entreprise peuvent les bloquer. Prevoyez un fallback SSE automatique pour ces cas.

Comment securiser les connexions WebSocket ?

Utilisez WSS (WebSocket Secure) en production, avec authentification par token JWT dans le handshake initial. Validez le token avant d'accepter la connexion. Implementez un rate limiting par utilisateur et par IP pour eviter les abus.

Conclusion

L'architecture RAG temps reel avec WebSocket offre une experience utilisateur incomparable :

  • TTFT < 800ms en p50 (vs 2s+ en HTTP classique)
  • Streaming natif des reponses token par token
  • Mises a jour live des documents via event-driven pipeline
  • Communication bidirectionnelle pour un vrai chat interactif

La cle est de combiner WebSocket pour le transport, le parallelisme pour le retrieval, et le caching pour les requetes frequentes.

Ailog integre nativement le WebSocket dans son widget et son API. Deployez un chatbot RAG temps reel en 5 minutes sur app.ailog.fr.


Voir aussi : Streaming des reponses RAG | Strategies de caching RAG | Orchestration d'agents RAG

Tags

RAGWebSocketSSEstreamingtemps reelFastAPIlatencearchitecturedeployment

Articles connexes

Ailog Assistant

Ici pour vous aider

Salut ! Pose-moi des questions sur Ailog et comment intégrer votre RAG dans vos projets !