AnleitungExperte

Echtzeit-RAG: WebSocket-Architekturen fuer sofortige Antworten

7. August 2026
21 Minuten Lesezeit
Ailog Team

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

AnsatzTTFT (erster Token)GesamtzeitBenutzerwahrnehmung
Klassisches HTTP2-5s2-5s"Es ist langsam"
HTTP + Streaming0,5-1,5s3-6s"Es ist schnell" (erste Woerter sichtbar)
WebSocket + Streaming0,3-0,8s2-5s"Es ist sofort"
WebSocket + Cache0,05-0,2s0,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

KriteriumHTTP-PollingSSEWebSocket
RichtungClient -> ServerServer -> ClientBidirektional
TTFT500ms+ (Intervall)300-800ms200-500ms
StreamingNeinJaJa
WiederverbindungManuellAutomatischManuell
OverheadHoch (HTTP-Header)NiedrigSehr niedrig
KompatibilitaetUniversalAusgezeichnetGut
Proxies/CDNKein ProblemManchmal blockiertManchmal blockiert
Multi-NachrichtenN/ANeinJa
Tipp-IndikatorenUnmoeglichMoeglichNatuerlich
Datei-UploadSeparate AnfrageUnmoeglichMoeglich
SkalierbarkeitHochMittelHandhabbar

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

MetrikKlassisches HTTPHTTP + SSEWebSocketWS + Cache
TTFT p502.100ms850ms620ms120ms
TTFT p954.200ms1.800ms1.200ms350ms
Gesamtzeit p503.500ms3.200ms2.800ms1.500ms
Gesamtzeit p956.800ms5.500ms4.800ms2.800ms
Verbindungen/ServerUnbegrenzt~5.000~10.000~10.000
BandbreiteHochMittelNiedrigNiedrig

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

RAGWebSocketSSEStreamingEchtzeitFastAPILatenzArchitekturDeployment

Verwandte Artikel

Ailog Assistant

Ici pour vous aider

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