通用Agent设计:实现并发tools工具调用与AgentMiddleware补充

一、背景与问题

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如下:

**image-20260410012640803

问题:

查询下北京、杭州天气以及这个禅道问题:http://zenpms.dtstack.cn/zentao/bug-view-148911.html分析下给我

场景验证:

1)没有开启并发tool执行情况:

image-20260410012731401

三组工具串行调用时间为3s

2)开启并发tool执行情况

image-20260410012931406

并发执行情况实际上只需要1.4s即可,能够进行一定的速度加快实现对话。


整理者:长路 时间:2026.4.10

评论区请在客户端页面查看