详述生产环境中轮询模式的四大失效场景,并给出用Redis Streams和幂等Worker构建低延迟响应系统的架构方案。
如今运行在生产环境中的大多数 AI Agent,最初都是简单的轮询循环。一个后台 worker 定期查询数据库或第三方 API:
# 这个朴素的反模式正运行在数百个生产服务中
while True:
records = db.query("SELECT * FROM invoices WHERE status = 'pending_review'")
for record in records:
agent.process(record)
time.sleep(30)
在真实生产流量下,这个模式会在四种截然不同的故障模式下崩溃:
尾延迟瓶颈:一个紧急事件在轮询周期结束后 100 毫秒到达,却要等待 29.9 秒才能被处理。对于时间敏感的工作流(如安全修复、实时客户路由或自动化交易),这会带来不可接受的延迟。
级联 API 速率限制:当你从 10 个扩展到 1,000 个 Agent,每个 Agent 独立检查下游工具(如 Salesforce、GitHub、Slack)时,你的基础设施每分钟发出数万次空 GET 请求。在实际工作负载执行之前,上游系统就会对你的 IP 地址进行限流或封禁。
计算资源浪费与内存压力:运行数千个阻塞的 Python 线程或事件循环,在空闲内存中持有数据库连接池,会降低节点性能并推高基础设施成本。
竞态条件和脑裂执行:水平扩展轮询 worker 时,如果不安复杂的分布式锁(SELECT FOR UPDATE SKIP LOCKED),两个实例经常会同时获取同一条记录。这会导致重复的 LLM 调用、重复付款或状态损坏。
事件驱动的 Agent 架构将事件摄入与 Agent 认知解耦。Agent 不再主动询问是否有工作可做,而是通过富化的事件载荷由环境通知 Agent。
+------------------+ +------------------+ +------------------+
| Event Producer | | Event Producer | | Event Producer |
| (API / Webhooks) | | (CDC / Postgres) | | (IoT / Sensors) |
+--------+---------+ +--------+---------+ +--------+---------+
| | |
+-------------------------+-------------------------+
|
v
+---------------------------------------------+
| Message Broker (Redis Streams / Kafka) |
| Stream: `agent:events:v1` |
+---------------------+-----------------------+
|
+--------------+--------------+
| |
v v
+-------------------------+ +-------------------------+
| Consumer Group: Worker 1| | Consumer Group: Worker 2|
+------------+------------+ +------------+------------+
| |
v v
+-------------------------+ +-------------------------+
| Idempotency Check | | Idempotency Check |
| (Redis SET key EX NX) | | (Redis SET key EX NX) |
+------------+------------+ +------------+------------+
| |
v v
+-------------------------+ +-------------------------+
| Agent Execution Engine | | Agent Execution Engine |
| (LangGraph / DSPy / LLM)| | (LangGraph / DSPy / LLM)|
+------------+------------+ +------------+------------+
| |
+--------------+--------------+
|
v
+---------------------------------------------+
| State Checkpoints / Write-Back (Postgres) |
+---------------------------------------------+
Event Producers(事件生产者):摄入点(FastAPI webhooks、Change Data Capture 流水线、内部系统事件)发布严格类型化的模式,包含变更增量及必要的上下文状态。
Streaming Ingestion & Consumer Groups(流式摄入与消费者组):Redis Streams 维护一个仅追加的日志,具备持久化的消费者组。如果 worker 进程在推理中途崩溃,消息保持未确认状态(XACK),并通过死信/待处理机制重新分配给健康的 worker。
确定性幂等性门控:由于网络传输保证至少一次投递,每个事件必须在 Agent 初始化其上下文窗口之前,通过基于确定性事件哈希的原子锁门控。
Agent Execution Worker(Agent 执行 Worker):Worker 解析预填充的事件载荷,运行推理链(LangGraph、原始工具调用或自定义状态机),更新持久化存储,并确认消息处理。
以下是完整、可运行的生产级 Python 实现。它使用 redis-py 配合 Redis Streams 和消费者组来管理有状态的、幂等的、事件驱动的 Agent 调用。
import asyncio
import hashlib
import json
import logging
import os
import sys
from typing import Any, Callable, Dict, Optional
from pydantic import BaseModel, Field
import redis.asyncio as redis
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] (%(name)s) %(message)s",
handlers=[logging.StreamHandler(sys.stdout)],
)
logger = logging.getLogger("EventDrivenAgent")
# ---------------------------------------------------------------------------
# 1. SCHEMAS
# ---------------------------------------------------------------------------
class AgentEvent(BaseModel):
event_id: str
event_type: str
source: str
payload: Dict[str, Any]
timestamp: int
def generate_idempotency_key(self) -> str:
"""Generates a deterministic hash based on event payload and type."""
raw_signature = f"{self.event_type}:{self.event_id}:{json.dumps(self.payload, sort_keys=True)}"
return f"idemp:{hashlib.sha256(raw_signature.encode()).hexdigest()}"
# ---------------------------------------------------------------------------
# 2. CORE EVENT CONSUMER & AGENT RUNNER
# ---------------------------------------------------------------------------
class EventDrivenAgentWorker:
def __init__(
self,
redis_url: str,
stream_name: str,
group_name: str,
consumer_name: str,
lock_ttl_seconds: int = 300,
):
self.redis_url = redis_url
self.stream_name = stream_name
self.group_name = group_name
self.consumer_name = consumer_name
self.lock_ttl = lock_ttl_seconds
self.redis_client: Optional[redis.Redis] = None
self._running = False
async def connect(self):
"""Initializes the Redis connection and sets up consumer groups."""
self.redis_client = redis.from_url(self.redis_url, decode_responses=True)
try:
# Create the stream consumer group if it doesn't already exist
await self.redis_client.xgroup_create(
name=self.stream_name,
groupname=self.group_name,
id="0",
mkstream=True,
)
logger.info(f"Consumer group '{self.group_name}' initialized on stream '{self.stream_name}'.")
except redis.exceptions.ResponseError as e:
if "BUSYGROUP" in str(e):
logger.info(f"Consumer group '{self.group_name}' already exists.")
else:
raise e
async def acquire_idempotency_lock(self, key: str) -> bool:
"""
Uses Redis SET key NX to prevent duplicate execution of the same event.
Returns True if lock was acquired, False if the event was already processed.
"""
is_new = await self.redis_client.set(key, "PROCESSING", ex=self.lock_ttl, nx=True)
return bool(is_new)
async def process_stream(self, agent_handler: Callable[[AgentEvent], Any]):
"""Continuous event consumption loop."""
self._running = True
logger.info(f"Worker {self.consumer_name} listening for events...")
while self._running:
try:
# Read new messages assigned to this consumer group
# Block for 2000ms if no messages exist
response = await self.redis_client.xreadgroup(
groupname=self.group_name,
consumername=self.consumer_name,
streams={self.stream_name: ">"},
count=10,
block=2000,
)
if not response:
await asyncio.sleep(0.01)
continue
for stream, messages in response:
for message_id, raw_data in messages:
await self._handle_single_message(message_id, raw_data, agent_handler)
except asyncio.CancelledError:
logger.info("Worker shutdown initiated.")