从请求-响应到事件驱动
传统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)提供解耦、缓冲和可靠投递。
Agent + MQ架构
import asyncio
import json
from aiokafka import AIOKafkaConsumer, AIOKafkaProducer
class EventDrivenAgent:
def __init__(self, agent_id: str,
input_topic: str,
output_topic: str,
kafka_servers: str):
self.agent_id = agent_id
self.input_topic = input_topic
self.output_topic = output_topic
self.kafka_servers = kafka_servers
self.consumer = None
self.producer = None
self._running = False
async def start(self):
self.consumer = AIOKafkaConsumer(
self.input_topic,
bootstrap_servers=self.kafka_servers,
group_id=f"agent-{self.agent_id}",
enable_auto_commit=False,
value_deserializer=lambda v: json.loads(v.decode())
)
self.producer = AIOKafkaProducer(
bootstrap_servers=self.kafka_servers,
value_serializer=lambda v: json.dumps(v).encode()
)
await self.consumer.start()
await self.producer.start()
self._running = True
await self._consume_loop()
async def _consume_loop(self):
while self._running:
try:
batch = await self.consumer.getmany(
timeout_ms=1000, max_records=10
)
for topic_partition, messages in batch.items():
tasks = [self._process_message(msg) for msg in messages]
results = await asyncio.gather(*tasks, return_exceptions=True)
# 批量提交offset
if all(not isinstance(r, Exception) for r in results):
await self.consumer.commit()
else:
# 有失败,不提交,重试
for r in results:
if isinstance(r, Exception):
logger.error(f"处理失败: {r}")
except Exception as e:
logger.error(f"消费循环异常: {e}")
async def _process_message(self, msg):
event = msg.value
result = await self._handle_event(event)
# 发布结果到输出topic
if result and self.output_topic:
await self.producer.send_and_wait(
self.output_topic,
value={
"source": self.agent_id,
"event_id": event.get("event_id"),
"result": result,
"timestamp": datetime.now().isoformat()
}
)
async def _handle_event(self, event: dict) -> dict:
"""子类实现具体处理逻辑"""
raise NotImplementedError
async def stop(self):
self._running = False
await self.consumer.stop()
await self.producer.stop()
事件Schema设计
from pydantic import BaseModel, Field
from datetime import datetime
from enum import Enum
import uuid
class EventPriority(Enum):
LOW = "low"
NORMAL = "normal"
HIGH = "high"
CRITICAL = "critical"
class AgentEvent(BaseModel):
event_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
event_type: str
source: str # 发起方Agent ID
target: str | None = None # 目标Agent ID(None=广播)
payload: dict
priority: EventPriority = EventPriority.NORMAL
timestamp: str = Field(default_factory=lambda: datetime.now().isoformat())
trace_id: str # 链路追踪ID
retry_count: int = 0
max_retries: int = 3
class CodeReviewRequested(AgentEvent):
event_type: str = "code_review.requested"
payload: dict # {"repo": str, "commit_sha": str, "files": list[str]}
class ReviewCompleted(AgentEvent):
event_type: str = "code_review.completed"
payload: dict # {"issues": list, "approved": bool, "comments": str}
第三层:事件流处理
事件溯源与CEP
对于复杂的多Agent协作,需要复杂事件处理(CEP)——在事件流中检测模式:
from typing import Callable
class EventPatternMatcher:
"""检测事件流中的模式,触发新事件"""
def __init__(self):
self.patterns: list[dict] = []
self.event_buffer: list[dict] = []
def add_pattern(self, name: str, conditions: list[dict],
action: Callable, window_ms: int = 60000):
self.patterns.append({
"name": name,
"conditions": conditions,
"action": action,
"window_ms": window_ms
})
async def on_event(self, event: dict):
self.event_buffer.append(event)
self._prune_buffer()
for pattern in self.patterns:
if await self._match_pattern(pattern):
await pattern["action"](self.event_buffer[-len(pattern["conditions"]):])
# 匹配后清除相关事件,防止重复触发
self.event_buffer = self.event_buffer[-1:]
async def _match_pattern(self, pattern: dict) -> bool:
conditions = pattern["conditions"]
if len(self.event_buffer) < len(conditions):
return False
relevant = self.event_buffer[-len(conditions):]
for i, cond in enumerate(conditions):
event = relevant[i]
if not self._check_condition(event, cond):
return False
return True
def _check_condition(self, event: dict, cond: dict) -> bool:
if cond.get("event_type") and event.get("event_type") != cond["event_type"]:
return False
if cond.get("source") and event.get("source") != cond["source"]:
return False
return True
# 示例:代码提交 + 测试通过 → 触发部署审查
matcher = EventPatternMatcher()
async def trigger_deploy(relevant_events):
print(f"触发部署审查: {[e['event_type'] for e in relevant_events]}")
matcher.add_pattern(
name="auto_deploy_check",
conditions=[
{"event_type": "code_pushed"},
{"event_type": "tests_passed"}
],
action=trigger_deploy,
window_ms=300000 # 5分钟窗口
)
架构选型对比
| 维度 | Webhook | 消息队列 | 事件流(CEP) |
|---|---|---|---|
| 复杂度 | 低 | 中 | 高 |
| 吞吐量 | 低 | 高 | 高 |
| 可靠性 | 低 | 高 | 高 |
| 延迟 | 最低 | 低-中 | 中 |
| 顺序保证 | 无 | 分区级 | 时间窗口 |
| 适用规模 | 单Agent | 多Agent | 多Agent+模式检测 |
演进路线
Webhook → 消息队列 → 事件流
建议路径:
1. MVP阶段: Webhook + 异步处理
2. 多Agent协作: 引入Kafka/RabbitMQ
3. 复杂编排: 加CEP引擎(Python流处理或Flink)
总结
事件驱动架构让Agent系统从"被动响应"变为"主动感知"。三层演进——Webhook(简单通知)→ 消息队列(可靠解耦)→ 事件流(模式检测)——对应着系统复杂度的不同阶段。核心设计原则:
- 事件不可变:发布后不可修改,便于追溯
- 幂等处理:同一事件被消费多次结果一致
- 背压控制:消费速度跟不上生产速度时优雅降级
- 全链路追踪:trace_id串联整个事件链路
不要过度设计。从Webhook开始,当痛点出现时再引入更高层级的方案。