从请求-响应到事件驱动

传统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开始,当痛点出现时再引入更高层级的方案。