Batch vs Echtzeit RAG: Die Architekturentscheidung die Ihr System macht oder bricht
Batch vs Echtzeit RAG: Wann welchen Ansatz verwenden, hybride Architektur, Queue-Systeme (Kafka, RabbitMQ, Redis Streams), Entscheidungsmatrix und Kostenvergleich.
TL;DR
Ein produktives RAG-System muss zwei grundlegend verschiedene Fluesse bewaeltigen: Batch (Dokumenten-Ingestion, Re-Indexierung, Analytics) und Echtzeit (Benutzerabfragen, Live-Updates). Die falsche Wahl fuehrt entweder zu inakzeptabler Latenz oder explodierenden Kosten. Dieser Leitfaden stellt Batch-, Echtzeit- und Hybrid-Architekturen mit Kostenvergleichen, Diagrammen und einer Entscheidungsmatrix vor.
Das grundlegende Dilemma
Jedes RAG-System steht vor demselben Kompromiss:
BATCH ECHTZEIT
───── ────────
Latenz: Minuten/Stunden Millisekunden
Durchsatz: Hoch (1M+ Docs/h) Mittel (100-1K Req/s)
Kosten: Niedrig pro Dokument Hoch pro Anfrage
Konsistenz: Eventuell Sofort
Komplexitaet: Niedrig Hoch
Anwendungsfall: Ingestion, Analytics Abfragen, Chat
Die richtige Architektur ist nicht das eine ODER das andere - es ist ein hybrides System, das jeden Ansatz am richtigen Ort einsetzt.
Batch-Architektur: Wann und wie
Batch-Anwendungsfaelle
| Operation | Haeufigkeit | Volumen | Akzeptable Latenz |
|---|---|---|---|
| Dokumenten-Ingestion | Taeglich/woechentlich | 1K-1M Docs | Minuten bis Stunden |
| Vollstaendige Re-Indexierung | Monatlich | Gesamter Korpus | Stunden |
| Embedding-Modell-Update | Bei Modellwechsel | Gesamter Korpus | Stunden |
| Analytics und Berichte | Taeglich | Alle Konversationen | Minuten |
| Bereinigung und Deduplizierung | Woechentlich | Gesamter Korpus | Stunden |
| Datenexport | Auf Anfrage | Variabel | Minuten |
Klassische Batch-Architektur
┌──────────────────────────────────────────────────────────────┐
│ BATCH-PIPELINE │
├──────────────────────────────────────────────────────────────┤
│ │
│ Quellen Queue Worker │
│ ┌─────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ S3/GCS │───>│ │───>│ Worker 1 │ │
│ │ Ext API │───>│ Redis Queue │───>│ Worker 2 │──> Qdrant│
│ │ Webhook │───>│ oder Celery │───>│ Worker 3 │ │
│ │ Upload │───>│ │───>│ Worker N │ │
│ └─────────┘ └──────────────┘ └──────────────┘ │
│ │
│ Schritte pro Worker: │
│ 1. Dokument herunterladen │
│ 2. Parsen (PDF, HTML, DOCX) │
│ 3. Chunking (semantisch/fest) │
│ 4. Embedding (Batch-API) │
│ 5. Upsert in Qdrant │
│ 6. Metadaten aktualisieren │
└──────────────────────────────────────────────────────────────┘
Implementierung mit Celery + Redis
DEVELOPERpythonfrom celery import Celery from celery.utils.log import get_task_logger app = Celery('rag_batch', broker='redis://localhost:6379/0') logger = get_task_logger(__name__) @app.task(bind=True, max_retries=3, default_retry_delay=60) def process_document(self, document_id: str, source_url: str): """Verarbeitet ein Dokument im Batch.""" try: # 1. Herunterladen content = download_document(source_url) logger.info(f"Heruntergeladen {document_id}: {len(content)} Bytes") # 2. Parsen parsed = parse_document(content, detect_format(source_url)) # 3. Chunking chunks = semantic_chunking(parsed.text, max_tokens=512, overlap=50) logger.info(f"{len(chunks)} Chunks erstellt fuer {document_id}") # 4. Embedding (Batch fuer Effizienz) embeddings = embed_batch( [c.text for c in chunks], model="text-embedding-3-large", batch_size=100 # 100 Chunks gleichzeitig ) # 5. Upsert in Qdrant points = [ { "id": f"{document_id}_{i}", "vector": emb, "payload": { "document_id": document_id, "chunk_index": i, "text": chunk.text, "metadata": chunk.metadata } } for i, (chunk, emb) in enumerate(zip(chunks, embeddings)) ] qdrant_client.upsert("documents", points) # 6. Status aktualisieren update_document_status(document_id, "indexed", chunks_count=len(chunks)) return {"document_id": document_id, "chunks": len(chunks)} except Exception as exc: logger.error(f"Verarbeitung fehlgeschlagen {document_id}: {exc}") self.retry(exc=exc) @app.task def batch_ingest(document_ids: list[str]): """Startet Batch-Ingestion mehrerer Dokumente.""" from celery import group tasks = group( process_document.s(doc_id, get_source_url(doc_id)) for doc_id in document_ids ) result = tasks.apply_async() return {"job_id": result.id, "total": len(document_ids)}
Optimierung: Batch-Embedding
Der wichtigste Batch-Trick: Embedding-APIs im Batch nutzen statt einzeln.
DEVELOPERpythonfrom openai import OpenAI client = OpenAI() def embed_batch(texts: list[str], model: str = "text-embedding-3-large", batch_size: int = 100) -> list[list[float]]: """Batch-Embedding - 10x schneller als einzeln.""" all_embeddings = [] for i in range(0, len(texts), batch_size): batch = texts[i:i + batch_size] response = client.embeddings.create(model=model, input=batch) batch_embeddings = [item.embedding for item in response.data] all_embeddings.extend(batch_embeddings) return all_embeddings # Leistungsvergleich # Einzeln: 1000 Docs x 200ms = 200 Sekunden # Batch 100: 10 Anfragen x 2s = 20 Sekunden (10x schneller)
Echtzeit-Architektur: Wann und wie
Echtzeit-Anwendungsfaelle
| Operation | Ziellatenz | Volumen | Kritikalitaet |
|---|---|---|---|
| Benutzerabfrage (Chatbot) | < 200ms Retrieval | 10-1K Req/s | Hoch |
| Live-Dokumentenaktualisierung | < 5s | 1-100/Min | Mittel |
| Antwort-Streaming | < 500ms TTFB | 10-500 Req/s | Hoch |
| Echtzeit-Vorschlaege | < 100ms | 100-10K Req/s | Hoch |
| Aenderungsbenachrichtigungen | < 1s | Variabel | Mittel |
Echtzeit-Architektur
┌──────────────────────────────────────────────────────────────┐
│ ECHTZEIT-PIPELINE │
├──────────────────────────────────────────────────────────────┤
│ │
│ Client API-Gateway RAG-Pipeline │
│ ┌───────┐ ┌──────────────┐ ┌──────────────────┐ │
│ │Widget │───>│ FastAPI │───>│ 1. Query-Embed │ │
│ │ API │───>│ + Rate Limit │───>│ 2. Vektorsuche │ │
│ │ Chat │───>│ + Auth │───>│ 3. Reranking │──> SSE│
│ └───────┘ └──────────────┘ │ 4. LLM-Generierung│ │
│ │ 5. Token-Streaming│ │
│ └──────────────────┘ │
│ │
│ Cache-Schichten: │
│ ┌─────────────────────────────────────────────┐ │
│ │ L1: In-Memory (exakte Uebereinstimmung) → 1ms│ │
│ │ L2: Redis (semantischer Cache) → 5ms │ │
│ │ L3: Qdrant (Vektorsuche) → 10-50ms │ │
│ └─────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────┘
Implementierung mit mehrstufigem Cache
DEVELOPERpythonimport hashlib import redis.asyncio as redis from qdrant_client import AsyncQdrantClient class RAGRealtimePipeline: def __init__(self): self.redis = redis.Redis() self.qdrant = AsyncQdrantClient() self.local_cache = {} # In-Memory LRU-Cache async def query(self, question: str, user_id: str) -> dict: # L1: Exakter Cache (In-Memory) cache_key = hashlib.md5(question.lower().strip().encode()).hexdigest() if cache_key in self.local_cache: return self.local_cache[cache_key] # L2: Semantischer Cache (Redis) cached = await self.redis.get(f"rag:cache:{cache_key}") if cached: result = json.loads(cached) self.local_cache[cache_key] = result return result # L3: Vollstaendige RAG-Pipeline query_embedding = await self.embed_query(question) # Vektorsuche search_results = await self.qdrant.search( collection_name="documents", query_vector=query_embedding, limit=10, score_threshold=0.7 ) # Reranking reranked = await self.rerank(question, search_results) # LLM-Generierung (Streaming) response = await self.generate(question, reranked[:5]) result = { "answer": response.text, "sources": [self._format_source(s) for s in reranked[:5]], "confidence": self._compute_confidence(reranked) } # Cache (TTL 1h) await self.redis.setex( f"rag:cache:{cache_key}", 3600, json.dumps(result) ) self.local_cache[cache_key] = result return result
Hybride Architektur: Das Beste aus beiden Welten
Die hybride Architektur ist der Produktionsstandard. Sie kombiniert Batch und Echtzeit mit einem Event-System.
Diagramm der hybriden Architektur
┌──────────────────────────────────────────────────────────────────┐
│ HYBRIDE RAG-ARCHITEKTUR │
├──────────────────────────────────────────────────────────────────┤
│ │
│ QUELLEN EVENT-BUS CONSUMER │
│ ┌──────────┐ ┌──────────────┐ │
│ │ Upload │─────────>│ │──> [Batch-Worker] │
│ │ API Sync │─────────>│ Kafka / │ → Ingestion │
│ │ Webhook │─────────>│ Redis │ → Re-Indexierung │
│ │ Cron Job │─────────>│ Streams │ → Analytics │
│ └──────────┘ │ │ │
│ │ │──> [Echtzeit-Worker] │
│ BENUTZER │ │ → Abfrageverarbeitung│
│ ┌──────────┐ │ │ → Live-Updates │
│ │ Widget │─────────>│ │ → Streaming │
│ │ API │─────────>│ │ │
│ │ Chat │─────────>│ │ │
│ └──────────┘ └──────────────┘ │
│ │
│ SPEICHER │
│ ┌───────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ PostgreSQL│ │ Qdrant │ │ Redis │ │ S3 │ │
│ │ (Metadata)│ │(Vektoren)│ │ (Cache) │ │ (Dateien)│ │
│ └───────────┘ └──────────┘ └──────────┘ └──────────┘ │
└──────────────────────────────────────────────────────────────────┘
Event-gesteuerte Dokumenten-Updates
DEVELOPERpythonimport redis.asyncio as redis import json from datetime import datetime class EventBus: def __init__(self, redis_url: str = "redis://localhost:6379"): self.redis = redis.from_url(redis_url) async def publish(self, event_type: str, data: dict): """Veroeffentlicht ein Event auf dem Bus.""" event = { "type": event_type, "data": data, "timestamp": datetime.utcnow().isoformat(), "id": str(uuid4()) } await self.redis.xadd( f"events:{event_type}", {"payload": json.dumps(event)} ) async def subscribe(self, event_type: str, consumer_group: str, consumer_name: str): """Abonniert einen Event-Typ.""" try: await self.redis.xgroup_create( f"events:{event_type}", consumer_group, id="0", mkstream=True ) except redis.ResponseError: pass # Gruppe existiert bereits while True: messages = await self.redis.xreadgroup( consumer_group, consumer_name, {f"events:{event_type}": ">"}, count=10, block=5000 ) for stream, msgs in messages: for msg_id, msg_data in msgs: event = json.loads(msg_data[b"payload"]) yield event await self.redis.xack( f"events:{event_type}", consumer_group, msg_id ) # Events veroeffentlichen event_bus = EventBus() # Wenn ein Dokument hochgeladen wird await event_bus.publish("document.uploaded", { "document_id": "doc_123", "source": "api_upload", "size_bytes": 45000 }) # Wenn ein Dokument aktualisiert wird await event_bus.publish("document.updated", { "document_id": "doc_123", "changed_sections": ["section_2", "section_5"], "update_type": "partial" })
Vergleich der Queue-Systeme
| Kriterium | Kafka | RabbitMQ | Redis Streams |
|---|---|---|---|
| Durchsatz | 100K+ Msg/s | 20K Msg/s | 50K Msg/s |
| Latenz | 5-15ms | 1-5ms | < 1ms |
| Persistenz | Disk (dauerhaft) | Memory + Disk | Memory + AOF |
| Reihenfolge | Pro Partition | Pro Queue | Pro Stream |
| Consumer Groups | Ja | Ja | Ja |
| Replay | Ja (Offset) | Nein | Ja (ID) |
| Ops-Komplexitaet | Hoch (KRaft-Cluster) | Mittel | Niedrig |
| Ideal fuer | Event Sourcing, hohes Volumen | Async Tasks, Routing | Cache + Queue, niedrige Latenz |
| RAG-Empfehlung | 10K+ Events/s | < 10K Events/s | Beste Wahl fuer RAG |
Warum Redis Streams fuer die meisten RAG-Systeme
Redis Streams bietet die beste Balance fuer ein typisches RAG-System:
- Bereits da fuer Caching: Redis dient als L2-Cache, keine zusaetzliche Komponente
- Sub-Millisekunden-Latenz: Perfekt fuer Live-Updates
- Consumer Groups: Lastverteilung zwischen Workern
- Replay: Wiederaufbereitung bei Fehler moeglich
- Geringe Komplexitaet: Kein Kafka-Cluster (KRaft-Controller, Broker) zu verwalten
Entscheidungsmatrix
Wann was verwenden?
| Szenario | Ansatz | Begruendung |
|---|---|---|
| Initiale Ingestion (1M+ Docs) | Batch | Volumen zu hoch fuer Echtzeit |
| Benutzerabfrage | Echtzeit | Latenz kritisch |
| Vom Benutzer hochgeladenes Dokument | Hybrid (Event -> Batch) | Hintergrund-Ingestion, Benachrichtigung wenn bereit |
| FAQ-Update | Echtzeit | Kleine Aenderung, sofortige Auswirkung |
| Embedding-Modellwechsel | Batch | Vollstaendige Korpus-Re-Indexierung |
| Analytics / Berichte | Batch | Keine Latenz-Einschraenkung |
| Live-Chat | Echtzeit | Streaming erforderlich |
| CRM-Synchronisation | Hybrid (Cron -> Batch) | Periodische Sync, mittleres Volumen |
| E-Commerce-Webhook | Hybrid (Event -> Echtzeit) | Schnelle Katalogaktualisierung |
Entscheidungsflussdiagramm
Ist die Latenz kritisch (<1s)?
├── JA → Uebersteigt das Volumen 100 Docs/Min?
│ ├── JA → HYBRIDE Architektur
│ │ (Queue + Echtzeit-Worker)
│ └── NEIN → Reine ECHTZEIT-Architektur
│ (Synchrone API + Cache)
└── NEIN → Uebersteigt das Volumen 10K Docs?
├── JA → Reine BATCH-Architektur
│ (Celery + verteilte Worker)
└── NEIN → Einfache BATCH-Architektur
(Cron Job + Python-Skript)
Kostenvergleich
Monatliche Infrastruktur (10M Dokumente, 100K Abfragen/Tag)
| Komponente | Nur Batch | Nur Echtzeit | Hybrid |
|---|---|---|---|
| Compute (Worker) | $200 | $500 | $400 |
| Redis | $50 | $150 | $100 |
| Qdrant | $320 | $320 | $320 |
| Kafka/Queue | $0 | $0 | $80 |
| LLM (Embeddings) | $150 | $300 | $200 |
| LLM (Generierung) | $0 | $2.000 | $2.000 |
| Gesamt | $720 | $3.270 | $3.100 |
Kosten pro Abfrage
| Ansatz | Kosten/Abfrage | Durchschnittl. Latenz | Durchsatz |
|---|---|---|---|
| Batch (vorberechnet) | $0,001 | N/A (offline) | 10K+ Docs/h |
| Echtzeit (on-the-fly) | $0,025 | 1,5-3s | 100-1K Req/s |
| Hybrid (Cache + Live) | $0,008 | 0,5-2s | 500-5K Req/s |
Das hybride Modell reduziert die Kosten pro Abfrage um 68% gegenueber reiner Echtzeit dank semantischem Caching.
Fortgeschrittene Patterns
Pattern: Write-Behind-Cache
Der Write-Behind-Cache ermoeglicht sofortige Antworten, waehrend der Index im Hintergrund aktualisiert wird.
DEVELOPERpythonclass WriteBehindRAG: """Antwortet aus dem Cache, aktualisiert Index im Hintergrund.""" async def update_document(self, doc_id: str, new_content: str): # 1. Cache sofort aktualisieren await self.redis.hset(f"doc:{doc_id}", "content", new_content) await self.redis.hset(f"doc:{doc_id}", "status", "pending_index") # 2. Event fuer Re-Indexierung veroeffentlichen await self.event_bus.publish("document.updated", { "document_id": doc_id, "update_type": "full" }) # 3. Abfragen nutzen in der Zwischenzeit den Cache return {"status": "updated", "indexing": "in_progress"} async def query(self, question: str): # Normale Vektorsuche results = await self.qdrant.search(question_embedding) # Mit Cache anreichern (kuerzlich geaenderte Dokumente) for result in results: cached = await self.redis.hgetall(f"doc:{result.id}") if cached and cached.get("status") == "pending_index": result.text = cached["content"] return results
Pattern: Backpressure
System-Ueberlastung verhindern bei Volumenspitzen.
DEVELOPERpythonclass BackpressureController: def __init__(self, max_queue_size: int = 10000): self.max_queue_size = max_queue_size async def should_accept(self, queue_name: str) -> bool: """Prueft, ob die Queue weitere Nachrichten akzeptieren kann.""" queue_size = await self.redis.xlen(f"events:{queue_name}") if queue_size > self.max_queue_size: return False if queue_size > self.max_queue_size * 0.8: await asyncio.sleep(0.1) return True async def ingest_with_backpressure(self, documents: list[dict]): """Ingestion mit Backpressure-Kontrolle.""" accepted, rejected = [], [] for doc in documents: if await self.should_accept("document.uploaded"): await self.event_bus.publish("document.uploaded", doc) accepted.append(doc["id"]) else: rejected.append(doc["id"]) return { "accepted": len(accepted), "rejected": len(rejected), "retry_after": 30 if rejected else None }
Unsere Architektur bei Ailog
Bei Ailog verwenden wir eine hybride Architektur:
- Batch: Kunden-Dokumenten-Ingestion (Celery + Redis), naechtliche Re-Indexierung
- Echtzeit: Widget-Abfragen mit semantischem Cache (Redis) + Qdrant
- Event-Bus: Redis Streams fuer Batch/Echtzeit-Koordination
- Mehrstufiger Cache: L1 In-Memory, L2 Redis, L3 Qdrant
Das Ergebnis: p50-Latenz von 800ms fuer eine vollstaendige Abfrage (Retrieval + LLM), mit Dokumenten die in weniger als 30 Sekunden nach dem Upload indexiert sind.
Entdecken Sie unseren Leitfaden zu RAG-Antwort-Streaming und Caching-Strategien.
FAQ
Fazit
Die Wahl zwischen Batch und Echtzeit ist nicht binaer. Die besten RAG-Systeme verwenden eine hybride Architektur:
- Batch fuer Massen-Ingestion und Re-Indexierung
- Echtzeit fuer Benutzerabfragen und kritische Updates
- Event-Bus (Redis Streams) zur Koordination beider
- Mehrstufiger Cache zur Reduzierung von Kosten und Latenz
Fangen Sie einfach an (Batch + synchrone API), dann fuegen Sie Komplexitaet hinzu, wenn das Volumen es erfordert.
Moechten Sie ein hybrides RAG-System ohne Infrastrukturverwaltung? Testen Sie Ailog - wir kuemmern uns um Batch, Echtzeit, Cache und Event-Bus fuer Sie.
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.