Echtzeit-RAG: WebSocket-Architekturen fuer sofortige Antworten
Vollstaendiger Leitfaden fuer Echtzeit-RAG-Architekturen: WebSocket vs SSE vs HTTP-Streaming. Event-driven Pipeline, Live-Dokumentenaktualisierungen, Latenzoptimierung mit FastAPI.
TL;DR
Echtzeit-RAG kombiniert WebSocket fuer bidirektionale Kommunikation, SSE fuer unidirektionales Streaming und eine event-driven Pipeline fuer sofortige Dokumentenaktualisierungen. Ergebnis: Antworten in p50 < 800ms (TTFT) gegenueber 2-5s mit einer klassischen Architektur. Dieser Leitfaden behandelt die 3 Architekturmuster, ein vollstaendiges FastAPI + WebSocket-Beispiel und Latenz-Benchmarks nach Ansatz.
Warum klassisches RAG nicht schnell genug ist
Das Latenzproblem im RAG
Eine klassische RAG-Pipeline (HTTP Request-Response) hat inhaerent hohe Latenz:
Benutzer -> [HTTP Request]
-> Query-Embedding (50-100ms)
-> Vektorsuche (20-50ms)
-> Reranking (100-200ms)
-> LLM-Generierung (2000-5000ms)
-> [Komplette HTTP Response]
Benutzer erhaelt ALLES auf einmal nach 2-5 Sekunden
Der Benutzer wartet auf die gesamte Antwort. Mit Streaming sieht er die ersten Woerter in < 1 Sekunde.
Auswirkungen auf die Benutzererfahrung
| Ansatz | TTFT (erster Token) | Gesamtzeit | Benutzerwahrnehmung |
|---|---|---|---|
| Klassisches HTTP | 2-5s | 2-5s | "Es ist langsam" |
| HTTP + Streaming | 0,5-1,5s | 3-6s | "Es ist schnell" (erste Woerter sichtbar) |
| WebSocket + Streaming | 0,3-0,8s | 2-5s | "Es ist sofort" |
| WebSocket + Cache | 0,05-0,2s | 0,5-2s | "Wow" |
Wahrnehmung ist alles. Selbst wenn die Gesamtzeit aehnlich ist, veraendert Streaming die Erfahrung radikal.
Die 3 Architekturmuster
Muster 1: HTTP-Polling (vermeiden)
DEVELOPERpython# Anti-Pattern: HTTP-Polling # Der Client fragt den Server regelmaessig ab # Client (JavaScript) """ setInterval(async () => { const response = await fetch('/api/chat/status/' + taskId); if (response.data.status === 'complete') { displayAnswer(response.data.answer); } }, 500); // Alle 500ms abfragen """ # Probleme: # - Bandbreitenverschwendung (unnoetige Anfragen) # - Latenz = Polling-Intervall (mindestens 500ms) # - Hohe Serverlast # - Kein Streaming moeglich
Muster 2: Server-Sent Events (SSE)
DEVELOPERpython# Gut fuer unidirektionales Streaming # Der Server sendet Events an den Client from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio app = FastAPI() async def rag_stream(query: str): """RAG-Pipeline mit SSE-Streaming""" # Phase 1: Retrieval (Fortschrittssignal senden) yield f"data: {json.dumps({'type': 'status', 'message': 'Suche...'})}\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': 'Generierung...'})}\n\n" # Phase 2: Generierung mit 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: Quellen 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", } )
SSE-Vorteile: Einfach, HTTP/2-kompatibel, automatische Wiederverbindung. SSE-Einschraenkungen: Unidirektional (nur Server -> Client), nicht ideal fuer Chat.
Muster 3: WebSocket (empfohlen fuer Chat)
DEVELOPERpython# Ideal fuer bidirektionalen RAG-Chat 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: Verarbeitung gestartet await websocket.send_json({ "type": "thinking", "message": "Analyse Ihrer Frage..." }) # Schritt 1: Paralleles Retrieval 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) }) # Schritt 2: Streaming-Generierung 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 }) # Schritt 3: Metadaten senden 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 getrennt")
Detaillierter Vergleich der 3 Ansaetze
Vergleichstabelle
| Kriterium | HTTP-Polling | SSE | WebSocket |
|---|---|---|---|
| Richtung | Client -> Server | Server -> Client | Bidirektional |
| TTFT | 500ms+ (Intervall) | 300-800ms | 200-500ms |
| Streaming | Nein | Ja | Ja |
| Wiederverbindung | Manuell | Automatisch | Manuell |
| Overhead | Hoch (HTTP-Header) | Niedrig | Sehr niedrig |
| Kompatibilitaet | Universal | Ausgezeichnet | Gut |
| Proxies/CDN | Kein Problem | Manchmal blockiert | Manchmal blockiert |
| Multi-Nachrichten | N/A | Nein | Ja |
| Tipp-Indikatoren | Unmoeglich | Moeglich | Natuerlich |
| Datei-Upload | Separate Anfrage | Unmoeglich | Moeglich |
| Skalierbarkeit | Hoch | Mittel | Handhabbar |
Wann was verwenden?
DEVELOPERpython# Entscheidungsframework def choose_architecture(requirements): if requirements.get("bidirectional"): return "WebSocket" # Chat, Kollaboration if requirements.get("streaming") and not requirements.get("bidirectional"): return "SSE" # Benachrichtigungen, Dashboards if requirements.get("simple") and not requirements.get("real_time"): return "HTTP" # Klassische REST-API # Typische Anwendungsfaelle use_cases = { "WebSocket": [ "Interaktiver RAG-Chatbot", "Team-Chat mit KI", "Echtzeit-Kollaboration", "Gaming / Interaktive Anwendungen", ], "SSE": [ "LLM-Antwort-Streaming (wie ChatGPT)", "Dokumenten-Update-Benachrichtigungen", "Echtzeit-Dashboards", "Fortschritt langer Aufgaben", ], "HTTP": [ "Batch-Processing-API", "Webhooks", "Drittanbieter-Integrationen", "Einmalige Anfragen", ], }
Event-Driven-Architektur fuer Live-Updates
Das Problem der Dokumentenaktualisierung
Wenn ein Dokument geaendert wird, wie macht man es sofort verfuegbar?
Dokument geaendert
-> Webhook empfangen (50ms)
-> Re-Chunking (200-500ms)
-> Re-Embedding (100-300ms pro Chunk)
-> Upsert in Qdrant (50ms)
-> Verfuegbar fuer naechste Anfragen
Gesamt: 500ms - 2s
Vollstaendige Event-Driven-Architektur
DEVELOPERpython# Event-driven Update-Pipeline 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 empfangen: Dokument geaendert""" doc = event.document # Schritt 1: Re-Chunking chunks = await self.chunker.chunk(doc.content, doc.metadata) # Schritt 2: Alte Chunks loeschen await self.vector_db.delete( filter={"document_id": doc.id} ) # Schritt 3: Batch-Embeddings embeddings = await self.embedder.embed_batch( [c.text for c in chunks] ) # Schritt 4: Upsert in 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) # Schritt 5: Verbundene Clients benachrichtigen await self.event_bus.publish("document_updated", { "document_id": doc.id, "chunks_updated": len(chunks), })
Echtzeit-Benachrichtigung der Clients
DEVELOPERpython# Verbundene Benutzer benachrichtigen, dass die Wissensbasis aktualisiert wurde class ConnectionManager: def __init__(self): self.active_connections: dict[str, WebSocket] = {} async def broadcast_update(self, document_id: str): """Clients informieren, dass neue Daten verfuegbar sind""" message = { "type": "knowledge_updated", "document_id": document_id, "message": "Die Wissensbasis wurde aktualisiert.", "timestamp": datetime.utcnow().isoformat() } for ws in self.active_connections.values(): try: await ws.send_json(message) except Exception: pass # Getrennter Client
Latenzoptimierung: Paralleles Retrieval
Sequenzielle vs. parallele Pipeline
DEVELOPERpython# Sequenzielle Pipeline (langsam) async def sequential_rag(query: str): # Gesamt: 50 + 50 + 150 + 100 = 350ms vor der Generierung 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) # Parallele Pipeline (schnell) async def parallel_rag(query: str): # Gesamt: max(50+50+150, 100) = 250ms vor der Generierung # Das sind 100ms gespart (29% schneller) # Unabhaengige Aufgaben parallel starten embedding_task = asyncio.create_task(embed_query(query)) history_task = asyncio.create_task(get_conversation_history()) # Auf Embedding warten (fuer Suche benoetigt) embedding = await embedding_task # Vektorsuche starten chunks = await vector_search(embedding) # Reranking und History parallel rerank_task = asyncio.create_task(rerank(query, chunks)) history = await history_task # Wahrscheinlich bereits fertig reranked = await rerank_task return await generate(query, reranked, history)
Fortgeschrittene Optimierungen
DEVELOPERpython# Spekulative Retrieval: Generierung vor dem Reranking starten async def speculative_rag(query: str): embedding = await embed_query(query) chunks = await vector_search(embedding, top_k=20) # Generierung mit rohem Top-3 starten (ohne Reranking) # waehrend das Reranking laeuft quick_context = build_context(chunks[:3]) gen_task = asyncio.create_task( generate_stream(query, quick_context) ) # Reranking parallel reranked = await rerank(query, chunks, top_n=5) # Wenn Reranking-Ergebnisse signifikant abweichen if reranked_differs_significantly(chunks[:3], reranked[:3]): gen_task.cancel() # Mit richtigem Kontext neu starten return generate_stream(query, build_context(reranked)) # Andernfalls mit aktueller Generierung fortfahren return gen_task
Vollstaendiges Beispiel: FastAPI + WebSocket + RAG-Streaming
Vollstaendiges Backend
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="Echtzeit-RAG-API") app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"], ) class RealtimeRAG: """Echtzeit-RAG-Pipeline mit 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: Streaming-Generierung await ws.send_json({"type": "phase", "phase": "generation"}) context = "\n\n".join([r["text"] for r in reranked]) prompt = f"Kontext:\n{context}\n\nFrage: {query}\n\nAntwort:" full_response = "" t0 = time.perf_counter() first_token = True stream = await self.llm.chat.completions.create( model="gpt-4o", messages=[ {"role": "system", "content": "Antworte basierend auf dem bereitgestellten Kontext."}, {"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: Abschluss 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
JavaScript-Client
DEVELOPERjavascript// client.js - WebSocket-Client fuer Echtzeit-RAG 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; } }
Latenz-Benchmarks nach Ansatz
Testprotokoll
DEVELOPERpython# Benchmark ueber 1000 Anfragen, 10K Dokumentenbasis benchmark = { "queries": 1000, "documents": 10_000, "vector_db": "Qdrant", "llm": "GPT-4o", "embedding": "text-embedding-3-small", "server": "4 vCPU, 8GB RAM", }
p50/p95-Ergebnisse
| Metrik | Klassisches HTTP | HTTP + SSE | WebSocket | WS + Cache |
|---|---|---|---|---|
| TTFT p50 | 2.100ms | 850ms | 620ms | 120ms |
| TTFT p95 | 4.200ms | 1.800ms | 1.200ms | 350ms |
| Gesamtzeit p50 | 3.500ms | 3.200ms | 2.800ms | 1.500ms |
| Gesamtzeit p95 | 6.800ms | 5.500ms | 4.800ms | 2.800ms |
| Verbindungen/Server | Unbegrenzt | ~5.000 | ~10.000 | ~10.000 |
| Bandbreite | Hoch | Mittel | Niedrig | Niedrig |
Latenzaufschluesselung (WebSocket, p50)
Gesamt: 620ms TTFT
|- Query-Embedding : 45ms (7%)
|- Qdrant-Suche : 35ms (6%)
|- Cohere-Reranking : 140ms (23%)
|- Netzwerk-Overhead : 20ms (3%)
+- LLM TTFT : 380ms (61%)
Der Engpass ist das LLM. Deshalb ist Caching so effektiv: Es eliminiert die LLM-Zeit fuer wiederholte Anfragen vollstaendig.
WebSocket-Skalierbarkeitsmanagement
Die Herausforderung persistenter Verbindungen
DEVELOPERpython# Skalierbare Architektur mit Redis Pub/Sub # aioredis wurde in redis-py 4.2 integriert: Import ueber redis.asyncio from redis import asyncio as aioredis class ScalableWebSocketManager: """Verwaltet WebSockets ueber mehrere Serverinstanzen""" 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 # Redis-Channel des Benutzers abonnieren 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): """Redis-Nachrichten an WebSocket weiterleiten""" 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): """An einen Benutzer senden (auch auf einem anderen Server)""" await self.redis.publish( f"user:{user_id}", json.dumps(data) )
Load Balancing mit Sticky Sessions
DEVELOPERnginx# nginx.conf fuer WebSocket mit Sticky Sessions upstream rag_backend { ip_hash; # IP-basierte Sticky Sessions 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 fuer lange Verbindungen } }
FAQ
WebSocket oder SSE fuer einen RAG-Chatbot?
WebSocket wird fuer einen Chatbot empfohlen, da die Kommunikation bidirektional ist: Der Benutzer sendet Nachrichten, der Server streamt Antworten. SSE eignet sich, wenn Sie nur Server -> Client-Streaming benoetigen (z.B. Benachrichtigungen). In der Praxis verwenden die meisten modernen RAG-Chatbots WebSocket.
Wie behandelt man die WebSocket-Wiederverbindung?
Implementieren Sie ein exponentielles Backoff auf der Client-Seite mit Gespraechszustandsspeicherung. Bei der Wiederverbindung senden Sie die letzte empfangene message_id, damit der Server dort weitermachen kann, wo er aufgehoert hat. Genau das macht das Ailog-Widget nativ.
Welche Auswirkungen hat es auf den Serververbrauch?
Jede WebSocket-Verbindung verbraucht ~50KB RAM. Ein Server mit 8GB kann ~100.000 inaktive Verbindungen verarbeiten. In der Praxis rechnen Sie mit RAG-Verarbeitung mit ~5.000-10.000 gleichzeitigen Verbindungen pro Server. Verwenden Sie Redis Pub/Sub fuer horizontale Skalierung.
Unterstuetzen CDNs WebSocket?
Ja, die meisten modernen CDNs (Cloudflare, AWS CloudFront, Fastly) unterstuetzen WebSocket. Allerdings koennen einige Unternehmens-Proxies sie blockieren. Planen Sie ein automatisches SSE-Fallback fuer diese Faelle ein.
Wie sichert man WebSocket-Verbindungen?
Verwenden Sie WSS (WebSocket Secure) in der Produktion mit JWT-Token-Authentifizierung beim initialen Handshake. Validieren Sie das Token, bevor Sie die Verbindung akzeptieren. Implementieren Sie Rate-Limiting pro Benutzer und pro IP, um Missbrauch zu verhindern.
Fazit
Echtzeit-RAG-Architektur mit WebSocket liefert eine unuebertroffene Benutzererfahrung:
- TTFT < 800ms bei p50 (vs 2s+ mit klassischem HTTP)
- Natives Streaming der Antworten Token fuer Token
- Live-Updates von Dokumenten ueber Event-driven Pipeline
- Bidirektionale Kommunikation fuer echten interaktiven Chat
Der Schluessel ist die Kombination von WebSocket fuer den Transport, Parallelismus fuer das Retrieval und Caching fuer haeufige Anfragen.
Ailog integriert WebSocket nativ in sein Widget und seine API. Stellen Sie einen Echtzeit-RAG-Chatbot in 5 Minuten bereit auf app.ailog.fr.
Siehe auch: Streaming von RAG-Antworten | RAG-Caching-Strategien | RAG-Agenten-Orchestrierung
Tags
Verwandte Artikel
Sicherheit und Compliance für RAG: DSGVO, AI Act und Best Practices
Sichern Sie Ihr RAG-System: DSGVO-Konformität, europäischer AI Act, Datenschutz und Audit. Umfassender Leitfaden für Unternehmen.
RAG für KMU: Kompletter Leitfaden ohne Data-Team
Setzen Sie ein leistungsfähiges RAG-System in Ihrem KMU ein, ganz ohne fortgeschrittene technische Kenntnisse: No-code-Lösungen, kontrolliertes Budget und schnellen ROI.
Souveräner RAG: Hosting in Frankreich und europäische Daten
Setzen Sie einen souveränen RAG in Frankreich ein: lokales Hosting, DSGVO‑Konformität, Alternativen zu GAFAM und Best Practices für europäische Daten.