一、背景与问题
1.1 业务场景引入
在开发钉钉机器人集成服务时,我们就遇到了这样一个具体问题:
当用户询问"今天北京和上海的天气如何"时,传统的Agent执行流程是这样的:
用户提问 → AI思考 → 调用北京天气工具 → 获取结果 → AI再次思考 → 调用上海天气工具 → 获取结果 → AI总结回复
这意味着需要2次AI模型调用,且工具是串行执行的,总耗时较长。同时,在实际业务场景中,我们还需要在执行的不同阶段插入自定义逻辑,比如日志记录、性能监控、流式输出、权限控制、数据增强等。
问题来了: 如何优化工具执行性能,同时灵活支持业务逻辑的扩展呢?
1.2 核心概念解析
并发工具调用
并发工具调用是指AI模型能够在一次交互中同时选择多个工具调用,后台底层实现也支持并发执行这些工具。
改进后的流程:
用户提问 → AI思考 → 并发调用北京和上海天气工具 → 并发获取结果 → AI总结回复
只需要1次AI模型调用,大大缩短了总体执行时间。
Agent中间件
Agent中间件是一种设计模式,允许我们在Agent执行的不同阶段插入自定义逻辑,而不需要修改核心代码。参考LangChain的Middleware设计模式,我们实现了一套灵活的Agent中间件系统。
1.3 方案对比分析
针对性能优化和业务扩展问题,我们对比了三种方案:
方案一:传统串行执行 + 硬编码业务逻辑
优点: 实现简单 缺点:
- 工具串行执行,性能差
- 业务逻辑与核心流程耦合严重
- 代码难以维护和扩展
方案二:并发执行 + 中间件机制 ✅
优点:
- 工具并发执行,性能显著提升
- 中间件机制灵活,支持精确控制
- 业务逻辑与核心流程完全解耦
- 支持中间件组合使用
缺点: 实现复杂度稍高
适用场景: 需要高性能和灵活扩展性的企业级应用
二、系统架构设计
2.1 整体执行流程
┌─────────────────────────────────────────────────────────────┐
│ Agent执行流程 │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ Step 1: 模型调用前 │
│ - beforeModelCall() 中间件回调 │
│ - 可以修改请求或中断调用 │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ Step 2: AI模型调用 │
│ - 调用LLM模型(支持并发工具选择) │
│ - 如果失败,触发onModelCallError中间件 │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ Step 3: 模型调用后 │
│ - afterModelCall() 中间件回调 │
│ - 解析工具调用请求 │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────┐
│ 是否需要执行工具? │
└─────────────────┘
│
┌───────────────┴───────────────┐
│ │
▼ ▼
[是] [否]
│ │
▼ │
┌─────────────────────┐ │
│ 判断是否并发执行? │ │
└─────────────────────┘ │
│ │
┌───────┴───────┐ │
│ │ │
▼ ▼ │
[并发] [串行] │
│ │ │
▼ ▼ │
┌─────────────┐ ┌─────────────┐ │
│并发执行前 │ │工具执行前 │ │
│中间件触发 │ │中间件触发 │ │
└─────────────┘ └─────────────┘ │
│ │ │
▼ ▼ │
┌─────────────┐ ┌─────────────┐ │
│并发执行工具 │ │串行执行工具 │ │
│(线程池并发) │ │(逐个执行) │ │
└─────────────┘ └─────────────┘ │
│ │ │
▼ ▼ │
┌─────────────┐ ┌─────────────┐ │
│并发执行后 │ │工具执行后 │ │
│中间件触发 │ │中间件触发 │ │
└─────────────┘ └─────────────┘ │
│ │ │
└───────┬───────┘ │
│ │
▼ │
┌─────────────────────┐ │
│ Agent停止前 │ │
│ - onStop() 中间件 │ │
└─────────────────────┘ │
│ │
└───────┬───────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 流程结束 │
└─────────────────────────────────────────────────────────────┘
2.2 关键组件设计
**并发执行引擎:**基于CompletableFuture和线程池实现,支持多个工具的并发执行和结果收集。
**中间件管理器:**管理和触发各个阶段的中间件回调,支持优先级排序和异常容错。
线程池管理器:ToolExecutionThreadPoolManager单例管理器,按会话隔离线程池,避免资源竞争。
2.3 中间件体系架构
中间件触发时机:
- 工具执行前:
beforeToolExecution()- 单个工具执行前 - 工具执行后:
afterToolExecution()- 单个工具执行后 - 并发执行前:
beforeConcurrentToolExecution()- 一组工具并发执行前 - 并发执行后:
afterConcurrentToolExecution()- 一组工具并发执行后 - 工具执行错误:
onToolExecutionError()- 工具执行失败时 - 模型调用前:
beforeModelCall()- LLM模型调用前 - 模型调用后:
afterModelCall()- LLM模型调用后 - 模型调用错误:
onModelCallError()- LLM模型调用失败时 - Agent停止时:
onStop()- Agent正常或异常停止时
中间件执行顺序:
通过PriorityMiddleware接口定义优先级,中间件按优先级顺序执行。
中间件组合使用:
支持同时注册多个中间件,形成中间件链,每个中间件专注单一职责。
三、并发工具调用实现
3.1 提示词引导与模型支持
系统提示词自动注入
我们在Agent的系统提示词中添加了并发执行的引导语:
// @author changlu @since 2026-04-09
public static final String CONCURRENT_TOOL_EXECUTION_PROMPT = """
# 并发执行工具
你可以同时调用多个工具来提高效率。当多个工具之间没有依赖关系时,请并发调用它们。
例如:如果需要同时获取天气和新闻,可以在一次响应中同时调用两个工具。
""";
在ReActAgent的构造函数中,当工具数量大于1时,自动注入并发执行提示词:
// @author changlu @since 2026-04-09
// 构建基础提示词
String basePrompt = buildBasePrompt();
// 如果工具数量大于1,添加并发执行提示词
int toolCount = this.toolService.toolSpecifications().size();
if (toolCount > 1) {
basePrompt += AgentGlobalConstant.CONCURRENT_TOOL_EXECUTION_PROMPT;
log.debug("[ReActAgent] 工具数量: {}, 已注入并发执行提示词", toolCount);
}
说明: 只有当Agent拥有多个工具时,才会注入并发执行提示词,避免不必要的提示干扰。
3.2 开关控制机制
AgentSettings配置
在AgentSettings中添加了并发执行的开关:
// @author changlu @since 2026-04-09
public class AgentSettings {
// 是否启用并发执行工具,默认为true
private final boolean enableConcurrentToolExecution;
// Builder模式配置
public Builder enableConcurrentToolExecution(boolean enableConcurrentToolExecution) {
this.enableConcurrentToolExecution = enableConcurrentToolExecution;
return this;
}
}
执行逻辑判断
在ReActAgent.act()方法中判断是否启用并发执行:
// @author changlu @since 2026-04-09
@Override
protected StepResult act(int currentStep, List<ToolExecutionRequest> curActTools, ChatContext chatContext) {
log.debug("[ReActAgent.act] 开始执行工具调用,工具数量: {}", curActTools.size());
// 判断是否需要并发执行
boolean enableConcurrent = agentSettings.isEnableConcurrentToolExecution();
boolean shouldConcurrent = enableConcurrent && curActTools.size() > 1;
List<ChatMessage> toolMessages;
if (shouldConcurrent) {
toolMessages = executeToolsConcurrently(curActTools, chatContext);
} else {
toolMessages = executeToolsSequentially(curActTools, chatContext);
}
String resultSummary = toolMessages.toString();
log.debug("[ReActAgent.act] 所有工具执行完成,结果汇总: {}", resultSummary);
return StepResult.toFinished(resultSummary);
}
重点: 只有当enableConcurrentToolExecution为true且工具数量大于1时,才会使用并发执行。
3.3 并发执行流程详解
并发执行的完整流程包含4个步骤,与中间件深度集成:
// @author changlu @since 2026-04-09
private List<ChatMessage> executeToolsConcurrently(List<ToolExecutionRequest> toolRequests, ChatContext chatContext) {
log.debug("[ReActAgent.executeToolsConcurrently] 启用并发执行,工具数量: {}", toolRequests.size());
Map<String, ToolExecutor> executorMap = getExecutorMap();
// 第一步:触发并发工具执行前的中间件
middlewareManager.triggerBeforeConcurrentToolExecution(toolRequests, chatContext);
// 获取当前会话的线程池
Object memoryId = chatContext.getMemoryId();
ExecutorService executor = ToolExecutionThreadPoolManager.getInstance()
.getExecutorService(this.agentName, memoryId);
// 第二步:并发执行所有工具
List<CompletableFuture<ToolExecutionResult>> futures = toolRequests.stream()
.map(toolRequest -> CompletableFuture.supplyAsync(() -> {
return executeToolWithoutMiddleware(toolRequest, executorMap, chatContext);
}, executor))
.collect(Collectors.toList());
// 等待所有工具执行完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
// 第三步:收集所有工具执行结果
List<String> toolResults = new ArrayList<>();
List<ChatMessage> toolMessages = new ArrayList<>();
boolean hasError = false;
for (int i = 0; i < toolRequests.size(); i++) {
ToolExecutionRequest toolRequest = toolRequests.get(i);
ToolExecutionResult result = futures.get(i).join();
String processedResult;
if (result.isSuccess()) {
processedResult = result.getResult();
log.debug("[ReActAgent.executeToolsConcurrently] 工具 {} 执行结果: {}",
toolRequest.name(), processedResult);
} else {
// 工具执行失败,触发error middleware
if (result.getException() != null) {
middlewareManager.triggerOnToolExecutionError(
toolRequest, result.getException(), chatContext);
}
processedResult = result.getError();
hasError = true;
log.warn("[ReActAgent.executeToolsConcurrently] 工具 {} 执行失败: {}",
toolRequest.name(), processedResult);
}
// 记录工具执行结果消息
if (StringUtils.isEmpty(processedResult)) {
processedResult = "tool exec success, but no result";
}
toolResults.add(processedResult);
ToolExecutionResultMessage toolMessage = ToolExecutionResultMessage.from(toolRequest, processedResult);
addMessage(chatContext, toolMessage);
toolMessages.add(toolMessage);
}
// 第四步:触发并发工具执行后的中间件
middlewareManager.triggerAfterConcurrentToolExecution(toolRequests, toolResults, chatContext);
return toolMessages;
}
重点: 并发执行流程与中间件深度集成,在执行前、执行后、错误时都会触发相应的中间件回调。
四、中间件机制实现
4.1 中间件接口设计
AgentMiddleware接口定义
// @author changlu @since 2026-04-09
public interface AgentMiddleware {
// ===================== 工具执行相关回调 ===================== //
/**
* 在工具执行之前调用(单个工具)
*/
default void beforeToolExecution(ToolExecutionRequest toolRequest, ChatContext chatContext) {}
/**
* 在工具执行之后调用(单个工具)
*/
default String afterToolExecution(ToolExecutionRequest toolRequest, String toolResult, ChatContext chatContext) {
return toolResult;
}
/**
* 在并发工具执行之前调用(一组工具)
*/
default void beforeConcurrentToolExecution(List<ToolExecutionRequest> toolRequests, ChatContext chatContext) {}
/**
* 在并发工具执行之后调用(一组工具)
*/
default void afterConcurrentToolExecution(List<ToolExecutionRequest> toolRequests, List<String> toolResults, ChatContext chatContext) {}
/**
* 在工具执行发生错误时调用
*/
default void onToolExecutionError(ToolExecutionRequest toolRequest, Throwable error, ChatContext chatContext) {}
// ===================== 模型调用相关回调 ===================== //
/**
* 在模型调用之前调用
* @return 处理后的聊天请求(返回null将阻止调用)
*/
default ChatRequest beforeModelCall(int currentStep, ChatRequest chatRequest, ChatContext chatContext) {
return chatRequest;
}
/**
* 在模型调用之后调用
* @return 处理后的聊天响应(返回null将使用原始响应)
*/
default ChatResponse afterModelCall(int currentStep, ChatRequest chatRequest, ChatResponse chatResponse, ChatContext chatContext) {
return chatResponse;
}
/**
* 在模型调用发生错误时调用
*/
default void onModelCallError(int currentStep, ChatRequest chatRequest, Throwable error, ChatContext chatContext) {}
// ===================== Agent停止相关回调 ===================== //
/**
* 在Agent停止执行时调用(无论是正常结束还是达到最大步数)
*/
default void onStop(int currentStep, StopResult stopResult, ChatContext chatContext) {}
/**
* 在Agent因错误而停止时调用
*/
default void onStopWithError(int currentStep, Throwable error, ChatContext chatContext) {}
}
4.3 内置中间件实现
StreamingAgentMiddleware - 流式输出
功能:实时流式输出执行进度,展示AI思考过程和工具调用详情。
// @author changlu @since 2026-04-09
public class StreamingAgentMiddleware implements AgentMiddleware {
private static final ThreadLocal<SessionContext> SESSION_CONTEXT_HOLDER = new ThreadLocal<>();
@Override
public void beforeConcurrentToolExecution(List<ToolExecutionRequest> toolRequests, ChatContext chatContext) {
if (toolRequests == null || toolRequests.isEmpty()) return;
StringBuilder sb = new StringBuilder();
// 0. 添加简约分隔符
sb.append("▹ · · · · · · · · · · · · · · · · ↺\n");
// 1. 提示并发执行
int toolCount = toolRequests.size();
sb.append("⚡ 并发执行 ").append(toolCount).append(" 个工具\n");
// 2. 列出每个工具及其参数
for (int i = 0; i < toolRequests.size(); i++) {
ToolExecutionRequest toolRequest = toolRequests.get(i);
ThoughtExtractionResult extractionResult = extractAndRemoveThought(toolRequest.arguments());
sb.append(" ").append(i + 1).append(". ").append(mapToolName(toolRequest.name()));
// 显示参数
String paramDisplay = formatParameterForDisplay(extractionResult.cleanedArguments);
if (paramDisplay != null) {
sb.append(" [").append(paramDisplay).append("]");
}
sb.append("\n");
}
sb.append(" → 执行中 ...\n");
handleThinkingState(true, sb.toString());
}
@Override
public void afterConcurrentToolExecution(List<ToolExecutionRequest> toolRequests, List<String> toolResults, ChatContext chatContext) {
if (toolRequests == null || toolRequests.isEmpty()) return;
StringBuilder sb = new StringBuilder();
// 1. 总体完成提示
sb.append(" ✓ 并发完成\n");
// 2. 显示每个工具的结果预览
for (int i = 0; i < toolRequests.size(); i++) {
ToolExecutionRequest toolRequest = toolRequests.get(i);
String toolResult = toolResults != null && i < toolResults.size() ? toolResults.get(i) : "";
sb.append(" ↳ ").append(mapToolName(toolRequest.name())).append(": ");
sb.append(formatEnhancedResultPreview(toolResult)).append("\n");
}
sb.append("\n");
handleThinkingState(true, sb.toString());
}
private String formatEnhancedResultPreview(String result) {
if (result == null || result.isEmpty()) {
return "无结果";
}
// 截取前50个字符作为预览
return result.length() > 50 ? result.substring(0, 50) + "..." : result;
}
private void handleThinkingState(boolean isThinking, String message) {
SessionContext sessionContext = SESSION_CONTEXT_HOLDER.get();
if (sessionContext != null && sessionContext.streamResponseHandler != null) {
sessionContext.streamResponseHandler.handleThinking(isThinking, message);
}
}
}
五、场景测试验证
我配置的agent如下:
**
问题:
查询下北京、杭州天气以及这个禅道问题:http://zenpms.dtstack.cn/zentao/bug-view-148911.html分析下给我
场景验证:
1)没有开启并发tool执行情况:

三组工具串行调用时间为3s
2)开启并发tool执行情况

并发执行情况实际上只需要1.4s即可,能够进行一定的速度加快实现对话。
整理者:长路 时间:2026.4.10
评论区请在客户端页面查看