RAG-Agenten: Orchestrierung von Multi-Agenten-Systemen
Konzipieren Sie RAG-basierte Multi-Agenten-Systeme: Orchestrierung, Spezialisierung, Zusammenarbeit und Fehlerbehandlung für komplexe Assistenten.
Agents RAG : Orchestrieren von Multi-Agenten-Systemen
Ein einfacher RAG-Agent folgt einer linearen Pipeline: retrieval, generation, antwort. Ein Multi-Agenten-System orchestriert mehrere spezialisierte Agents, die zusammenarbeiten, um komplexe Aufgaben zu bearbeiten. Dieser Leitfaden untersucht Architekturen und Patterns zum Aufbau solcher Systeme.
Warum Multi-Agenten?
Grenzen des monolithischen RAG
Ein klassischer RAG stößt bei bestimmten Aufgaben an Grenzen:
| Aufgabe | Probleme bei einfachem RAG |
|---|---|
| Fragen mit mehreren Quellen | Ein einziges retrieval, begrenzter Kontext |
| Komplexes Reasoning | Keine Problemzerlegung |
| Mehrere Aktionen | Lineare Pipeline, keine Schleife |
| Faktenüberprüfung | Kein Double-Check |
| Lange Aufgaben | Keine Planung |
Der Multi-Agenten-Ansatz
Teilen, um besser zu herrschen:
Komplexe Frage
│
▼
┌─────────────────┐
│ ORCHESTRATOR │ ← Zerlegt und koordiniert
└────────┬────────┘
│
┌────┼────┬────────────┐
▼ ▼ ▼ ▼
┌──────┐┌──────┐┌──────┐┌──────┐
│Agent ││Agent ││Agent ││Agent │
│Search││Reason││Verify││Action│
└──────┘└──────┘└──────┘└──────┘
│ │ │ │
└────┴────┴────────────┘
│
▼
Finale Antwort
Orchestrierungsarchitekturen
Pattern 1 : Routeur
Der Orchestrator routet an den spezialisierten Agenten.
DEVELOPERpythonfrom enum import Enum from typing import Callable, Dict class AgentType(Enum): FAQ = "faq" TECHNICAL = "technical" SALES = "sales" ESCALATION = "escalation" class RouterOrchestrator: def __init__(self, agents: Dict[AgentType, Callable], classifier): self.agents = agents self.classifier = classifier async def process(self, query: str, context: dict = None) -> dict: """ Leitet die Anfrage an den passenden Agenten weiter """ # 1. Die Anfrage klassifizieren classification = await self.classifier.classify(query) agent_type = AgentType(classification["agent"]) # 2. Agenten aufrufen agent = self.agents.get(agent_type) if not agent: agent = self.agents[AgentType.ESCALATION] result = await agent(query, context) return { "response": result["answer"], "agent_used": agent_type.value, "confidence": classification["confidence"], "sources": result.get("sources", []) } class IntentClassifier: def __init__(self, llm): self.llm = llm async def classify(self, query: str) -> dict: prompt = f""" Klassifiziere diese Nutzeranfrage. Kategorien : - faq : Allgemeine Fragen, Richtlinien, Informationen - technical : Technische Probleme, Bugs, Konfigurationen - sales : Preise, Käufe, Abonnements - escalation : Beschwerden, Notfälle, sensible Anfragen Anfrage : {query} Antworte in JSON : {{"agent": "...", "confidence": 0.0-1.0}} """ result = await self.llm.generate(prompt, temperature=0) return self._parse_json(result)
Pattern 2 : Sequenzielle Pipeline
Die Agents werden nacheinander ausgeführt, jeder erweitert den Kontext.
DEVELOPERpythonfrom dataclasses import dataclass from typing import List @dataclass class AgentResult: agent_name: str output: dict success: bool error: str = None class SequentialPipeline: def __init__(self, agents: List[tuple]): """ agents: Liste von (Name, Agent-Callable) """ self.agents = agents async def process(self, query: str, initial_context: dict = None) -> dict: """ Führt die Agenten sequenziell aus """ context = initial_context or {} context["original_query"] = query results = [] for agent_name, agent in self.agents: try: result = await agent(query, context) agent_result = AgentResult( agent_name=agent_name, output=result, success=True ) # Kontext für den nächsten Agenten anreichern context[f"{agent_name}_result"] = result context["last_result"] = result except Exception as e: agent_result = AgentResult( agent_name=agent_name, output={}, success=False, error=str(e) ) results.append(agent_result) # Abbrechen, wenn ein Agent fehlschlägt (optional) if not agent_result.success and self.stop_on_failure: break return { "final_result": context.get("last_result"), "pipeline_results": results, "context": context } # Beispiel-Pipeline pipeline = SequentialPipeline([ ("query_analyzer", QueryAnalyzerAgent()), ("retriever", RetrievalAgent()), ("fact_checker", FactCheckAgent()), ("generator", GenerationAgent()), ("citation_adder", CitationAgent()) ])
Pattern 3 : Parallel mit Fusion
Mehrere Agents arbeiten parallel, danach werden die Ergebnisse zusammengeführt.
DEVELOPERpythonimport asyncio from typing import List, Callable class ParallelOrchestrator: def __init__( self, agents: List[tuple], fusion_agent: Callable ): self.agents = agents # (name, agent, weight) self.fusion = fusion_agent async def process(self, query: str, context: dict = None) -> dict: """ Führt die Agenten parallel aus und fusioniert danach """ # Alle Agenten parallel starten tasks = [ self._run_agent(name, agent, query, context) for name, agent, _ in self.agents ] results = await asyncio.gather(*tasks, return_exceptions=True) # Gültige Ergebnisse sammeln valid_results = [] for i, result in enumerate(results): name, _, weight = self.agents[i] if not isinstance(result, Exception): valid_results.append({ "agent": name, "result": result, "weight": weight }) # Ergebnisse fusionieren fused = await self.fusion(query, valid_results) return { "response": fused["answer"], "contributing_agents": [r["agent"] for r in valid_results], "fusion_confidence": fused.get("confidence") } async def _run_agent(self, name: str, agent: Callable, query: str, context: dict): try: return await asyncio.wait_for( agent(query, context), timeout=10.0 ) except asyncio.TimeoutError: return {"error": "timeout", "agent": name} class FusionAgent: def __init__(self, llm): self.llm = llm async def __call__(self, query: str, results: List[dict]) -> dict: """ Fusioniert die Antworten mehrerer Agenten """ responses_text = "\n\n".join([ f"Agent {r['agent']} (Gewicht {r['weight']}):\n{r['result'].get('answer', 'N/A')}" for r in results ]) prompt = f""" Du musst diese Antworten verschiedener Agenten zu einer einzigen kohärenten Antwort synthetisieren. Ursprüngliche Frage : {query} Antworten der Agenten : {responses_text} Regeln : 1. Priorisiere die Agenten mit dem höchsten Gewicht 2. Bei Widersprüchen beide Sichtweisen aufzeigen 3. Quellen zitieren, sofern verfügbar 4. Eine einzige, kohärente Antwort erzeugen Synthetisierte Antwort : """ answer = await self.llm.generate(prompt, temperature=0.3) return { "answer": answer, "confidence": self._calculate_confidence(results) }
Pattern 4 : ReAct (Reasoning + Acting)
Der Agent denkt nach, handelt, beobachtet und iteriert.
DEVELOPERpythonfrom typing import Dict, Any class ReActAgent: def __init__(self, llm, tools: Dict[str, Callable], max_iterations: int = 5): self.llm = llm self.tools = tools self.max_iterations = max_iterations async def process(self, query: str) -> dict: """ Führt das ReAct-Pattern aus: Thought -> Action -> Observation """ history = [] final_answer = None for i in range(self.max_iterations): # Den nächsten Gedanken/die nächste Aktion erzeugen step = await self._generate_step(query, history) if step["type"] == "thought": history.append({"thought": step["content"]}) elif step["type"] == "action": # Aktion ausführen tool_name = step["tool"] tool_input = step["input"] if tool_name in self.tools: observation = await self.tools[tool_name](tool_input) else: observation = f"Werkzeug '{tool_name}' nicht verfügbar" history.append({ "action": f"{tool_name}({tool_input})", "observation": observation }) elif step["type"] == "answer": final_answer = step["content"] break return { "answer": final_answer, "reasoning_trace": history, "iterations": i + 1 } async def _generate_step(self, query: str, history: list) -> dict: """ Erzeugt den nächsten Schritt (Gedanke, Aktion oder Antwort) """ history_text = self._format_history(history) tools_desc = self._format_tools() prompt = f""" Du bist ein Agent, der Probleme Schritt für Schritt löst. Frage : {query} Verfügbare Werkzeuge : {tools_desc} Verlauf : {history_text} Nächster Schritt (wähle EIN Format) : Gedanke : [dein Reasoning] ODER Aktion : [werkzeug_name]([parameter]) ODER Antwort : [deine finale Antwort] """ result = await self.llm.generate(prompt, temperature=0) return self._parse_step(result) # Werkzeuge für den ReAct-Agenten tools = { "search_kb": lambda q: rag_search(q), "calculate": lambda expr: eval(expr), "get_current_date": lambda _: datetime.now().isoformat(), "check_inventory": lambda product_id: inventory_api.check(product_id) }
Spezialisierte Agents
Search Agent
DEVELOPERpythonclass SearchAgent: def __init__(self, retrievers: Dict[str, Retriever], llm): self.retrievers = retrievers self.llm = llm async def __call__(self, query: str, context: dict) -> dict: """ Agent, spezialisiert auf die Suche über mehrere Quellen """ # 1. Relevante Quellen bestimmen sources = await self._select_sources(query) # 2. In jeder Quelle suchen all_results = [] for source_name in sources: retriever = self.retrievers.get(source_name) if retriever: results = await retriever.search(query) for r in results: r["source"] = source_name all_results.extend(results) # 3. Ergebnisse reranken ranked = await self._rerank(query, all_results) return { "documents": ranked[:10], "sources_used": sources, "total_found": len(all_results) } async def _select_sources(self, query: str) -> List[str]: """ Wählt die für die Anfrage relevanten Quellen aus """ prompt = f""" Welche Quellen sind für diese Frage relevant ? Verfügbare Quellen : - kb_general : Allgemeine Wissensdatenbank - kb_technical : Technische Dokumentation - kb_support : Support-Ticket-Historie - kb_products : Produktkatalog Frage : {query} Relevante Quellen (JSON-Liste) : """ result = await self.llm.generate(prompt, temperature=0) return self._parse_json(result)
Verification Agent
DEVELOPERpythonclass VerificationAgent: def __init__(self, llm, fact_db): self.llm = llm self.fact_db = fact_db async def __call__(self, query: str, context: dict) -> dict: """ Überprüft Aussagen anhand vertrauenswürdiger Quellen """ # Die zu überprüfende Antwort abrufen answer = context.get("last_result", {}).get("answer", "") # 1. Aussagen extrahieren claims = await self._extract_claims(answer) # 2. Jede Aussage überprüfen verifications = [] for claim in claims: verification = await self._verify_claim(claim) verifications.append(verification) # 3. Gesamtscore bestimmen verified_count = sum(1 for v in verifications if v["verified"]) total = len(verifications) return { "claims": verifications, "verification_score": verified_count / total if total > 0 else 1.0, "needs_correction": any(not v["verified"] for v in verifications) } async def _verify_claim(self, claim: str) -> dict: """ Überprüft eine spezifische Aussage """ # Widersprüchliche oder bestätigende Fakten suchen facts = await self.fact_db.search(claim, top_k=3) prompt = f""" Überprüfe diese Aussage anhand der folgenden Fakten. Aussage : {claim} Referenzfakten : {self._format_facts(facts)} Ist die Aussage : - verified : durch die Fakten gestützt - contradicted : durch die Fakten widerlegt - unverified : nicht überprüfbar Antworte in JSON mit "status" und "evidence". """ result = await self.llm.generate(prompt, temperature=0) parsed = self._parse_json(result) return { "claim": claim, "verified": parsed.get("status") == "verified", "status": parsed.get("status"), "evidence": parsed.get("evidence") }
Action Agent
DEVELOPERpythonclass ActionAgent: def __init__(self, action_registry: Dict[str, Callable], llm): self.actions = action_registry self.llm = llm async def __call__(self, query: str, context: dict) -> dict: """ Führt Aktionen basierend auf der Anfrage aus """ # 1. Notwendige Aktionen bestimmen action_plan = await self._plan_actions(query, context) # 2. Aktionen ausführen results = [] for action in action_plan: if action["name"] in self.actions: try: result = await self.actions[action["name"]](**action["params"]) results.append({ "action": action["name"], "success": True, "result": result }) except Exception as e: results.append({ "action": action["name"], "success": False, "error": str(e) }) return { "actions_executed": results, "all_successful": all(r["success"] for r in results) } async def _plan_actions(self, query: str, context: dict) -> List[dict]: """ Plant die auszuführenden Aktionen """ actions_desc = "\n".join([ f"- {name}: {func.__doc__ or 'No description'}" for name, func in self.actions.items() ]) prompt = f""" Plane die notwendigen Aktionen, um diese Anfrage zu erfüllen. Anfrage : {query} Kontext : {context} Verfügbare Aktionen : {actions_desc} Aktionsplan (JSON-Array) : [ {{"name": "action_name", "params": {{}}}}, ... ] """ result = await self.llm.generate(prompt, temperature=0) return self._parse_json(result) # Aktionsregister action_registry = { "create_ticket": create_support_ticket, "send_email": send_email, "update_order": update_order_status, "schedule_callback": schedule_callback, "apply_discount": apply_discount_code }
Fehlerbehandlung
Retry mit Backoff
DEVELOPERpythonimport asyncio from functools import wraps def with_retry(max_attempts: int = 3, backoff: float = 1.0): def decorator(func): @wraps(func) async def wrapper(*args, **kwargs): last_exception = None for attempt in range(max_attempts): try: return await func(*args, **kwargs) except Exception as e: last_exception = e if attempt < max_attempts - 1: wait_time = backoff * (2 ** attempt) await asyncio.sleep(wait_time) raise last_exception return wrapper return decorator
Sanfter Fallback
DEVELOPERpythonclass ResilientOrchestrator: def __init__(self, primary_agents, fallback_agent): self.primary = primary_agents self.fallback = fallback_agent async def process(self, query: str, context: dict = None) -> dict: """ Führt mit Fallback im Fehlerfall aus """ try: # Haupt-Pipeline versuchen result = await self._run_primary(query, context) if self._is_valid_result(result): return result except Exception as e: logger.warning(f"Primary pipeline failed: {e}") # Fallback zum einfachen Agenten return await self.fallback(query, context) def _is_valid_result(self, result: dict) -> bool: """ Prüft, ob das Ergebnis gültig ist """ return ( result.get("answer") and len(result["answer"]) > 10 and result.get("confidence", 0) > 0.5 )
Monitoring und Observability
DEVELOPERpythonimport time from dataclasses import dataclass @dataclass class AgentTrace: agent_name: str start_time: float end_time: float input_query: str output: dict success: bool error: str = None class TracingOrchestrator: def __init__(self, orchestrator, trace_store): self.orchestrator = orchestrator self.traces = trace_store async def process(self, query: str, context: dict = None) -> dict: """ Führt mit vollständigem Tracing aus """ trace_id = str(uuid.uuid4()) traces = [] # Wrapper zum Erfassen der Traces original_agents = self.orchestrator.agents.copy() for name, agent in original_agents.items(): self.orchestrator.agents[name] = self._wrap_agent(name, agent, traces) try: result = await self.orchestrator.process(query, context) # Traces speichern await self.traces.save(trace_id, traces) result["trace_id"] = trace_id result["execution_time_ms"] = sum( (t.end_time - t.start_time) * 1000 for t in traces ) return result finally: # Ursprüngliche Agenten wiederherstellen self.orchestrator.agents = original_agents def _wrap_agent(self, name: str, agent, traces: list): async def wrapped(query, context): start = time.time() try: result = await agent(query, context) traces.append(AgentTrace( agent_name=name, start_time=start, end_time=time.time(), input_query=query, output=result, success=True )) return result except Exception as e: traces.append(AgentTrace( agent_name=name, start_time=start, end_time=time.time(), input_query=query, output={}, success=False, error=str(e) )) raise return wrapped
Best Practices
1. Atomare Agents
Jeder Agent sollte eine klare, einzelne Verantwortung haben.
2. Timeouts überall
Setzen Sie Timeouts, um Blockierungen zu vermeiden.
3. Strukturierte Logs
Loggen Sie jeden Schritt für Debugging.
4. Unittests pro Agent
Testen Sie jeden Agent unabhängig vor der Integration.
5. Circuit Breaker
Temporäres Deaktivieren fehlerhafter Agents.
Weiterführende Links
- Konversationelles RAG - Gedächtnis und Kontext
- LLM-Generierung - Antworten optimieren
- RAG-Evaluation - Qualität messen
FAQ
Multi-Agenten-Orchestrierung mit Ailog
Der Aufbau einer robusten Multi-Agenten-Architektur erfordert tiefgehende Expertise. Mit Ailog profitieren Sie von einer vorkonfigurierten Agenten-Infrastruktur:
- Spezialisierte Agents : Suche, Verifikation, Aktion
- Konfigurierbare Orchestrierung ohne Code
- Integriertes Monitoring mit vollständigen Traces
- Automatische Fallbacks und Fehlerbehandlung
- Skalierbarkeit für hohe Lasten
Testen Sie Ailog kostenlos und stellen Sie Multi-Agenten-Assistenten in wenigen Klicks bereit.
Tags
Verwandte Artikel
AutoGen: Multiagentensysteme für RAG
Umfassender Leitfaden zum Aufbau von RAG-Multiagentensystemen mit Microsoft AutoGen. Konversationen zwischen Agenten, Orchestrierung und fortgeschrittene Anwendungsfälle.
Function calling : RAG mit Aktionen
Vollständiger Leitfaden zum Kombinieren von RAG und function calling: Agenten, die recherchieren UND handeln, Integration externer APIs, automatisierte Aktionen und interaktive Workflows.
CrewAI: Spezialisierte Teams von RAG-Agenten
Umfassender Leitfaden zum Aufbau von Teams aus RAG-Agenten mit CrewAI: Rollen, Aufgaben, Delegation, personalisierte Tools und Zusammenarbeit zwischen spezialisierten Agenten.