企业级 Agent 异步并发实战:从线上事故到高并发架构
1. 从一次线上事故说起企业级 Agent 的异步并发到底难在哪去年下半年我接手了一个企业级 Agent 平台的稳定性治理工作这个平台对外提供智能体编排、工具调用、多轮对话记忆、RAG 检索增强等能力日均请求量在百万级别。上线初期一切看起来都挺美好直到某天下午两点半监控突然报警整个 Agent 服务的 P99 延迟从 800ms 飙到 12s大量请求超时上游业务方电话直接打到了我们主管那里。排查过程相当痛苦。日志里没有明显的报错CPU 和内存看起来也不算离谱但就是所有请求都卡住了。最后定位到根因一个工具调用节点里用了同步的数据库驱动在异步框架里直接阻塞了事件循环。一个请求卡住整个 worker 里所有协程全部排队等待雪崩就这么发生了。这件事之后我花了大概三个月时间把整个 Agent 平台的异步并发模型从头到尾梳理了一遍踩了无数的坑也总结了不少经验。今天就把这些东西完整地分享出来从最基础的概念讲起一直讲到企业级场景下的实战方案。不管你是刚接触 Agent 开发的新人还是已经写过一些异步代码但总觉得心里没底的老手相信都能从里面找到对自己有用的东西。这篇文章会覆盖异步编程的底层原理、Agent 场景下特有的并发模型、数据库和外部工具调用的异步化改造、并发控制与限流、常见故障排查等内容。我会尽量用大白话把原理讲清楚同时给出可以直接抄作业的代码和配置。涉及 Python 生态的内容会多一些因为目前大部分 Agent 框架都是 Python 写的但核心思路对 Rust async、Java 并发同样适用。2. 异步编程的底层逻辑为什么 Agent 场景绕不开它2.1 同步、异步、并发、并行这四个词到底啥区别很多人刚入门的时候会被这几个词绕晕我用一个生活化的例子来解释。假设你开了一家奶茶店只有一个店员。同步模式就是顾客点单店员开始做做完递给顾客然后才接待下一位。中间做奶茶的几分钟店员就站在那儿等着什么都不干。这就是同步阻塞效率极低。异步模式是顾客点单店员把单子记下来告诉顾客好了叫你然后立刻接待下一位。做奶茶这件事被挂起了等奶茶机那边有结果了再回来处理。店员的时间被充分利用起来了。并发是指店员同时处理多个订单的能力比如一边等奶茶机出杯一边给新顾客点单。并行则是真的有两个店员同时在干活。在单核 CPU 上我们只能做到并发做不到真正的并行多核才能并行。放到 Agent 场景里一个用户请求进来Agent 需要调用大模型、查询知识库、调用外部工具、写数据库这些操作大部分时间都在等 IOCPU 其实是闲着的。如果用同步方式写一个 worker 同时只能服务一个请求吞吐量低得可怜。异步就是为了解决这个问题让一个 worker 在等待 IO 的时候去处理别的请求。2.2 事件循环异步的心脏Python 的 asyncio 核心是一个事件循环Event Loop。你可以把它想象成一个永不停歇的调度员手里维护着一个任务队列。每个协程coroutine在执行到await的时候会把控制权交还给事件循环事件循环就去执行队列里的下一个任务。等之前的 IO 操作完成了操作系统会通知事件循环事件循环再把对应的协程唤醒继续执行。这里有个关键点协程的切换是协作式的不是抢占式的。也就是说只有协程主动await的时候才会让出控制权。如果你在协程里写了一段纯 CPU 密集的代码中间没有任何await那事件循环就被你独占了其他所有协程都得干等着。这就是为什么在异步框架里调用同步阻塞函数是致命的。我见过太多人写这样的代码async def handle_request(request): # 这是同步阻塞调用会卡死整个事件循环 result requests.get(https://api.example.com/data) return process(result)requests.get是同步的它执行的时候整个线程都被占住了事件循环根本转不动。正确的做法是用httpx或aiohttp这类异步 HTTP 客户端import httpx async def handle_request(request): async with httpx.AsyncClient() as client: result await client.get(https://api.example.com/data) return process(result)2.3 async/await 的底层原理别只停留在会用async def定义的函数调用后返回的是一个协程对象不会立即执行。只有被await或者被包装成 Task 交给事件循环才会真正运行。await的本质是把当前协程挂起注册一个回调等被 await 的对象完成后事件循环再恢复这个协程。Python 的协程底层是基于生成器generator实现的yield和send是它的核心机制。理解这一点很重要因为它解释了为什么协程不能跨线程随意调度也解释了为什么某些操作在协程里做会出问题。Rust 的 async 则是另一套机制它编译成状态机没有运行时开销性能更极致但学习曲线也更陡。如果你做的是超高性能的 Agent 网关Rust async 是很好的选择如果追求开发效率Python asyncio 完全够用。2.4 Agent 场景为什么天然适合异步Agent 的工作流本质上是一连串的 IO 等待等大模型返回、等工具执行、等数据库查询、等向量检索。这些操作的耗时占比通常在 90% 以上CPU 真正计算的时间很少。这种 IO 密集型场景正是异步的主场。举个具体数字。假设一个 Agent 请求平均耗时 2 秒其中 1.8 秒在等 IO。同步模式下一个 worker 每秒只能处理 0.5 个请求。异步模式下如果并发度开到 100理论上每秒能处理 50 个请求提升 100 倍。当然实际受限于下游服务的承载能力但量级上的差距是实实在在的。3. 企业级 Agent 的并发模型设计从单机到集群3.1 一个请求进来Agent 内部发生了什么先梳理一下典型 Agent 请求的完整链路。用户发来一条消息网关接收后做鉴权、限流然后路由到 Agent 编排引擎。编排引擎根据配置加载对应的 Agent 定义开始执行工作流可能是先调用大模型做意图识别然后根据识别结果决定调用哪些工具工具可能包括数据库查询、外部 API 调用、代码执行等最后把结果汇总再交给大模型生成回复。这条链路上每一步都是潜在的并发瓶颈。大模型调用有速率限制数据库连接池有限外部 API 可能不稳定工具执行可能超时。如果这些环节没有做好并发控制很容易出现一个慢请求拖垮整个服务的情况。3.2 并发模型选型协程、线程池还是多进程在 Python 里做 Agent并发模型的选择主要有三种模型适用场景优势劣势asyncio 协程IO 密集型如 LLM 调用、HTTP 请求单线程高并发开销小阻塞调用会卡死CPU 密集不适用线程池需要调用同步库如某些 SDK兼容性好改造量小GIL 限制线程切换开销大多进程CPU 密集型如本地模型推理真正并行隔离性好内存开销大进程间通信复杂我的实践经验是主链路用 asyncio同步库用run_in_executor包一层线程池CPU 密集任务丢给独立进程或独立服务。这样既能享受异步的高并发又能兼容生态里大量的同步库。具体来说如果某个工具必须用同步的 SDK可以这样处理import asyncio from concurrent.futures import ThreadPoolExecutor executor ThreadPoolExecutor(max_workers20) async def call_sync_tool(params): loop asyncio.get_running_loop() result await loop.run_in_executor(executor, sync_tool_function, params) return result注意线程池的max_workers不能设太大否则线程切换开销会吃掉收益一般设为 CPU 核数的 2-4 倍比较合适。16C32G 的机器设 32-64 就差不多了。3.3 16C32G 服务器到底能扛多少并发这是被问得最多的问题但答案真的取决于你的业务特征。我拿我们平台的实测数据来说。一个典型的 Agent 请求平均耗时 2.5 秒其中 LLM 调用 1.5 秒工具调用 0.6 秒其他 0.4 秒。在 16C32G 的机器上用 uvicorn 起 8 个 worker 进程每个进程一个事件循环每个进程并发度开到 200理论上单机能扛 1600 并发。但实测下来稳定运行在 800-1000 并发左右P99 延迟能控制在 3 秒以内。为什么达不到理论值因为下游有瓶颈。LLM API 有 QPS 限制数据库连接池只有 100 个连接外部工具服务也有自己的承载上限。所以并发能力不是单机决定的而是整条链路最短的那块板决定的。我的建议是先做压测找到瓶颈点再针对性优化。不要盲目堆机器也不要盲目调大并发参数。3.4 多进程 多协程的部署架构生产环境我推荐这样的部署方式用 gunicorn 或 uvicorn 起多个 worker 进程每个进程内部跑一个 asyncio 事件循环。进程数一般设为 CPU 核数或核数加一16 核就起 16-17 个进程。每个进程的并发度根据压测结果调整。gunicorn app:app \ --worker-class uvicorn.workers.UvicornWorker \ --workers 16 \ --bind 0.0.0.0:8000 \ --timeout 120 \ --graceful-timeout 30 \ --max-requests 10000 \ --max-requests-jitter 1000max-requests是为了防止内存泄漏跑够一定请求数就重启 worker。Agent 服务里经常会有各种缓存和临时对象长时间运行内存容易涨定期重启是个简单有效的办法。4. 核心环节的异步化改造实战4.1 数据库层从同步驱动到异步驱动Agent 平台离不开数据库会话历史、用户配置、工具调用记录都要落库。Python 生态里同步的数据库操作是重灾区一个不小心就把事件循环卡死了。以 PostgreSQL 为例同步方案是 psycopg2 SQLAlchemy 同步 Session异步方案是 asyncpg 或 psycopg3 SQLAlchemy 异步 Session。我强烈建议新项目直接上异步方案。from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import sessionmaker engine create_async_engine( postgresqlasyncpg://user:passhost/db, pool_size20, max_overflow10, pool_pre_pingTrue, pool_recycle3600, ) AsyncSessionLocal sessionmaker( engine, class_AsyncSession, expire_on_commitFalse ) async def get_session(): async with AsyncSessionLocal() as session: yield session连接池参数很关键。pool_size是常驻连接数max_overflow是峰值时能额外创建的连接数。总数不要超过数据库的max_connections除以 worker 进程数。比如数据库最大连接 500起了 16 个 worker那每个 worker 的池子最多 30 左右留点余量给其他服务。pool_pre_ping一定要开它会在取连接前先 ping 一下避免拿到已经断开的死连接。pool_recycle设得比数据库或中间件的空闲超时短一些防止连接被服务端主动断开。4.2 LLM 调用并发控制与重试的艺术LLM 调用是 Agent 里最贵也最慢的环节。企业级场景下你不可能无限制地并发调用供应商那边有速率限制你自己的成本也要控制。我的做法是用信号量Semaphore做并发控制配合指数退避重试import asyncio from tenacity import retry, stop_after_attempt, wait_exponential llm_semaphore asyncio.Semaphore(50) retry( stopstop_after_attempt(3), waitwait_exponential(multiplier1, min1, max10), ) async def call_llm(prompt, modelgpt-4): async with llm_semaphore: async with httpx.AsyncClient(timeout60) as client: resp await client.post( LLM_ENDPOINT, json{model: model, messages: prompt}, ) resp.raise_for_status() return resp.json()信号量设为 50 意味着同时最多 50 个 LLM 调用在飞。这个数字要根据供应商的速率限制和你自己的配额来定。如果供应商限制 100 QPS而你的平均调用耗时 2 秒那并发度设 200 才能打满 QPS但为了安全起见设 100-150 比较稳妥。重试策略也有讲究。LLM 调用失败分几种情况网络超时、限流 429、服务端 5xx、内容审核拒绝。前三种可以重试最后一种重试也没用。所以重试逻辑里要判断错误类型别一股脑全重试。4.3 工具调用超时、熔断、降级三板斧Agent 会调用各种各样的工具外部 API、内部服务、代码执行环境等等。这些工具的质量参差不齐有的可能几秒不响应有的可能直接挂掉。如果不做保护一个坏工具能把整个 Agent 拖垮。我的经验是给每个工具调用都加上超时、熔断、降级import asyncio from circuitbreaker import circuit circuit(failure_threshold5, recovery_timeout30) async def call_external_tool(tool_name, params, timeout10): try: async with asyncio.timeout(timeout): result await do_call(tool_name, params) return result except asyncio.TimeoutError: logger.warning(ftool {tool_name} timeout) return {error: timeout, fallback: True}asyncio.timeout是 Python 3.11 引入的比wait_for更清晰。熔断器在连续失败 5 次后打开30 秒内直接拒绝请求给下游服务恢复的时间。降级就是返回一个兜底结果让 Agent 能继续往下走而不是整个请求失败。4.4 向量检索批量与并行的平衡RAG 场景下经常要同时查多个向量库或者做多路召回。串行查太慢全并行又可能打爆下游。我的做法是分批并行每批控制在合理数量async def batch_retrieve(queries, batch_size10): results [] for i in range(0, len(queries), batch_size): batch queries[i:ibatch_size] batch_results await asyncio.gather( *[retrieve(q) for q in batch], return_exceptionsTrue, ) results.extend(batch_results) return resultsreturn_exceptionsTrue很重要它保证一个查询失败不会影响其他查询。否则gather会在第一个异常时直接抛出其他任务的结果就丢了。5. 并发控制与限流的进阶技巧5.1 全局限流 vs 单租户限流企业级 Agent 平台通常服务多个租户限流要分两层。全局限流保护整个平台不被压垮单租户限流防止某个租户的突发流量影响其他人。全局限流可以用令牌桶算法单租户限流可以用滑动窗口。Python 里slowapi或者自己基于 Redis 实现都可以。核心是要把限流状态放在 Redis 这种共享存储里因为你有多个 worker 进程。import redis import time class SlidingWindowLimiter: def __init__(self, redis_client, limit, window): self.redis redis_client self.limit limit self.window window async def acquire(self, key): now time.time() pipe self.redis.pipeline() pipe.zremrangebyscore(key, 0, now - self.window) pipe.zadd(key, {str(now): now}) pipe.zcard(key) pipe.expire(key, self.window) _, _, count, _ await pipe.execute() return count self.limit5.2 背压机制让系统优雅地拒绝背压Backpressure是个容易被忽视的概念。当系统负载超过处理能力时与其让请求堆积到内存爆炸不如主动拒绝一部分。这听起来反直觉但对企业级系统来说是必须的。实现方式很简单监控当前正在处理的任务数超过阈值就直接返回 503让上游重试或降级。class BackpressureMiddleware: def __init__(self, max_concurrent1000): self.semaphore asyncio.Semaphore(max_concurrent) async def __call__(self, request, call_next): if self.semaphore.locked(): return JSONResponse( {error: service busy}, status_code503 ) async with self.semaphore: return await call_next(request)5.3 任务队列把长任务异步化有些 Agent 任务耗时很长比如批量文档处理、复杂工作流编排同步等待不现实。这时候要把任务丢到队列里异步执行客户端轮询或通过回调获取结果。Celery 是经典选择但它对 asyncio 支持一般。我更喜欢用arq或者基于 Redis Stream 自己实现。核心是把任务序列化后入队worker 从队列里取出来执行结果写回存储。import arq async def process_agent_task(ctx, task_id, payload): result await run_agent(payload) await ctx[redis].set(ftask:{task_id}, result, ex3600) return result5.4 优雅关闭别让重启丢请求服务重启时正在处理的请求不能直接掐断。要监听关闭信号停止接收新请求等存量请求处理完再退出。import signal async def shutdown(sig, loop): logger.info(freceived {sig.name}, shutting down) tasks [t for t in asyncio.all_tasks() if t is not asyncio.current_task()] for task in tasks: task.cancel() await asyncio.gather(*tasks, return_exceptionsTrue) loop.stop() loop asyncio.get_event_loop() for sig in (signal.SIGTERM, signal.SIGINT): loop.add_signal_handler(sig, lambda ssig: asyncio.create_task(shutdown(s, loop)))配合 gunicorn 的graceful-timeout给存量请求留出足够时间。Agent 请求可能跑几十秒这个超时要设得比最长请求时间还长。6. 常见故障与排查技巧实录6.1 事件循环被阻塞的典型症状与定位症状很典型所有请求突然变慢但 CPU 不高日志没有报错重启后短暂恢复又复发。这时候八成是某个地方有同步阻塞调用。定位方法用asyncio的调试模式开启慢回调检测。import asyncio loop asyncio.get_event_loop() loop.set_debug(True) loop.slow_callback_duration 0.1 # 超过100ms就告警或者用py-spy直接 dump 堆栈看哪个函数占着线程不放py-spy dump --pid worker_pid我踩过的一个坑是日志库。某些同步日志 handler 在写文件时会阻塞高并发下日志量一大就卡。解决方案是用QueueHandler把日志写入队列后台线程异步落盘。6.2 数据库连接池耗尽的排查报错通常是TimeoutError: QueuePool limit of size X overflow Y reached。原因可能是连接泄漏用完没归还、慢查询占着连接不放、或者池子设太小。排查步骤先看数据库端的pg_stat_activity看有多少活跃连接、都在跑什么 SQL。如果发现大量idle in transaction说明有代码开了事务没提交。如果发现慢查询加索引或优化 SQL。SELECT pid, state, query_start, query FROM pg_stat_activity WHERE state ! idle ORDER BY query_start;6.3 并发数上不去的原因分析有时候你明明把并发度调大了吞吐量却上不去。常见原因有几个下游有速率限制、连接池不够、GIL 争抢、或者有隐藏的锁竞争。用压测工具如 locust、wrk逐步加压观察 QPS 和延迟曲线。如果 QPS 到某个点就不涨了延迟却飙升说明到了瓶颈。这时候要逐层排查网关、应用、数据库、下游服务。6.4 常见问题速查表现象可能原因排查方法解决方案所有请求卡住事件循环被阻塞py-spy dump找出同步调用改异步连接池超时连接泄漏或池太小查 pg_stat_activity修复泄漏或调大池子内存持续上涨任务堆积或缓存泄漏看任务队列长度加背压或定期重启P99 延迟高慢请求拖累分布式追踪定位慢环节优化限流频繁触发并发度设置不当看限流日志调整信号量或加机器6.5 我踩过的几个印象深刻的坑第一个坑是asyncio.gather的异常处理。默认情况下一个任务抛异常gather 会立即抛出其他任务的结果全丢。必须加return_exceptionsTrue然后自己判断每个结果。第二个坑是asyncio.create_task创建的任务如果没人 await异常会被吞掉。要用task.add_done_callback处理异常或者维护一个任务集合定期清理。第三个坑是上下文变量contextvars在任务间传递。Agent 请求的 trace_id、用户信息这些要跨协程传递用 contextvars 是对的但要注意create_task会复制当前上下文如果任务创建时机不对上下文可能是错的。第四个坑是异步生成器的清理。Agent 流式输出常用异步生成器如果客户端提前断开生成器要正确关闭否则资源泄漏。要用aclose()或者async with管理。7. 一些关于 Agent 并发架构的个人体会做企业级 Agent 这一年多我最大的感受是异步并发不是炫技而是被业务逼出来的必然选择。当你的 QPS 从几十涨到几千当你的下游从一两个变成十几个同步模型根本撑不住。但异步也不是银弹。它带来了复杂度调试更难错误更隐蔽对开发者的要求更高。我的建议是如果你的业务量不大同步模型加多进程完全够用别为了异步而异步。只有当 IO 等待成为瓶颈同步模型的资源利用率低到无法接受时再考虑异步化改造。改造的时候要循序渐进别想着一次性全改。先把最耗时的环节异步化比如 LLM 调用和数据库查询观察效果再决定下一步。改造过程中一定要有完善的监控和压测否则你根本不知道改动是变好了还是变坏了。最后分享一个我常用的压测脚本模板基于 locust可以模拟 Agent 请求的并发场景from locust import HttpUser, task, between class AgentUser(HttpUser): wait_time between(0.1, 0.5) task def chat(self): self.client.post(/v1/agent/chat, json{ agent_id: test_agent, message: 帮我查一下今天的订单, stream: False, }, timeout60)压测的时候要逐步加压从 10 并发开始每次翻倍观察 QPS、延迟、错误率的变化。找到拐点那就是你系统的真实容量。留 30% 的余量应对突发流量剩下的就是你的安全水位。这套东西说起来简单做起来全是细节。每个参数背后都有取舍每个方案都有适用边界。希望我这些踩坑经验能帮你少走点弯路。