Skip to content

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)       │
                 └──────────────┘

队列深度与超时配置

python
# 基于 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)的连接复用,避免频繁创建和销毁连接的开销。

生活类比:连接池就像出租车车队。与其每次客人来了都临时叫一辆车(创建连接),不如保持一批车随时待命(连接池),用完归还。

python
# ============ 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_connectionsCPU 核数 × 50取决于下游服务承受能力
min_size5~10避免冷启动延迟
keepalive_expiry30~60s过短浪费连接,过长占用资源
connect_timeout5~10s快速失败优于长时间等待
read_timeout30~120sLLM 调用需要较长的读取超时

经验法则:连接池的 max_connections 不是越大越好。过大会导致下游服务(尤其是 LLM API)的连接数限制被打满,反而拖慢所有请求。建议从 CPU 核数 × 10 开始,逐步调优。

9.2.3 Rate Limit 处理策略

Rate Limit 是 LLM API 最常遇到的限制。不同的 LLM 提供商有不同的速率限制策略,核心思路都是在用户侧做限流,避免触发上游的限制。

策略一:令牌桶算法(Token Bucket)

生活类比:令牌桶就像一个漏水的水桶。水龙头以固定速率往桶里滴水(生成令牌),桶有最大容量(水满了就溢出)。每次处理请求就从桶里舀一瓢水(消耗令牌),桶空了就得等。

python
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 分钟内进了多少辆车,一旦满额就拦住后面的。时间窗口不断向前滑动,之前的记录自动过期。

python
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 分钟。加上随机抖动是为了避免所有被拒绝的人在同一时刻一起重试。

python
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),骑手送完了再领下一单。

python
# 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 适合需要精确路由和复杂消息模式的场景。

python
# 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)的流式处理场景。

python
# 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 + RedisRabbitMQKafka
吞吐量中等(~10K msg/s)中等(~50K msg/s)高(~1M msg/s)
消息持久化可选
消息顺序不保证保证(单队列)保证(单分区)
消息优先级支持支持不支持
运维复杂度
适用场景中小规模 Agent需要精确路由大规模流式处理

选型建议:如果你的 Agent 系统日请求量在 10 万以下,Celery + Redis 足够;如果需要复杂的消息路由和优先级队列,选 RabbitMQ;如果要做大规模实时流处理,选 Kafka。不要为了"高大上"而选 Kafka——它的运维成本远高于前两者。

常见误区

  1. "连接池越大越好"——过大的连接池会耗尽下游服务的连接配额。比如 OpenAI API 对单个 API Key 有并发限制,连接池再大也没用,反而增加排队时间。建议从 CPU 核数 × 10 开始,根据 P99 延迟逐步调优。

  2. "限流只在 API Gateway 做"——Gateway 层的限流是粗粒度的(按用户/IP),但 LLM API 层面的限流需要更细粒度。比如一个用户请求触发 3 次工具调用,每次工具调用又触发 LLM 调用,你需要在 LLM 调用层也做令牌桶限流。

  3. "死信队列可有可无"——没有死信队列,失败的任务就直接丢失了。死信队列不仅防止消息丢失,还是排查问题的重要入口——定期检查死信队列可以发现隐藏的 Bug。

  4. "指数退避不需要抖动"——如果不加随机抖动,所有被限流的请求会在同一时刻重试,形成"重试风暴",反而加重拥堵。抖动是指数退避的灵魂。

  5. "用 Celery 就不用考虑并发安全了"——Celery 的 Worker 是多进程模型,但如果你用了共享资源(如全局字典、文件),仍然需要加锁或使用分布式锁。

  6. "异步一定比同步快"——异步架构的引入增加了网络跳数和序列化开销。对于响应时间 < 2 秒的简单请求,同步架构可能反而更快。

本节小结

要点说明
瓶颈识别Agent 系统的瓶颈在 LLM API(长尾延迟)和工具调用(级联放大)
请求队列削峰填谷第一道防线,Redis List/BLPOP 实现简单可靠
连接池max_connections 很重要,过大耗尽下游资源,过小导致排队
Rate Limit令牌桶适合平滑突发,滑动窗口适合精确控制,指数退避处理限流响应
异步框架Celery 适合中小规模,RabbitMQ 适合需要精确路由,Kafka 适合大规模流式
死信队列处理失败的任务,防止丢失,支持后续人工介入或自动重试

本节讨论了高并发场景下的核心武器——请求队列削峰填谷、连接池复用连接、三种限流算法各有用武之地、三种异步框架按规模选型。但无论限流做得多么精密,如果每个请求都打到 LLM 上,成本和延迟仍然压不下来。下一节我们将进入另一个关键话题:缓存策略——如何让 30% 的请求在毫秒级返回,成本几乎为零。


下一节:9.3 缓存策略——从 Prefix Cache 到语义缓存,四层缓存让系统又快又省。