9.2 高并发处理
在上一节中,我们搭建了 Agent 系统的四层架构骨架,并讨论了同步与异步、有状态与无状态的选型。但架构图画的再漂亮,真正决定系统生死的是一个更接地气的问题:当一千个用户同时提问时,你的系统会怎样?
本节就来回答这个问题。我们将从高并发的特殊挑战讲起,深入请求队列、连接池、限流算法和异步任务管道的实现细节。
9.2.1 Agent 高并发的特殊挑战
与传统 Web 应用不同,Agent 系统的高并发面临独特的挑战。打个生活中的比方:传统 Web 应用的并发就像超市收银——每个顾客结账只要几秒钟,即便排队也很快轮到。而 Agent 系统的并发更像医院挂号——每个患者可能要看十几分钟,队伍一长就容易堵死。
挑战一:LLM API 的长尾延迟
LLM API 的响应时间非常不稳定。同一个请求,可能在 500ms 返回,也可能需要 30 秒。P99 延迟往往是 P50 的 10~50 倍。这种长尾效应意味着即使平均并发不高,也会间歇性地出现大量请求堆积。
生活类比:就像高速公路上的收费站,大部分车几秒就过去了,但偶尔有一辆大卡车要花两分钟,后面就排起了长龙。
挑战二:有状态会话的亲和性
有状态 Agent 的会话数据通常存储在特定节点上。如果负载均衡器将同一会话的请求路由到不同节点,会导致上下文丢失。这需要"会话亲和性(Sticky Session)"的支持。
挑战三:Token 消耗的突发性
用户行为具有很强的突发性。一个热门话题可能导致并发请求瞬间飙升,而 LLM API 的 Rate Limit 往往比自建服务更严格。
挑战四:工具调用的级联放大
一个用户请求可能触发 Agent 的多次工具调用,每次工具调用又可能触发新的 LLM 调用。这种级联放大效应会使实际并发远超表面请求数。表面上 100 个并发用户,实际可能产生 300~500 个并发 LLM 调用。
9.2.2 请求队列与连接池
请求队列设计
请求队列是高并发系统的第一道防线。它的核心作用是将瞬时高峰"削峰填谷",平滑后端的处理压力。
生活类比:请求队列就像银行取号机。不管你多着急,先拿一个号排着,窗口按顺序叫号。即使突然涌入 50 个人,系统也不会崩——只是号排得长一点。
┌──────────┐ ┌──────────────┐ ┌──────────────┐
│ Client │────▶│ 请求队列 │────▶│ Worker Pool │
│ Requests│ │ (Redis/RMQ) │ │ (Celery) │
└──────────┘ └──────────────┘ └──────────────┘
│
▼
┌──────────────┐
│ 死信队列 │
│ (DLQ) │
└──────────────┘队列深度与超时配置:
# 基于 Redis 的请求队列实现
import asyncio # 导入异步框架
import redis.asyncio as aioredis # 导入异步 Redis 客户端
from typing import Optional # 导入可选类型注解
import json # 导入 JSON 处理
import time # 导入时间模块
class RequestQueue:
"""基于 Redis List 的请求队列"""
def __init__(self, redis_url: str = "redis://localhost"):
self.redis = aioredis.from_url(redis_url) # 创建 Redis 连接
self.queue_key = "agent:requests" # 请求队列的 Redis Key
self.dlq_key = "agent:dead_letter" # 死信队列的 Redis Key
self.max_queue_size = 1000 # 最大队列深度
self.request_timeout = 300 # 请求超时 5 分钟
async def enqueue(self, request_data: dict) -> str:
"""入队:将请求放入队列"""
# 检查队列深度,防止积压过多
queue_len = await self.redis.llen(self.queue_key) # 获取当前队列长度
if queue_len >= self.max_queue_size: # 队列已满
raise QueueFullError(f"队列已满: {queue_len}/{self.max_queue_size}")
request_id = f"req_{int(time.time()*1000)}" # 生成请求 ID
request_data["request_id"] = request_id # 写入请求 ID
request_data["created_at"] = time.time() # 记录创建时间
# LPUSH 从左侧推入(新请求在队头)
await self.redis.lpush(
self.queue_key,
json.dumps(request_data, ensure_ascii=False) # 序列化为 JSON
)
return request_id # 返回请求 ID 供查询
async def dequeue(self, timeout: int = 0) -> Optional[dict]:
"""出队:从队列中取出一个请求(阻塞式)"""
# BRPOP 从右侧弹出(先进先出),timeout=0 表示永久阻塞
result = await self.redis.brpop(self.queue_key, timeout=timeout)
if result is None: # 超时未获取到
return None
_, data = result # result = (key, value)
return json.loads(data) # 反序列化
async def move_to_dlq(self, request_data: dict, error: str):
"""将失败请求移入死信队列"""
request_data["error"] = error # 记录错误信息
request_data["failed_at"] = time.time() # 记录失败时间
await self.redis.lpush(
self.dlq_key,
json.dumps(request_data, ensure_ascii=False) # 推入死信队列
)连接池配置
连接池管理对下游服务(LLM API、数据库、Redis)的连接复用,避免频繁创建和销毁连接的开销。
生活类比:连接池就像出租车车队。与其每次客人来了都临时叫一辆车(创建连接),不如保持一批车随时待命(连接池),用完归还。
# ============ PostgreSQL 连接池 ============
import asyncpg # 导入异步 PostgreSQL 驱动
async def create_db_pool():
"""创建 PostgreSQL 连接池"""
pool = await asyncpg.create_pool(
dsn="postgresql://localhost/agent_db", # 数据库连接字符串
min_size=5, # 最小连接数:保持 5 个常驻连接
max_size=50, # 最大连接数:上限 50
max_queries=50000, # 单个连接最大查询数(防止内存泄漏)
max_inactive_connection_lifetime=300.0, # 空闲连接最大生命周期 5 分钟
command_timeout=30 # 单条命令超时 30 秒
)
return pool
# ============ Redis 连接池 ============
import redis.asyncio as aioredis # 导入异步 Redis 客户端
async def create_redis_pool():
"""创建 Redis 连接池"""
pool = aioredis.ConnectionPool(
host="localhost", # Redis 主机地址
port=6379, # Redis 端口
db=0, # 数据库编号
max_connections=100, # 最大连接数
socket_keepalive=True, # 启用 TCP 保活
socket_connect_timeout=5, # 连接超时 5 秒
retry_on_timeout=True # 超时自动重试
)
return aioredis.Redis(connection_pool=pool)连接池配置的关键参数:
| 参数 | 建议值 | 说明 |
|---|---|---|
| max_connections | CPU 核数 × 50 | 取决于下游服务承受能力 |
| min_size | 5~10 | 避免冷启动延迟 |
| keepalive_expiry | 30~60s | 过短浪费连接,过长占用资源 |
| connect_timeout | 5~10s | 快速失败优于长时间等待 |
| read_timeout | 30~120s | LLM 调用需要较长的读取超时 |
经验法则:连接池的
max_connections不是越大越好。过大会导致下游服务(尤其是 LLM API)的连接数限制被打满,反而拖慢所有请求。建议从 CPU 核数 × 10 开始,逐步调优。
9.2.3 Rate Limit 处理策略
Rate Limit 是 LLM API 最常遇到的限制。不同的 LLM 提供商有不同的速率限制策略,核心思路都是在用户侧做限流,避免触发上游的限制。
策略一:令牌桶算法(Token Bucket)
生活类比:令牌桶就像一个漏水的水桶。水龙头以固定速率往桶里滴水(生成令牌),桶有最大容量(水满了就溢出)。每次处理请求就从桶里舀一瓢水(消耗令牌),桶空了就得等。
import asyncio # 导入异步框架
import time # 导入时间模块
class TokenBucket:
"""令牌桶限流器"""
def __init__(self, rate: float, capacity: int):
"""
rate: 每秒生成的令牌数
capacity: 桶的最大容量
"""
self.rate = rate # 令牌生成速率
self.capacity = capacity # 桶容量
self.tokens = capacity # 初始令牌数 = 满桶
self.last_refill = time.monotonic() # 上次补充时间
self._lock = asyncio.Lock() # 异步锁,保证线程安全
async def acquire(self, tokens: int = 1) -> bool:
"""尝试获取令牌,返回是否成功"""
async with self._lock: # 加锁防止并发竞争
self._refill() # 先补充令牌
if self.tokens >= tokens: # 令牌充足
self.tokens -= tokens # 消耗令牌
return True # 获取成功
return False # 令牌不足
async def wait_and_acquire(self, tokens: int = 1, timeout: float = None):
"""等待直到获取令牌或超时"""
deadline = time.monotonic() + timeout if timeout else None # 计算截止时间
while True: # 循环等待
async with self._lock: # 加锁
self._refill() # 补充令牌
if self.tokens >= tokens: # 令牌充足
self.tokens -= tokens # 消耗令牌
return True # 获取成功
if deadline and time.monotonic() >= deadline: # 检查超时
return False # 超时返回失败
await asyncio.sleep(0.1) # 短暂休眠后重试
def _refill(self):
"""补充令牌:根据经过的时间按速率补充"""
now = time.monotonic() # 当前时间
elapsed = now - self.last_refill # 经过的秒数
self.tokens = min(self.capacity, self.tokens + elapsed * self.rate) # 补充但不超过容量
self.last_refill = now # 更新补充时间
# 使用示例
bucket = TokenBucket(rate=10, capacity=20) # 每秒生成 10 个令牌,最多积压 20 个
async def call_with_rate_limit():
if await bucket.acquire(): # 尝试获取令牌
return await call_llm() # 获取成功,调用 LLM
else:
raise Exception("Rate limit exceeded") # 令牌不足,抛出异常策略二:滑动窗口算法
生活类比:滑动窗口就像停车场的计数器。它记录过去 1 分钟内进了多少辆车,一旦满额就拦住后面的。时间窗口不断向前滑动,之前的记录自动过期。
import time # 导入时间模块
from collections import deque # 导入双端队列
import asyncio # 导入异步框架
class SlidingWindowRateLimiter:
"""滑动窗口限流器"""
def __init__(self, max_requests: int, window_seconds: float):
self.max_requests = max_requests # 窗口内最大请求数
self.window = window_seconds # 时间窗口大小(秒)
self.requests = deque() # 用双端队列记录请求时间戳
self._lock = asyncio.Lock() # 异步锁
async def is_allowed(self) -> bool:
"""检查是否允许新请求"""
async with self._lock: # 加锁
now = time.monotonic() # 当前时间
# 移除窗口外的过期时间戳
while self.requests and self.requests[0] < now - self.window:
self.requests.popleft() # 从队头移除过期记录
if len(self.requests) < self.max_requests: # 窗口内未满
self.requests.append(now) # 记录当前请求时间
return True # 允许通过
return False # 窗口已满,拒绝
# 使用示例
limiter = SlidingWindowRateLimiter(
max_requests=60, # 每分钟最多 60 个请求
window_seconds=60 # 窗口 60 秒
)策略三:智能退避(Exponential Backoff)
当遇到 Rate Limit 时,不能简单重试(会再次被拒),需要使用指数退避策略——每次重试间隔翻倍,加上随机抖动。
生活类比:就像打客服电话占线,第一次等 1 分钟再打,第二次等 2 分钟,第三次等 4 分钟。加上随机抖动是为了避免所有被拒绝的人在同一时刻一起重试。
import asyncio # 导入异步框架
import random # 导入随机数模块
class RateLimitHandler:
"""Rate Limit 处理与智能退避"""
def __init__(self, max_retries: int = 5, base_delay: float = 1.0):
self.max_retries = max_retries # 最大重试次数
self.base_delay = base_delay # 基础延迟(秒)
async def call_with_retry(self, func, *args, **kwargs):
"""带重试的函数调用"""
for attempt in range(self.max_retries): # 最多重试 max_retries 次
try:
return await func(*args, **kwargs) # 尝试调用
except RateLimitError as e: # 遇到限流错误
if attempt == self.max_retries - 1: # 最后一次也失败
raise # 抛出异常
# 指数退避:延迟 = 基础延迟 × 2^attempt
delay = self.base_delay * (2 ** attempt) # 1, 2, 4, 8, 16...
# 随机抖动:避免重试风暴
jitter = random.uniform(0, delay * 0.3) # 抖动量 0~30%
total_delay = delay + jitter # 总延迟
# 如果 API 返回了 Retry-After 头,优先使用
if hasattr(e, 'retry_after') and e.retry_after:
total_delay = e.retry_after # 用服务端建议的等待时间
print(f"Rate limited, retrying in {total_delay:.1f}s "
f"(attempt {attempt + 1}/{self.max_retries})")
await asyncio.sleep(total_delay) # 等待后重试9.2.4 Celery / RabbitMQ / Kafka 异步处理
Celery + Redis 方案
Celery 是 Python 生态中最成熟的分布式任务队列,适合中小规模的 Agent 系统。
生活类比:Celery 就像一个外卖调度中心。订单来了,调度中心把任务分配给空闲的骑手(Worker),骑手送完了再领下一单。
# celery_config.py
from celery import Celery # 导入 Celery 框架
celery_app = Celery( # 创建 Celery 实例
'agent_worker', # 应用名称
broker='redis://localhost:***@celery_app.task( # 注册为 Celery 任务
bind=True, # 绑定 self,可访问任务实例
max_retries=3, # 最大重试次数
default_retry_delay=5, # 默认重试间隔 5 秒
rate_limit='10/m', # 限流:每分钟最多 10 个任务
queue='agent_tasks', # 指定队列名称
priority=5 # 任务优先级
)
def process_agent_request(self, user_message: str, user_id: str):
"""处理 Agent 请求的 Celery 任务"""
try:
client = openai.OpenAI() # 创建 OpenAI 客户端
response = client.chat.completions.create( # 调用 LLM
model="gpt-4", # 使用 GPT-4
messages=[{"role": "user", "content": user_message}], # 用户消息
stream=False # 非流式
)
return { # 返回结果
"reply": response.choices[0].message.content, # 提取回复
"usage": { # Token 用量
"prompt_tokens": response.usage.prompt_tokens,
"completion_tokens": response.usage.completion_tokens,
}
}
except openai.RateLimitError as e: # 遇到限流
raise self.retry(exc=e, countdown=10) # 10 秒后自动重试
except Exception as e: # 其他异常
raise self.retry(exc=e, countdown=5) # 5 秒后自动重试RabbitMQ 方案
RabbitMQ 适合需要精确路由和复杂消息模式的场景。
# rabbitmq_worker.py
import aio_pika # 导入异步 RabbitMQ 客户端
import asyncio # 导入异步框架
import json # 导入 JSON 处理
class RabbitMQAgentWorker:
"""基于 RabbitMQ 的 Agent Worker"""
def __init__(self, amqp_url: str = "amqp://guest:***@localhost/"):
self.amqp_url = amqp_url # AMQP 连接地址
self.connection = None # 连接对象
self.channel = None # 通道对象
async def connect(self):
"""建立连接并配置队列"""
self.connection = await aio_pika.connect_robust(self.amqp_url) # 建立鲁棒连接
self.channel = await self.connection.channel() # 创建通道
# 设置 QoS(每次只取一个任务,处理完再取下一个)
await self.channel.set_qos(prefetch_count=1) # 公平调度
# 声明持久化队列
self.task_queue = await self.channel.declare_queue(
"agent_tasks", # 队列名称
durable=True # 持久化:重启不丢消息
)
# 声明死信交换机(处理失败的消息)
self.dlx = await self.channel.declare_exchange(
"agent_dlx", # 死信交换机名称
aio_pika.ExchangeType.DIRECT # 直连模式
)
self.dlq = await self.channel.declare_queue(
"agent_dead_letter", # 死信队列名称
durable=True # 持久化
)
await self.dlq.bind(self.dlx, "dead") # 绑定路由键
async def publish_task(self, task_data: dict, priority: int = 5):
"""发布任务到队列"""
message = aio_pika.Message(
body=json.dumps(task_data).encode(), # 消息体(JSON 编码)
delivery_mode=aio_pika.DeliveryMode.PERSISTENT, # 持久化投递
priority=priority, # 优先级
expiration="300000" # 5 分钟过期
)
await self.channel.default_exchange.publish(
message, # 消息对象
routing_key="agent_tasks" # 路由到任务队列
)
async def consume_tasks(self, callback):
"""消费任务"""
async with self.task_queue.iterator() as queue_iter: # 创建迭代器
async for message in queue_iter: # 逐条消费
async with message.process(): # 自动 ACK
task_data = json.loads(message.body.decode()) # 解析消息
try:
await callback(task_data) # 执行回调
except Exception as e: # 处理失败
# 发送到死信队列
await self.channel.default_exchange.publish(
aio_pika.Message(
body=json.dumps({
"original": task_data, # 原始数据
"error": str(e) # 错误信息
}).encode()
),
routing_key="agent_dead_letter" # 路由到死信
)Kafka 方案
Kafka 适合超大规模(百万级 TPS)的流式处理场景。
# kafka_worker.py
from aiokafka import AIOKafkaProducer, AIOKafkaConsumer # 导入异步 Kafka 客户端
import json # 导入 JSON 处理
import asyncio # 导入异步框架
class KafkaAgentPipeline:
"""基于 Kafka 的 Agent 处理管道"""
def __init__(self, bootstrap_servers: str = "localhost:9092"):
self.bootstrap_servers = bootstrap_servers # Kafka 集群地址
async def create_producer(self):
"""创建生产者"""
producer = AIOKafkaProducer(
bootstrap_servers=self.bootstrap_servers, # Kafka 地址
value_serializer=lambda v: json.dumps(v).encode('utf-8'), # 值序列化器
compression_type='gzip', # 压缩传输,节省带宽
max_request_size=1048576, # 单条消息最大 1MB
linger_ms=10, # 批量发送延迟 10ms
)
await producer.start() # 启动生产者
return producer
async def create_consumer(self, group_id: str = "agent_workers"):
"""创建消费者"""
consumer = AIOKafkaConsumer(
'agent_requests', # 订阅请求 Topic
'agent_results', # 订阅结果 Topic
bootstrap_servers=self.bootstrap_servers, # Kafka 地址
group_id=group_id, # 消费者组
value_deserializer=lambda v: json.loads(v.decode('utf-8')), # 值反序列化
auto_offset_reset='earliest', # 从最早开始消费
enable_auto_commit=False, # 手动提交,确保至少处理一次
max_poll_records=10, # 每次最多拉取 10 条
)
await consumer.start() # 启动消费者
return consumer
async def process_stream(self):
"""处理流式请求"""
consumer = await self.create_consumer() # 创建消费者
producer = await self.create_producer() # 创建生产者
try:
async for msg in consumer: # 逐条消费消息
request_data = msg.value # 获取请求数据
try:
# 处理 Agent 请求
result = await self._process_agent(request_data)
# 发送成功结果到结果 Topic
await producer.send(
'agent_results', # 结果 Topic
value={
"request_id": request_data.get("request_id"),
"result": result,
"status": "success"
}
)
except Exception as e:
# 发送失败结果
await producer.send(
'agent_results',
value={
"request_id": request_data.get("request_id"),
"error": str(e),
"status": "failed"
}
)
# 手动提交 offset(确保消息已处理)
await consumer.commit()
finally:
await consumer.stop() # 关闭消费者
await producer.stop() # 关闭生产者9.2.5 技术选型对比
| 维度 | Celery + Redis | RabbitMQ | Kafka |
|---|---|---|---|
| 吞吐量 | 中等(~10K msg/s) | 中等(~50K msg/s) | 高(~1M msg/s) |
| 消息持久化 | 可选 | 是 | 是 |
| 消息顺序 | 不保证 | 保证(单队列) | 保证(单分区) |
| 消息优先级 | 支持 | 支持 | 不支持 |
| 运维复杂度 | 低 | 中 | 高 |
| 适用场景 | 中小规模 Agent | 需要精确路由 | 大规模流式处理 |
选型建议:如果你的 Agent 系统日请求量在 10 万以下,Celery + Redis 足够;如果需要复杂的消息路由和优先级队列,选 RabbitMQ;如果要做大规模实时流处理,选 Kafka。不要为了"高大上"而选 Kafka——它的运维成本远高于前两者。
常见误区
"连接池越大越好"——过大的连接池会耗尽下游服务的连接配额。比如 OpenAI API 对单个 API Key 有并发限制,连接池再大也没用,反而增加排队时间。建议从 CPU 核数 × 10 开始,根据 P99 延迟逐步调优。
"限流只在 API Gateway 做"——Gateway 层的限流是粗粒度的(按用户/IP),但 LLM API 层面的限流需要更细粒度。比如一个用户请求触发 3 次工具调用,每次工具调用又触发 LLM 调用,你需要在 LLM 调用层也做令牌桶限流。
"死信队列可有可无"——没有死信队列,失败的任务就直接丢失了。死信队列不仅防止消息丢失,还是排查问题的重要入口——定期检查死信队列可以发现隐藏的 Bug。
"指数退避不需要抖动"——如果不加随机抖动,所有被限流的请求会在同一时刻重试,形成"重试风暴",反而加重拥堵。抖动是指数退避的灵魂。
"用 Celery 就不用考虑并发安全了"——Celery 的 Worker 是多进程模型,但如果你用了共享资源(如全局字典、文件),仍然需要加锁或使用分布式锁。
"异步一定比同步快"——异步架构的引入增加了网络跳数和序列化开销。对于响应时间 < 2 秒的简单请求,同步架构可能反而更快。
本节小结
| 要点 | 说明 |
|---|---|
| 瓶颈识别 | Agent 系统的瓶颈在 LLM API(长尾延迟)和工具调用(级联放大) |
| 请求队列 | 削峰填谷第一道防线,Redis List/BLPOP 实现简单可靠 |
| 连接池 | max_connections 很重要,过大耗尽下游资源,过小导致排队 |
| Rate Limit | 令牌桶适合平滑突发,滑动窗口适合精确控制,指数退避处理限流响应 |
| 异步框架 | Celery 适合中小规模,RabbitMQ 适合需要精确路由,Kafka 适合大规模流式 |
| 死信队列 | 处理失败的任务,防止丢失,支持后续人工介入或自动重试 |
本节讨论了高并发场景下的核心武器——请求队列削峰填谷、连接池复用连接、三种限流算法各有用武之地、三种异步框架按规模选型。但无论限流做得多么精密,如果每个请求都打到 LLM 上,成本和延迟仍然压不下来。下一节我们将进入另一个关键话题:缓存策略——如何让 30% 的请求在毫秒级返回,成本几乎为零。
下一节:9.3 缓存策略——从 Prefix Cache 到语义缓存,四层缓存让系统又快又省。