事件驱动的Agent架构:从Webhook到消息队列
从请求-响应到事件驱动 传统LLM应用采用请求-响应模式:用户发消息→LLM处理→返回结果。但在生产环境中,Agent需要处理大量异步、多源、解耦的场景: 用户上传文件后自动触发分析Agent 多个Agent完成各自任务后汇总 外部系统(GitHub、Jira)状态变更触发Agent响应 这些场景的共同特征是:事件发生的时间不确定,处理者可能不唯一。事件驱动架构(EDA)是解决这类问题的自然选择。 第一层:Webhook模式 适用场景 Webhook是最简单的EDA实现,适合外部系统→Agent的单向通知。 from fastapi import FastAPI, Request, HTTPException import hmac import hashlib import asyncio app = FastAPI() class WebhookHandler: def __init__(self): self.handlers: dict[str, callable] = {} self.secrets: dict[str, str] = {} def register(self, event_type: str, handler: callable, secret: str = ""): self.handlers[event_type] = handler if secret: self.secrets[event_type] = secret async def handle(self, event_type: str, payload: dict, signature: str = ""): # 验证签名 if event_type in self.secrets: if not self._verify_signature(payload, signature, self.secrets[event_type]): raise HTTPException(401, "Invalid signature") handler = self.handlers.get(event_type) if not handler: raise HTTPException(404, f"No handler for {event_type}") # 异步处理,不阻塞响应 asyncio.create_task(handler(payload)) def _verify_signature(self, payload: dict, signature: str, secret: str) -> bool: expected = hmac.new( secret.encode(), str(payload).encode(), hashlib.sha256 ).hexdigest() return hmac.compare_digest(expected, signature) webhook_handler = WebhookHandler() @app.post("/webhook/{event_type}") async def webhook(event_type: str, request: Request): payload = await request.json() signature = request.headers.get("X-Signature", "") await webhook_handler.handle(event_type, payload, signature) return {"status": "accepted"} # 注册处理器 @webhook_handler.register("github.push", secret=os.getenv("GITHUB_WEBHOOK_SECRET")) async def handle_push(payload: dict): changes = payload.get("commits", []) # 触发代码审查Agent await code_review_agent.analyze(changes) Webhook的局限 局限 影响 解决方案 无重试机制 网络抖动导致丢失 上游需重试 + 幂等设计 无背压 高峰期Agent过载 加消息队列缓冲 同步确认 处理慢导致超时 异步处理 + 立即返回 无顺序保证 事件乱序处理 加序列号/时间戳 第二层:消息队列模式 为什么需要消息队列 当事件量增大或处理时间变长时,Webhook模式会崩溃。消息队列(MQ)提供解耦、缓冲和可靠投递。 ...