LLM 工作流高并发防线实战:当请求并发拉满,工程上先守住哪条线
LLM 工作流在测试环境中可用,不代表高并发下仍能保持相同的延迟和资源占用。它包含 Prompt 拼接、外部工具调用和流式输出,请求生命周期通常比普通 API 更长。
上游响应变慢时,等待中的任务会占用连接和并发配额。应在服务入口设置明确的并发上限、排队策略和降级路径,而不是把容量风险留给上游服务。
1. 现场抓包与堆栈剖析:到底是什么拖垮了 Task 协程池
告警触发后的第一反应,是确认瓶颈到底卡在哪一步。
我们在网关与 LLM 工作流调度引擎之间挂了链路追踪。抓包后发现,当并发请求从 50 上升到 300 时,系统瓶颈并非出在向量数据库查询,也不是本地业务代码逻辑,而是以下两个互相叠加的连锁反应:
- Tool Calling 陷入递归循环:当模型返回的 JSON 格式偶尔出现字段缺失时,工作流没有立即切断,而是尝试让模型“自我纠错”。模型在极高并发下响应变慢,导致这种重试重定向变成了长达数秒的死锁循环。
- 连接池与 Semaphore 粒度失控:全局虽然设置了限流,但限流粒度只作用在 HTTP 入口,没有对底层 LLM API 请求的 HTTP Client 连接池做细粒度隔离。上游请求被放进来后,全部阻塞在等待 LLM 输出的 await 句柄上,进而把上游网关的连接池彻底耗尽。
知道了根因,治理思路就清晰了:必须在工作流入口、工具调用边界和 LLM 通信层建立三层硬硬的闸门。
2. 三层并发防线的设计与流转
在重新设计的工作流引擎中,我们取消了全局一刀切的简单计数器限流,改用基于令牌桶(Token Bucket)与动态 Semaphore 的分层防御架构。
流量进来的第一关是硬并发闸门,限制当前同时处于 await 状态的模型调用总数;第二关是单会话 Tool Calling 深度与时间预算控制;第三关则是输出端的结构体强校验。
flowchart TD A[客户端高并发请求] --> B{第一层: 动态 Semaphore 流量闸门} B -- 超过最大并发上限 --> C[立即触发降级/返回 HTTP 429 忙] B -- 许可通过 --> D[进入 LLM 工作流引擎] D --> E{第二层: 单 Request 时间与轮次预算控制} E -- 超出预算/死循环 --> F[硬切断并返回 Fallback 兜底结果] E -- 预算充足 --> G[调用 LLM API / 工具链] G --> H{第三层: Pydantic 结构体校验} H -- 校验失败且达到重试上限 --> F H -- 校验成功 --> I[流式吐字输出给客户端]这套结构的核心逻辑是:宁可直接拒绝一部分过载请求,也绝不允许把长耗时任务挂在服务器上消耗连接数。
3. 基于 Python 3.11 的生产级并发治理防线实现
下面是基于 Python 3.11 异步asyncio、pydantic校验与 Semaphore 隔离的防线实现代码。代码中包含了具体的超时切断、异常捕获与降级熔断处理。
import asyncio import time from typing import Dict, Any, Optional from pydantic import BaseModel, Field, ValidationError # 定义输出强约束 Schema,防范非确定性输出拖垮下游 class StepResultSchema(BaseModel): action: str = Field(..., description="下一步动作指令") data: Dict[str, Any] = Field(default_factory=dict, description="业务载荷") confidence: float = Field(..., ge=0.0, le=1.0, description="置信度评分") class LLMExecutionBudgetExceeded(Exception): """自定义异常:工作流执行超出时间或轮次预算""" pass class WorkflowConcurrencyGuard: def __init__(self, max_concurrent_llm_calls: int = 50, request_timeout_seconds: float = 5.0): # 限制同时发往 LLM 的底层请求数,守住连接池 self._semaphore = asyncio.Semaphore(max_concurrent_llm_calls) self._timeout = request_timeout_seconds self._active_calls = 0 async def execute_step_with_guard( self, request_id: str, prompt: str, max_tool_rounds: int = 3 ) -> Dict[str, Any]: """ 带防线保护的工作流节点执行入口 """ start_time = time.perf_counter() # 1. 第一层防线:并发限流切断 try: # 尝试在短时间内获取许可,拿不到直接切断,防范连接无限堆积 await asyncio.wait_for(self._semaphore.acquire(), timeout=0.5) except asyncio.TimeoutError: return { "status": "rejected", "reason": "系统当前并发过高,流量保护激活", "request_id": request_id } try: self._active_calls += 1 # 2. 第二层防线:时间预算总超时控制 return await asyncio.wait_for( self._run_tool_loop(request_id, prompt, max_tool_rounds), timeout=self._timeout ) except asyncio.TimeoutError: # 硬超时切断,防止挂起耗尽内存 return { "status": "degraded", "reason": f"执行超时(限额 {self._timeout}s),系统已自动熔断", "request_id": request_id, "fallback_payload": {"text": "系统繁忙,已采用默认文本兜底回复"} } except LLMExecutionBudgetExceeded as e: return { "status": "degraded", "reason": str(e), "request_id": request_id } except Exception as err: # 兜底未知异常,防止服务 Crash return { "status": "error", "reason": f"未预期工程错误: {type(err).__name__}", "request_id": request_id } finally: self._semaphore.release() self._active_calls -= 1 async def _run_tool_loop(self, request_id: str, prompt: str, max_rounds: int) -> Dict[str, Any]: current_round = 0 while current_round < max_rounds: current_round += 1 # 模拟发往 LLM 的异步通信 raw_response = await self._mock_llm_api_call(prompt, current_round) # 3. 第三层防线:结构体硬校验 try: validated_data = StepResultSchema.model_validate(raw_response) # 校验成功,返回确定的数据格式 return { "status": "success", "data": validated_data.model_dump(), "rounds_used": current_round } except ValidationError as ve: # 若校验失败且已达最大轮次,抛出预算超限,不继续无限纠错 if current_round >= max_rounds: raise LLMExecutionBudgetExceeded( f"模型输出格式连续 {max_rounds} 次不符合规范: {ve.errors()[0]['msg']}" ) # 未超轮次,记录日志后进入下一轮微调尝试 await asyncio.sleep(0.05) raise LLMExecutionBudgetExceeded("到达工具链调用最大轮次安全红线") async def _mock_llm_api_call(self, prompt: str, round_num: int) -> Dict[str, Any]: """模拟外部 LLM 网络响应与可能格式异常""" await asyncio.sleep(0.1) # 模拟网络 IO if round_num == 1: # 模拟首次返回异常字段引发 ValidationError return {"action": "continue", "confidence": "high_score"} # 第二次返回符合 Schema 的格式 return {"action": "finish", "data": {"answer": "处理完成"}, "confidence": 0.95}4. 压测验收清单
重构完成后,我们在压测环境使用 Locust 模拟了 500 到 2000 的梯度并发请求。
压测应比较加防线前后的 P99、超时比例、拒绝比例和资源曲线。并发数、超时和阈值需要随模型、连接池和供应商配额调整,不能把示例数值当成结论。一个可发布的验收清单包括:
- 延迟与拒绝行为:记录各并发档位的 P99、HTTP 状态码和排队时间,并区分正常完成、超时和主动拒绝。
- 连带影响:观察连接池等待、CPU、内存和下游错误;出现异常时保留可复现的输入与配置,不能仅凭一次平稳曲线宣称风险消失。
工程实战给出的教训非常直接:在涉及 LLM 这类高延迟、非确定性外部依赖的架构设计中,宁可快速失败,也绝不能让长请求无限占用系统连接资源。入口控制住了,底线也就守住了。