前面的文章中,我们已经讨论过 HITL、Guardrails 等 Agent 的可控性问题。解决了“能不能安全、可控地执行”之后,接下来很自然会遇到另一个工程问题:当任务量上来以后,Agent 怎么避免被串行执行拖慢?
LangChain 通过 Runnable 提供了 batch、abatch、*_as_completed 和 max_concurrency 等能力,用来处理批量输入和并发控制;而在一次 Agent 执行内部,还可能同时出现多个 Tool Call。本文就围绕这些能力,看看 LangChain 如何让 Agent 从“逐个执行”走向更高吞吐。
一、Runnable 介绍
Runnable 可以理解为 LangChain 中统一的可执行对象接口。
一个 Runnable 接收输入、执行某项工作,再产生输出。模型、部分 Tool、由多个组件组合出的流程,以及编译后的 LangGraph,都可以按照这套接口运行。
Runnable 提供的几个常用方法是:
invoke / ainvoke
单次同步 / 异步调用
batch / abatch
多份输入的批量调用
batch_as_completed / abatch_as_completed
批量执行,并按实际完成顺序返回结果
stream / astream
流式执行
Runnable 对并发很重要,因为 batch 并不是 Agent 独有的特殊功能,而是 Runnable 执行模型的一部分。Runnable 默认的 batch 会在线程池中并行执行多次 invoke;abatch 默认通过 asyncio.gather 并发执行多次 ainvoke。
而 create_agent() 返回的是 CompiledStateGraph,编译后的 Graph 同样实现 Runnable 接口,所以我们可以这样调用:
agent.invoke(input)
agent.batch(inputs)
await agent.ainvoke(input)
await agent.abatch(inputs)
二、Agent为什么需要批量调用
假设现在有一组软件工单:
工单 T1024:登录接口频繁出现 500 错误
工单 T1025:用户反馈导出 PDF 后格式错乱
工单 T1026:搜索功能响应时间明显变长
工单 T1027:深色模式下代码块文字看不清
我们希望 Agent 根据工单描述,自行决定是否查询工单详情,并给出处理建议。
先定义一个简单的 Tool:
from langchain.agents import create_agent
from langchain.tools import tool
@tool
def get_ticket(ticket_id: str) -> str:
"""查询工单状态和相关信息。"""
tickets = {
"T1024": "高优先级,已分配给后端团队",
"T1025": "中优先级,等待复现",
"T1026": "高优先级,正在进行性能分析",
"T1027": "低优先级,已进入 UI 优化计划",
}
return tickets.get(ticket_id, "未找到工单")
agent = create_agent(
model="openai:gpt-5.5",
tools=[get_ticket],
system_prompt=(
"你负责分析软件工单。"
"必要时先查询工单详情,然后给出简短的处理建议。"
),
)
单次处理时,我们一般这样写:
result = agent.invoke({
"messages": [
{
"role": "user",
"content": "工单 T1024:登录接口频繁出现 500 错误。",
}
]
})
如果有几十个请求,最直接的写法可能是:
results = []
for request in requests:
result = agent.invoke(request)
results.append(result)
不难看出,这段代码存在效率问题。
因为每个请求都对应一次完整的 Agent 执行:
用户请求
↓
模型
↓
Tool
↓
模型
↓
结果
不同工单之间并不存在依赖,却被 for 循环强制串行执行。
三、使用 batch 批量执行 Agent
对于这种“多份独立输入交给同一个 Agent”的场景,可以直接使用:
agent.batch()
准备四个 Agent 输入:
requests = [
{
"messages": [{
"role": "user",
"content": "工单 T1024:登录接口频繁出现 500 错误。",
}]
},
{
"messages": [{
"role": "user",
"content": "工单 T1025:用户反馈导出 PDF 后格式错乱。",
}]
},
{
"messages": [{
"role": "user",
"content": "工单 T1026:搜索功能最近响应时间明显变长。",
}]
},
{
"messages": [{
"role": "user",
"content": "工单 T1027:深色模式下代码块文字看不清。",
}]
},
]
批量执行:
results = agent.batch(requests)
Runnable 默认的 batch 会并行执行这些 Agent 调用,而不是简单地在内部替我们写一个串行 for 循环。
运行后,模型生成内容可能略有不同,结合大致如下:
T1024:工单为高优先级,已分配后端团队,建议优先检查近期登录接口变更和服务日志。
T1025:当前等待复现,建议补充导出样例和异常文档,以便定位格式问题。
T1026:工单正在进行性能分析,建议重点检查搜索接口耗时和下游查询链路。
T1027:问题已进入 UI 优化计划,建议进一步确认代码块前景色和背景色的对比度。
虽然各个 Agent 任务实际完成时间可能不同,但 batch 最终返回的结果顺序仍然对应输入顺序。
因此它比较适合:
批量分析工单
批量处理客服请求
批量分析文档
批量执行研究任务
批量运行 Agent 评测
这里的重点是:
batch 批量执行的是完整的 Agent,而不只是批量调用一次模型。
每个输入都会独立跑一遍 Agent Loop,其中可能包含一次模型调用,也可能包含多次模型调用和 Tool 执行。
四、异步批量调用
如果应用本身运行在异步环境中,可以使用:
await agent.abatch(requests)
例如:
async def process_tickets():
results = await agent.abatch(requests)
for result in results:
print(result["messages"][-1].content)
Runnable 的 abatch 默认使用 asyncio.gather 并发执行多个 ainvoke。
在 FastAPI、异步 Worker 等环境中,通常更适合使用下面这些接口:
ainvoke
abatch
abatch_as_completed
需要注意,并发不会让一次 Agent 执行本身变快。
如果一次 Agent Loop 本来需要 4 秒,abatch 并不会把它缩短到 2 秒。它减少的是不同 Agent 请求之间没有必要的等待。
例如:
串行
Agent A ────→ Agent B ────→ Agent C ────→
改成并发后:
Agent A ────→
Agent B ────→
Agent C ────→
优化的是整批任务的完成时间和吞吐量。
五、按完成顺序处理
batch 适合“等全部完成,再统一处理”。
但后台任务经常有另一种需求:
哪个 Agent 先完成,就先保存哪个结果。
例如一次运行 100 个资料审核 Agent,其中有些只需要一次模型调用,有些需要连续查询多个 Tool。
它们的执行时间可能差异很大。
这时候可以使用:
batch_as_completed()
例如:
for index, result in agent.batch_as_completed(requests):
final_message = result["messages"][-1]
print(
f"任务 {index} 完成:",
final_message.content,
)
save_result(index, result)
batch_as_completed 会并行执行多个 invoke,但不会等所有任务结束后再统一返回,而是谁先结束就先产出谁的结果,同时返回它在原输入列表中的索引。
异步环境对应:
async for index, result in agent.abatch_as_completed(requests):
await save_result(index, result)
两类接口的区别可以简单理解为:
batch / abatch
全部完成 → 统一拿结果
*_as_completed
完成一个 → 处理一个
需要做任务进度、实时落库或者大量离线 Agent 处理时,后一种往往更实用。
六、使用 max_concurrency
能够并发,不意味着应该无限并发。
假设一次需要处理 500 个 Agent 请求,而每个 Agent 内部又可能调用两三次模型和多个 Tool。
如果同时放出 500 个 Agent,则可能会出现下列这些情况:
模型 API 出现 429
Tool 服务被打满
数据库连接池耗尽
HTTP 请求大量超时
失败重试进一步放大流量
RunnableConfig 提供了:
max_concurrency
用于限制同时执行的任务数量。
例如:
results = agent.batch(
requests,
config={
"max_concurrency": 8,
},
)
含义是:
任务总数:100
最大并发:8
任意时刻
最多执行 8 个 Agent
实际项目中,我们可以先设置一个保守值,再观察模型限流、Tool 延迟、数据库连接数和失败重试情况逐步调整。
七、Agent内部的并行 Tool Call
上面内容着重的是多个 Agent 请求之间并发。Agent 内部还有另一层并发:
一次 Agent 执行中的多个 Tool Call
例如用户问:
同时查询 T1024、T1025、T1026 三个工单的当前状态。
支持并行 Tool Call 的模型可以在一次响应中生成:
get_ticket("T1024")
get_ticket("T1025")
get_ticket("T1026")
很多支持 Tool Calling 的模型默认允许一次生成多个 Tool Call;部分模型也可以通过 parallel_tool_calls=False 禁止这种行为。
使用 create_agent 时,Agent 会负责 Tool 执行循环。当前实现内部使用 ToolNode,并通过 Send 机制分发多个 Tool Call,从而支持并行工具执行。
因此这里实际上存在两层并发:
第一层
请求 A ─┐
请求 B ─┼→ agent.batch()
请求 C ─┘
第二层
某个 Agent
↓
Model
↓
Tool Call A ─→ Tool
Tool Call B ─→ Tool
Tool Call C ─→ Tool
这也是为什么 Agent 系统比普通的批量模型调用更需要关注并发上限。
一个 max_concurrency=20 的 Agent 批量任务,并不意味着下游永远只有 20 个 HTTP 请求。
每个 Agent 内部还可能产生多轮模型调用和多个并行 Tool Call。
八、什么时候不要并发
是否适合并发,最先要判断的是任务之间有没有依赖。
例如一次生产发布流程:
查询发布任务
↓
确认测试状态
↓
检查是否满足发布条件
↓
执行生产发布
这几个步骤存在明确的数据依赖,不能因为“并发更快”就同时执行。
同样,如果 Tool 会修改共享资源,例如:
修改配置
执行部署
调整资源
变更权限
删除文件
就需要进一步考虑幂等性、重复执行、并发写入和失败恢复。
还有一个容易忽略的问题是 Memory。
如果 create_agent 配置了 Checkpointer,同一个 thread_id 表示同一条会话线程。批量运行多个彼此独立的 Agent 请求时,应给它们使用不同的 thread_id,不要让独立任务同时写入同一个会话状态。
batch 本身支持为不同输入传入不同的 RunnableConfig,因此有状态 Agent 批量执行时,更合理的做法是:
configs = [
{"configurable": {"thread_id": "task-001"}},
{"configurable": {"thread_id": "task-002"}},
{"configurable": {"thread_id": "task-003"}},
]
results = agent.batch(
requests,
config=configs,
)
并发解决的是执行效率,不负责自动解决状态隔离问题。
九、如何选择
实际开发中,可以按照执行对象来判断。
如果是一个 Agent 处理多份彼此独立的输入,使用 batch / abatch。
如果希望谁完成谁先处理,使用:batch_as_completed、abatch_as_completed。
如果需要限制同时运行的 Agent 数量,使用 max_concurrency 参数控制。
如果是一次 Agent 推理产生多个互不依赖的工具调用,使用 Parallel Tool Calls。
而如果任务存在明确的前后依赖,则应按 Agent Loop 或 Workflow 的执行顺序处理,不需要为了并发而并发。
总结
并发与 Batch 的价值,不是让单次 Agent 推理更快,而是减少多个独立任务之间的串行等待,让 Agent 在批量任务中获得更高吞吐。
实际落地时,建议先区分两类并发:一类是多个独立 Agent 请求之间的并发,另一类是单次 Agent 执行内部的并行 Tool Call。前者通过 batch、abatch 和 max_concurrency 控制,后者重点关注 Tool 是否存在依赖、共享状态和副作用。
对于带 Memory、数据库写入或外部操作的 Agent,还应进一步考虑 thread_id 隔离、幂等、限流、超时和重试,避免为了提高吞吐而引入新的稳定性问题。
社区讨论
参与讨论
有问题或想法?欢迎继续讨论。