AnleitungExperte

Batch vs Echtzeit RAG: Die Architekturentscheidung die Ihr System macht oder bricht

31. August 2026
22 Min. Lesezeit
Ailog Team

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

OperationHaeufigkeitVolumenAkzeptable Latenz
Dokumenten-IngestionTaeglich/woechentlich1K-1M DocsMinuten bis Stunden
Vollstaendige Re-IndexierungMonatlichGesamter KorpusStunden
Embedding-Modell-UpdateBei ModellwechselGesamter KorpusStunden
Analytics und BerichteTaeglichAlle KonversationenMinuten
Bereinigung und DeduplizierungWoechentlichGesamter KorpusStunden
DatenexportAuf AnfrageVariabelMinuten

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

DEVELOPERpython
from 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.

DEVELOPERpython
from 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

OperationZiellatenzVolumenKritikalitaet
Benutzerabfrage (Chatbot)< 200ms Retrieval10-1K Req/sHoch
Live-Dokumentenaktualisierung< 5s1-100/MinMittel
Antwort-Streaming< 500ms TTFB10-500 Req/sHoch
Echtzeit-Vorschlaege< 100ms100-10K Req/sHoch
Aenderungsbenachrichtigungen< 1sVariabelMittel

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

DEVELOPERpython
import 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

DEVELOPERpython
import 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

KriteriumKafkaRabbitMQRedis Streams
Durchsatz100K+ Msg/s20K Msg/s50K Msg/s
Latenz5-15ms1-5ms< 1ms
PersistenzDisk (dauerhaft)Memory + DiskMemory + AOF
ReihenfolgePro PartitionPro QueuePro Stream
Consumer GroupsJaJaJa
ReplayJa (Offset)NeinJa (ID)
Ops-KomplexitaetHoch (KRaft-Cluster)MittelNiedrig
Ideal fuerEvent Sourcing, hohes VolumenAsync Tasks, RoutingCache + Queue, niedrige Latenz
RAG-Empfehlung10K+ Events/s< 10K Events/sBeste Wahl fuer RAG

Warum Redis Streams fuer die meisten RAG-Systeme

Redis Streams bietet die beste Balance fuer ein typisches RAG-System:

  1. Bereits da fuer Caching: Redis dient als L2-Cache, keine zusaetzliche Komponente
  2. Sub-Millisekunden-Latenz: Perfekt fuer Live-Updates
  3. Consumer Groups: Lastverteilung zwischen Workern
  4. Replay: Wiederaufbereitung bei Fehler moeglich
  5. Geringe Komplexitaet: Kein Kafka-Cluster (KRaft-Controller, Broker) zu verwalten

Entscheidungsmatrix

Wann was verwenden?

SzenarioAnsatzBegruendung
Initiale Ingestion (1M+ Docs)BatchVolumen zu hoch fuer Echtzeit
BenutzerabfrageEchtzeitLatenz kritisch
Vom Benutzer hochgeladenes DokumentHybrid (Event -> Batch)Hintergrund-Ingestion, Benachrichtigung wenn bereit
FAQ-UpdateEchtzeitKleine Aenderung, sofortige Auswirkung
Embedding-ModellwechselBatchVollstaendige Korpus-Re-Indexierung
Analytics / BerichteBatchKeine Latenz-Einschraenkung
Live-ChatEchtzeitStreaming erforderlich
CRM-SynchronisationHybrid (Cron -> Batch)Periodische Sync, mittleres Volumen
E-Commerce-WebhookHybrid (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)

KomponenteNur BatchNur EchtzeitHybrid
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

AnsatzKosten/AbfrageDurchschnittl. LatenzDurchsatz
Batch (vorberechnet)$0,001N/A (offline)10K+ Docs/h
Echtzeit (on-the-fly)$0,0251,5-3s100-1K Req/s
Hybrid (Cache + Live)$0,0080,5-2s500-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.

DEVELOPERpython
class 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.

DEVELOPERpython
class 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

Fuer die meisten RAG-Systeme (< 10K Events/Sekunde) ist **Redis Streams** die beste Wahl, weil es bereits fuer das Caching vorhanden ist, Sub-Millisekunden-Latenz bietet und einfach zu betreiben ist. Kafka wird relevant jenseits von 10K Events/Sekunde oder wenn Sie langfristige Aufbewahrung und Event Sourcing benoetigen. Bei Ailog deckt Redis Streams 100% unserer Beduerfnisse ab.
Verwenden Sie das Write-Behind-Pattern: Cache sofort aktualisieren, dann ein Event fuer die Re-Indexierung veroeffentlichen. Waehrend des Update-Fensters (normalerweise < 30s) Suchergebnisse mit Cache-Daten anreichern. TTL auf den Cache setzen, um veraltete Daten zu vermeiden.
Die hybride Architektur kostet etwa 4x mehr als reiner Batch (Infrastruktur), reduziert aber die Kosten pro Abfrage um 68% gegenueber reiner Echtzeit dank Caching. Fuer 100K Abfragen/Tag rechnen Sie mit etwa $3.100/Monat hybrid gegenueber $3.270 bei reiner Echtzeit. Der echte Vorteil ist die Latenz: 800ms hybrid gegenueber 2-3s Echtzeit ohne Cache.
Einfache Regel: 1 Worker kann etwa 100-500 Dokumente/Stunde verarbeiten (abhaengig von Groesse und Komplexitaet). Fuer eine initiale Ingestion von 1M Dokumenten planen Sie 20-50 Worker fuer 10-20 Stunden. Im laufenden Betrieb reichen 2-5 Worker fuer die meisten Faelle. Verwenden Sie Auto-Scaling basierend auf der Queue-Groesse.
Mit aktivierter AOF-Persistenz (appendfsync=everysec) kann Redis Streams bei einem Crash maximal 1 Sekunde an Daten verlieren. Fuer kritische Systeme verwenden Sie Redis Sentinel oder Redis Cluster Replikation. Consumer Groups garantieren, dass eine Nachricht mindestens einmal verarbeitet wird (At-Least-Once-Delivery mit ACK). Fuer Exactly-Once fuegen Sie einen Idempotenz-Schluessel in Ihrem Consumer hinzu. ---

Fazit

Die Wahl zwischen Batch und Echtzeit ist nicht binaer. Die besten RAG-Systeme verwenden eine hybride Architektur:

  1. Batch fuer Massen-Ingestion und Re-Indexierung
  2. Echtzeit fuer Benutzerabfragen und kritische Updates
  3. Event-Bus (Redis Streams) zur Koordination beider
  4. 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

RAGArchitekturBatchEchtzeitKafkaRedisEvent-DrivenPipeline

Verwandte Artikel

Ailog Assistant

Ici pour vous aider

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