Batch vs Temps Reel RAG : Le Choix d'Architecture qui Fait ou Defait votre Systeme
Batch vs temps reel RAG : quand utiliser chaque approche, architecture hybride, systemes de queues (Kafka, RabbitMQ, Redis Streams), matrice de decision et comparaison des couts.
TL;DR
Un systeme RAG en production doit gerer deux flux radicalement differents : le batch (ingestion de documents, re-indexation, analytics) et le temps reel (requetes utilisateurs, mises a jour live). Le mauvais choix entraine soit une latence inacceptable, soit des couts explosifs. Ce guide presente les architectures batch, temps reel et hybride, avec des comparaisons de couts, des diagrammes et une matrice de decision.
Le dilemme fondamental
Chaque systeme RAG fait face au meme compromis :
BATCH TEMPS REEL
───── ──────────
Latence : Minutes/heures Millisecondes
Debit : Eleve (1M+ docs/h) Moyen (100-1K req/s)
Cout : Faible par document Eleve par requete
Coherence : Eventuellement Immediatement
Complexite : Faible Elevee
Cas d'usage : Ingestion, analytics Requetes, chat
La bonne architecture n'est pas l'une OU l'autre - c'est un systeme hybride qui utilise chaque approche au bon endroit.
Architecture Batch : quand et comment
Cas d'usage du batch
| Operation | Frequence | Volume | Latence acceptable |
|---|---|---|---|
| Ingestion de documents | Quotidien/hebdo | 1K-1M docs | Minutes a heures |
| Re-indexation complete | Mensuel | Tout le corpus | Heures |
| Mise a jour embeddings | Lors d'un changement de modele | Tout le corpus | Heures |
| Analytics et rapports | Quotidien | Toutes les conversations | Minutes |
| Nettoyage et deduplication | Hebdomadaire | Tout le corpus | Heures |
| Export de donnees | A la demande | Variable | Minutes |
Architecture batch classique
┌──────────────────────────────────────────────────────────────┐
│ PIPELINE BATCH │
├──────────────────────────────────────────────────────────────┤
│ │
│ Sources Queue Workers │
│ ┌─────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ S3/GCS │───>│ │───>│ Worker 1 │ │
│ │ API ext │───>│ Redis Queue │───>│ Worker 2 │──> Qdrant│
│ │ Webhook │───>│ ou Celery │───>│ Worker 3 │ │
│ │ Upload │───>│ │───>│ Worker N │ │
│ └─────────┘ └──────────────┘ └──────────────┘ │
│ │
│ Etapes par worker : │
│ 1. Telecharger le document │
│ 2. Parser (PDF, HTML, DOCX) │
│ 3. Chunking (semantic/fixed) │
│ 4. Embedding (batch API) │
│ 5. Upsert dans Qdrant │
│ 6. Mettre a jour les metadonnees │
└──────────────────────────────────────────────────────────────┘
Implementation avec 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): """Traite un document en batch.""" try: # 1. Telecharger content = download_document(source_url) logger.info(f"Downloaded {document_id}: {len(content)} bytes") # 2. Parser parsed = parse_document(content, detect_format(source_url)) # 3. Chunking chunks = semantic_chunking(parsed.text, max_tokens=512, overlap=50) logger.info(f"Created {len(chunks)} chunks for {document_id}") # 4. Embedding (batch pour efficacite) embeddings = embed_batch( [c.text for c in chunks], model="text-embedding-3-large", batch_size=100 # 100 chunks a la fois ) # 5. Upsert dans 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. Mettre a jour le statut 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"Failed to process {document_id}: {exc}") self.retry(exc=exc) @app.task def batch_ingest(document_ids: list[str]): """Lance l'ingestion batch de plusieurs documents.""" 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)}
Optimisation : batch embedding
L'astuce la plus importante en batch : utiliser les API d'embedding en batch plutot qu'un par un.
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]]: """Embedding par batch - 10x plus rapide que unitaire.""" 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 # Comparaison de performance # Unitaire : 1000 docs × 200ms = 200 secondes # Batch 100 : 10 requetes × 2s = 20 secondes (10x plus rapide)
Architecture Temps Reel : quand et comment
Cas d'usage temps reel
| Operation | Latence cible | Volume | Criticite |
|---|---|---|---|
| Requete utilisateur (chatbot) | < 200ms retrieval | 10-1000 req/s | Haute |
| Mise a jour live de document | < 5s | 1-100/min | Moyenne |
| Streaming de reponse | < 500ms TTFB | 10-500 req/s | Haute |
| Suggestion temps reel | < 100ms | 100-10K req/s | Haute |
| Notification de changement | < 1s | Variable | Moyenne |
Architecture temps reel
┌──────────────────────────────────────────────────────────────┐
│ PIPELINE TEMPS REEL │
├──────────────────────────────────────────────────────────────┤
│ │
│ Client API Gateway RAG Pipeline │
│ ┌───────┐ ┌──────────────┐ ┌──────────────────┐ │
│ │Widget │───>│ FastAPI │───>│ 1. Query embed │ │
│ │ API │───>│ + Rate limit │───>│ 2. Vector search │ │
│ │ Chat │───>│ + Auth │───>│ 3. Rerank │──> SSE│
│ └───────┘ └──────────────┘ │ 4. LLM generate │ │
│ │ 5. Stream tokens │ │
│ └──────────────────┘ │
│ │
│ Cache layers : │
│ ┌─────────────────────────────────────────────┐ │
│ │ L1: In-memory (exact match) → 1ms │ │
│ │ L2: Redis (semantic cache) → 5ms │ │
│ │ L3: Qdrant (vector search) → 10-50ms │ │
│ └─────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────┘
Implementation avec cache multi-niveaux
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 = {} # LRU cache in-memory async def query(self, question: str, user_id: str) -> dict: # L1 : Cache exact (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 : Cache semantique (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 : Pipeline RAG complet query_embedding = await self.embed_query(question) # Recherche vectorielle 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) # Generation LLM (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) } # Mettre en cache (TTL 1h) await self.redis.setex( f"rag:cache:{cache_key}", 3600, json.dumps(result) ) self.local_cache[cache_key] = result return result
Architecture Hybride : le meilleur des deux mondes
L'architecture hybride est la norme en production. Elle combine batch et temps reel avec un systeme d'events.
Diagramme d'architecture hybride
┌──────────────────────────────────────────────────────────────────┐
│ ARCHITECTURE HYBRIDE RAG │
├──────────────────────────────────────────────────────────────────┤
│ │
│ SOURCES EVENT BUS CONSUMERS │
│ ┌──────────┐ ┌──────────────┐ │
│ │ Upload │─────────>│ │──> [Batch Worker] │
│ │ API sync │─────────>│ Kafka / │ → Ingestion │
│ │ Webhook │─────────>│ Redis │ → Re-indexation │
│ │ Cron job │─────────>│ Streams │ → Analytics │
│ └──────────┘ │ │ │
│ │ │──> [Realtime Worker] │
│ UTILISATEURS │ │ → Query processing │
│ ┌──────────┐ │ │ → Live updates │
│ │ Widget │─────────>│ │ → Streaming │
│ │ API │─────────>│ │ │
│ │ Chat │─────────>│ │ │
│ └──────────┘ └──────────────┘ │
│ │
│ STOCKAGE │
│ ┌───────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ PostgreSQL│ │ Qdrant │ │ Redis │ │ S3 │ │
│ │ (metadata)│ │ (vectors)│ │ (cache) │ │ (files) │ │
│ └───────────┘ └──────────┘ └──────────┘ └──────────┘ │
└──────────────────────────────────────────────────────────────────┘
Event-driven document 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): """Publie un evenement sur le 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): """Souscrit a un type d'evenement.""" # Creer le groupe si inexistant try: await self.redis.xgroup_create( f"events:{event_type}", consumer_group, id="0", mkstream=True ) except redis.ResponseError: pass # Groupe existe deja 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 # ACK le message await self.redis.xack( f"events:{event_type}", consumer_group, msg_id ) # Publication d'evenements event_bus = EventBus() # Quand un document est uploade await event_bus.publish("document.uploaded", { "document_id": "doc_123", "source": "api_upload", "size_bytes": 45000 }) # Quand un document est modifie await event_bus.publish("document.updated", { "document_id": "doc_123", "changed_sections": ["section_2", "section_5"], "update_type": "partial" }) # Quand un document est supprime await event_bus.publish("document.deleted", { "document_id": "doc_123" })
Consumer : mise a jour incrementale
DEVELOPERpythonasync def incremental_update_consumer(): """Consumer qui met a jour les vecteurs incrementalement.""" async for event in event_bus.subscribe( "document.updated", "indexing_group", "worker_1" ): doc_id = event["data"]["document_id"] update_type = event["data"]["update_type"] if update_type == "partial": # Mise a jour partielle : re-indexer uniquement les sections modifiees changed_sections = event["data"]["changed_sections"] doc = await fetch_document(doc_id) for section_id in changed_sections: section_content = doc.get_section(section_id) chunks = semantic_chunking(section_content) embeddings = await embed_batch([c.text for c in chunks]) # Supprimer les anciens chunks de cette section await qdrant_client.delete( collection_name="documents", points_selector=FilterSelector( filter=Filter( must=[ FieldCondition( key="document_id", match=MatchValue(value=doc_id) ), FieldCondition( key="section_id", match=MatchValue(value=section_id) ) ] ) ) ) # Inserer les nouveaux chunks await qdrant_client.upsert("documents", [ PointStruct( id=uuid4().int >> 64, vector=emb, payload={ "document_id": doc_id, "section_id": section_id, "text": chunk.text } ) for chunk, emb in zip(chunks, embeddings) ]) elif update_type == "full": # Re-indexation complete du document await process_document.delay(doc_id)
Comparaison des systemes de queues
| Critere | Kafka | RabbitMQ | Redis Streams |
|---|---|---|---|
| Debit | 100K+ msg/s | 20K msg/s | 50K msg/s |
| Latence | 5-15ms | 1-5ms | < 1ms |
| Persistance | Disk (durable) | Memory + disk | Memory + AOF |
| Ordering | Par partition | Par queue | Par stream |
| Consumer groups | Oui | Oui | Oui |
| Replay | Oui (offset) | Non | Oui (ID) |
| Complexite ops | Elevee (cluster KRaft) | Moyenne | Faible |
| Ideal pour | Event sourcing, haute volumetrie | Taches async, routing | Cache + queue, faible latence |
| Recommandation RAG | 10K+ events/s | < 10K events/s | Meilleur choix pour RAG |
Pourquoi Redis Streams pour la plupart des RAG
Redis Streams offre le meilleur equilibre pour un systeme RAG typique :
- Deja la pour le cache : Redis sert de cache L2, pas de composant supplementaire
- Latence sub-milliseconde : parfait pour les mises a jour live
- Consumer groups : repartition de charge entre workers
- Replay : re-traitement possible en cas d'erreur
- Faible complexite : pas de cluster Kafka (controleurs KRaft, brokers) a gerer
Matrice de decision
Quand utiliser quoi ?
| Scenario | Approche | Justification |
|---|---|---|
| Ingestion initiale (1M+ docs) | Batch | Volume trop eleve pour le temps reel |
| Requete utilisateur | Temps reel | Latence critique |
| Document uploade par utilisateur | Hybride (event -> batch) | Ingestion en arriere-plan, notification quand pret |
| FAQ mise a jour | Temps reel | Changement petit, impact immediat |
| Changement de modele d'embedding | Batch | Re-indexation complete du corpus |
| Analytics / rapports | Batch | Pas de contrainte de latence |
| Chat en direct | Temps reel | Streaming obligatoire |
| Synchronisation CRM | Hybride (cron -> batch) | Sync periodique, volume moyen |
| Webhook e-commerce | Hybride (event -> realtime) | Mise a jour catalogue rapide |
Decision flowchart
La latence est-elle critique (<1s) ?
├── OUI → Le volume depasse 100 docs/min ?
│ ├── OUI → Architecture HYBRIDE
│ │ (queue + workers temps reel)
│ └── NON → Architecture TEMPS REEL pure
│ (API synchrone + cache)
└── NON → Le volume depasse 10K docs ?
├── OUI → Architecture BATCH pure
│ (Celery + workers distribues)
└── NON → Architecture BATCH simple
(cron job + script Python)
Comparaison des couts
Infrastructure mensuelle (10M documents, 100K requetes/jour)
| Composant | Batch seul | Temps reel seul | Hybride |
|---|---|---|---|
| Compute (workers) | $200 | $500 | $400 |
| Redis | $50 | $150 | $100 |
| Qdrant | $320 | $320 | $320 |
| Kafka/Queue | $0 | $0 | $80 |
| LLM (embeddings) | $150 | $300 | $200 |
| LLM (generation) | $0 | $2 000 | $2 000 |
| Total | $720 | $3 270 | $3 100 |
Cout par requete
| Approche | Cout/requete | Latence moyenne | Throughput |
|---|---|---|---|
| Batch (pre-calcule) | $0.001 | N/A (offline) | 10K+ docs/h |
| Temps reel (a la volee) | $0.025 | 1.5-3s | 100-1K req/s |
| Hybride (cache + live) | $0.008 | 0.5-2s | 500-5K req/s |
Le modele hybride reduit le cout par requete de 68% par rapport au temps reel pur grace au cache semantique.
Patterns avances
Pattern : Write-behind cache
Le write-behind cache permet de repondre immediatement tout en mettant a jour l'index en arriere-plan.
DEVELOPERpythonclass WriteBehindRAG: """Repond depuis le cache, met a jour l'index en background.""" async def update_document(self, doc_id: str, new_content: str): # 1. Mettre a jour le cache immediatement await self.redis.hset(f"doc:{doc_id}", "content", new_content) await self.redis.hset(f"doc:{doc_id}", "status", "pending_index") # 2. Publier l'evenement pour re-indexation await self.event_bus.publish("document.updated", { "document_id": doc_id, "update_type": "full" }) # 3. Les requetes utilisent le cache en attendant return {"status": "updated", "indexing": "in_progress"} async def query(self, question: str): # Recherche vectorielle normale results = await self.qdrant.search(question_embedding) # Enrichir avec le cache (documents recemment modifies) for result in results: cached = await self.redis.hgetall(f"doc:{result.id}") if cached and cached.get("status") == "pending_index": # Utiliser la version cachee (plus recente) result.text = cached["content"] return results
Pattern : Backpressure
Eviter de surcharger le systeme quand le volume explose.
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: """Verifie si la queue peut accepter plus de messages.""" queue_size = await self.redis.xlen(f"events:{queue_name}") if queue_size > self.max_queue_size: # Backpressure : rejeter avec 503 return False if queue_size > self.max_queue_size * 0.8: # Warning : ralentir await asyncio.sleep(0.1) return True async def ingest_with_backpressure(self, documents: list[dict]): """Ingestion avec controle de backpressure.""" 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 }
Notre architecture chez Ailog
Chez Ailog, nous utilisons une architecture hybride :
- Batch : ingestion de documents clients (Celery + Redis), re-indexation nocturne
- Temps reel : requetes widget avec cache semantique (Redis) + Qdrant
- Event bus : Redis Streams pour la coordination batch/realtime
- Cache multi-niveaux : L1 in-memory, L2 Redis, L3 Qdrant
Le resultat : latence p50 de 800ms pour une requete complete (retrieval + LLM), avec des documents indexes en moins de 30 secondes apres upload.
Decouvrez notre guide sur le streaming de reponses RAG et les strategies de caching.
FAQ
Conclusion
Le choix entre batch et temps reel n'est pas binaire. Les meilleurs systemes RAG utilisent une architecture hybride :
- Batch pour l'ingestion volumineuse et la re-indexation
- Temps reel pour les requetes utilisateur et les mises a jour critiques
- Event bus (Redis Streams) pour coordonner les deux
- Cache multi-niveaux pour reduire les couts et la latence
Commencez simple (batch + API synchrone), puis ajoutez la complexite quand le volume l'exige.
Vous voulez un systeme RAG hybride sans gerer l'infrastructure ? Testez Ailog - nous gerons le batch, le temps reel, le cache et l'event bus pour vous.
Tags
Articles connexes
Securite et Conformite RAG : RGPD, AI Act et bonnes pratiques
Securisez votre systeme RAG : conformite RGPD, AI Act europeen, protection des donnees et audit. Guide complet pour les entreprises.
RAG pour PME : Guide complet sans équipe data
Déployez un système RAG performant dans votre PME sans compétences techniques avancées : solutions no-code, budget maîtrisé et ROI rapide.
RAG Souverain : Hebergement France et donnees europeennes
Deployez un RAG souverain en France : hebergement local, conformite RGPD, alternatives aux GAFAM et bonnes pratiques pour les donnees europeennes.