从零搭建 AI Agent 调度中台(二):Agent 对话主链路——从用户提问到 LLM 回答的 9 步闭环

第 3 / 14 章
从零搭建 AI Agent 调度中台(二):Agent 对话主链路——从用户提问到 LLM 回答的 9 步闭环

一、ChatServiceImpl:平台的心脏

代码结构 ChatServiceImpl 是整个 Agent-Core 模块的核心,674 行代码,完成一次对话需要串联 4 个外部服务:Agent-Core 自己、RAG-Knowledge、Model-Gateway、Tool-Plugin。它不调任何人,但所有人都得调它。

先把主方法 chat() 的 9 步流程摆出来:

// ChatServiceImpl.java — 主方法骨架(简化后)
public ChatResponse chat(ChatRequest request) {
    Long tenantId = resolveTenantId(request);
    String sessionId = request.getSessionId();

    // 1. 加载 Agent 配置(走 Redis 缓存)
    AgentInfo agent = agentService.loadAgent(request.getAgentId());

    // 2. Redis 历史会话 + Token 滑动窗口裁剪
    List<ChatMessageVO> history = sessionContextService.slidingWindow(
        tenantId, agent.getId(), sessionId, windowTokens, reserveTokens);

    // 3. RAG 多路召回(失败降级为空列表,不阻断主对话)
    List<RagRecallResponse.RecallFragment> ragFragments = doRagRecall(agent, query, tenantId);

    // 4. 构建 LLM messages:system + 历史 + 当前提问
    List<ModelChatRequest.Message> messages = buildMessages(agent, history, query, ragFragments);

    // 5. 构建 ModelChatRequest(模型名/temperature/tools schema)
    // ...

    // 6. 工具白名单 + 可委派 Agent → tools schema
    modelReq.setTools(buildToolsSchema(toolWhitelist, delegateAgentIds));

    // 7. LLM 调用 + Function Calling 工具循环(核心)
    LlmLoopResult loopResult = callLlmWithToolLoop(modelReq, ...);

    // 8. 持久化会话(user + assistant)
    sessionContextService.appendMessage(tenantId, agentId, sessionId, userMsg);
    sessionContextService.appendMessage(tenantId, agentId, sessionId, assistantMsg);

    // 9. 组装响应
    ChatResponse resp = new ChatResponse();
    resp.setAnswer(loopResult.answer);
    resp.setRagFragments(ragFragments);
    resp.setTokenUsage(loopResult.usage);
    return resp;
}

这是一个编排器——它不做任何计算,只负责调别人、拼结果、处理异常。下面逐段拆解。

二、第 1-2 步:Agent 配置 + Redis 会话

// 1. 加载 Agent 配置
AgentInfo agent = agentService.loadAgent(request.getAgentId());

AgentInfo 里存什么?不是一个简单的"人设",而是一整套对话运行时参数:

public class AgentInfo {
    private Long id;
    private String agentName;
    private String systemPrompt;      // System Prompt
    private String modelName;          // 默认模型(可被请求覆盖)
    private Double temperature;        // 温度
    private String knowledgeBaseIds;   // 逗号分隔的知识库 ID
    private String toolIds;            // 逗号分隔的工具白名单
    private String delegateAgentIds;   // 可委派的子 Agent
    private Integer contextWindowTokens; // 上下文窗口大小(默认 4096)
    private Integer maxTokens;         // 单次最大输出
    private Boolean enabled;
}

这个设计的关键点:AgentInfo 里的 knowledgeBaseIds/toolIds/delegateAgentIds 是逗号分隔的字符串,不是外键关联。为什么?因为 Agent 配置可能被频繁增删改,用关联表每次 CRUD Agent 都要多表操作;而逗号字符串在 AgentService 层一次性加载,ChatService 里用 parseLongList() 转成 Set 就行——简单、够用、读多写少。

Redis 会话:双结构存储

// 2. 读取历史 + 滑动窗口裁剪
List<ChatMessageVO> history = sessionContextService.slidingWindow(
    tenantId, agent.getId(), sessionId, windowTokens, reserveTokens);

会话怎么存?Redis 三个 Key:

Hash: aap:agent:session:{tenantId}:{agentId}:{sessionId}
       field=seq → value=消息JSON

ZSet: aap:agent:tokens:{tenantId}:{agentId}:{sessionId}
       member=seq → score=timestamp

String: aap:agent:session:seq:{tenantId}:{agentId}:{sessionId}
       value=自增序号

读写流程:

// 写入一条新消息(appendMessage)
Long seq = stringRedisTemplate.opsForValue().increment(seqKey);  // seq++
stringRedisTemplate.opsForHash().put(hashKey, String.valueOf(seq), json);   // Hash 存内容
stringRedisTemplate.opsForZSet().add(zsetKey, String.valueOf(seq), timestamp);  // ZSet 存时间序

// 读取历史(loadHistory)
Set<String> seqs = stringRedisTemplate.opsForZSet().range(zsetKey, 0, -1);  // ZSet 按时间取有序 seq
for (String seq : seqs) {
    String json = stringRedisTemplate.opsForHash().get(hashKey, seq);       // Hash 查内容
    messages.add(JSON.parseObject(json, ChatMessageVO.class));
}

为什么是 Hash + ZSet 两个结构? Hash 擅长 field→value 的随机读写,追加消息用 opsForHash().put() O(1);ZSet 擅长按分数范围查询,取历史消息用 range(zsetKey, 0, -1) O(N) 但保证时间有序。如果只用一个 Hash,你没法按时间排——Hash 的 field 是无序的。

TTL 统一刷新:每次 append 都给三个 Key 续 7 天 TTL(agent.sessionTtlSeconds = 604800),保证活跃会话不丢失、僵尸会话自动清理。

三、第 2.5 步:Token 滑动窗口裁剪

多轮对话的致命问题:上下文溢出。GPT-4o 的窗口是 128K,但企业场景下用的 deepseek-chat 是 32K——用户聊了 20 轮之后,历史消息可能占 25K,再加上 System Prompt 和 RAG 片段,就炸了。

解决方案:滑动窗口裁剪——从最早的消息开始丢,保证最近上下文能放进去。

TokenSlidingWindowManager.java

// 自后向前累加历史消息 Token,累计超过 budget 即停止
public List<ChatMessageVO> trim(List<ChatMessageVO> history, int maxTokens, int reserveTokens) {
    int budget = Math.max(0, maxTokens - reserveTokens);  // 预算 = 总窗口 - 预留额度

    int accumulated = 0;
    int keepFromIndex = history.size(); // 默认全裁(从后往前定位起点)
    for (int i = history.size() - 1; i >= 0; i--) {
        int tok = tokenOf(history.get(i));
        accumulated += tok;
        if (accumulated > budget) {
            keepFromIndex = i + 1;  // 从下一条开始保留
            break;
        }
        keepFromIndex = i;
    }
    // 至少保留最近一条
    if (keepFromIndex >= history.size()) {
        keepFromIndex = history.size() - 1;
    }
    if (keepFromIndex <= 0) {
        return new ArrayList<>(history);  // 未超限,全部保留
    }
    return new ArrayList<>(history.subList(keepFromIndex, history.size()));
}

算法逻辑(和 AI 生成的空壳完全不同——AI 写的是"从前往后累加直到超限",但那会把最新的消息裁掉!):

历史消息(时间升序,早 → 晚):[m1, m2, m3, m4, m5, m6]

自后向前累加:
  m6: 400 tokens → accumulated = 400
  m5: 300 tokens → accumulated = 700
  m4: 500 tokens → accumulated = 1200
  ... 继续加直到 accumulated > budget

假设 budget = 2048,加到 m2 时 accumulated = 2100 → 超限
→ keepFromIndex = m3 的索引
→ 保留 [m3, m4, m5, m6]

budget = maxTokens - reserveTokens。maxTokens 默认 4096,reserveTokens 默认 512——512 是给 System Prompt + 当前提问留的安全区,不让滑动窗口把这两块挤没了。

Token 估算器

精确算 Token 要用 tiktoken(OpenAI 官方),但引入成本大。我们用工程化近似:

// TokenEstimator.java
public static int estimate(String text) {
    int cjk = 0, ascii = 0;
    for (int i = 0; i < text.length(); i++) {
        char c = text.charAt(i);
        if (isCjk(c)) cjk++;
        else ascii++;
    }
    return cjk + (ascii + 3) / 4;  // 中文 1 字 ≈ 1 token,英文 4 字符 ≈ 1 token
}

private static boolean isCjk(char c) {
    return (c >= 0x4E00 && c <= 0x9FFF)   // CJK 统一表意文字
        || (c >= 0x3000 && c <= 0x30FF)   // 日文假名
        || (c >= 0xFF00 && c <= 0xFFEF);  // 全角符号
}

误差能接受吗?滑动窗口裁剪允许 ±20% 的估算误差——因为我们留了 512 的 reserve buffer,即使实际 Token 比估算多 20%,也只会多裁一两条早期消息,不会影响最新上下文。

四、第 3 步:RAG 多路召回(失败降级)

private List<RagRecallResponse.RecallFragment> doRagRecall(AgentInfo agent, String query, Long tenantId) {
    List<Long> kbIds = parseLongList(agent.getKnowledgeBaseIds());
    if (kbIds.isEmpty()) {
        return Collections.emptyList();  // Agent 没绑知识库,跳过
    }
    try {
        Result<RagRecallResponse> ragResp = ragKnowledgeFeignClient.recall(ragReq);
        if (ragResp.isSuccess()) {
            return ragResp.getData().getFragments();
        }
    } catch (Exception e) {
        log.warn("[ChatService] RAG 召回异常,降级无上下文,err={}", e.getMessage());
    }
    return Collections.emptyList();  // 失败降级为空列表
}

关键设计:RAG 失败不阻断主对话。如果 Milvus 挂了或者 ES 超时,Agent 只是失去知识库上下文,仍然用纯 LLM 能力回答用户——用户不会感知到底层有问题,只是回答可能没那么准。

这和 AI 生成的代码有本质区别:AI 写的 RAG 调用是直接抛出异常的(用 ragKnowledgeFeignClient.recall(...).getData() 然后 .getFragments(),中间任何一步 null 就 NPE)。我们改成 try-catch + 降级。

五、第 4 步:构建 LLM Messages

private List<ModelChatRequest.Message> buildMessages(AgentInfo agent,
                                                     List<ChatMessageVO> history,
                                                     String query,
                                                     List<RagRecallResponse.RecallFragment> ragFragments) {
    List<ModelChatRequest.Message> messages = new ArrayList<>();

    // 1. System Prompt + RAG 上下文
    StringBuilder system = new StringBuilder();
    system.append(agent.getSystemPrompt());
    if (ragFragments != null && !ragFragments.isEmpty()) {
        system.append("\n\n【知识库参考】\n");
        for (RagRecallResponse.RecallFragment f : ragFragments) {
            system.append("- 【").append(f.getFileName()).append("】")
                  .append(f.getChunkContent()).append("\n");
        }
    }
    messages.add(new ModelChatRequest.Message("system", system.toString()));

    // 2. 历史会话(滑动窗口裁剪后的)
    for (ChatMessageVO h : history) {
        messages.add(new ModelChatRequest.Message(h.getRole(), h.getContent()));
    }

    // 3. 当前提问
    messages.add(new ModelChatRequest.Message("user", query));
    return messages;
}

System Prompt 里拼接 RAG 片段——这不是最佳实践(好的 RAG 实现应该用独立的 role=system 或 role=context 消息块),但它对 LLM 来说足够清晰。生产级优化方向:把 RAG 片段作为独立的 system message 段落,加 clear separator 告诉 LLM"这是知识库内容,优先参考"。

六、第 5-6 步:构建 ModelChatRequest + Tools Schema

// 5. ModelChatRequest
ModelChatRequest modelReq = new ModelChatRequest();
modelReq.setModel(agent.getModelName());        // deepseek-chat
modelReq.setMessages(messages);                  // system + 历史 + 当前
modelReq.setTemperature(agent.getTemperature()); // 默认 0.7
modelReq.setStream(false);                       // Feign 走非流式
modelReq.setTenantId(tenantId);

// 6. Tools Schema
Set<Long> toolWhitelist = parseLongList(agent.getToolIds());
Set<Long> delegateAgentIds = parseLongList(agent.getDelegateAgentIds());
if (!toolWhitelist.isEmpty() || !delegateAgentIds.isEmpty()) {
    modelReq.setEnableToolCalling(true);
    modelReq.setTools(buildToolsSchema(toolWhitelist, delegateAgentIds));
}

BuildToolsSchema:两类工具

private List<Map<String, Object>> buildToolsSchema(Set<Long> toolIds, Set<Long> delegateAgentIds) {
    List<Map<String, Object>> tools = new ArrayList<>();

    // ===== 第一类:普通 HTTP/MCP 工具 =====
    // 从 Tool-Plugin 服务拉取元数据(描述 + JSON Schema 入参)
    List<ToolMetaDTO> metaList = fetchToolMeta(toolIds);
    for (ToolMetaDTO m : metaList) {
        Map<String, Object> function = new HashMap<>();
        function.put("name", String.valueOf(m.getId()));           // 工具名 = ID
        function.put("description", m.getDescription());          // LLM 看的描述
        function.put("parameters", parseParamsSchema(m.getParamsSchema()));  // 入参 JSON Schema
        Map<String, Object> tool = Map.of("type", "function", "function", function);
        tools.add(tool);
    }

    // ===== 第二类:Agent 委派工具 =====
    // 名字固定:delegate_agent_{目标AgentId}
    for (Long targetAgentId : delegateAgentIds) {
        AgentInfo target = agentService.loadAgent(targetAgentId);
        Map<String, Object> function = new HashMap<>();
        function.put("name", "delegate_agent_" + targetAgentId);
        function.put("description", "委派给【" + target.getAgentName() + "】处理");
        // 入参固定:{"query": "委派的问题"}
        Map<String, Object> params = Map.of(
            "type", "object",
            "properties", Map.of("query", Map.of("type", "string", "description", "委派给目标 Agent 的问题")),
            "required", List.of("query")
        );
        function.put("parameters", params);
        tools.add(Map.of("type", "function", "function", function));
    }
    return tools;
}

为什么工具名用 ID 而不是可读名称? 因为 Tool-Plugin 服务的 HTTP/SQL 执行器是按 ID 查的(treq.setToolId(toolId)),用 ID 省一次"名称→ID"的查询。Agent 委派工具名用 delegate_agent_ 前缀——LLM 看到这个前缀就知道这是"把问题扔给另一个 Agent 处理",而不是调一个 HTTP 接口。

七、第 7 步:LLM + Function Calling 工具循环(核心)

private LlmLoopResult callLlmWithToolLoop(ModelChatRequest modelReq, ...) {
    LlmLoopResult result = new LlmLoopResult();
    int maxRounds = agentProperties.getMaxToolCallRounds();  // 默认 5 轮
    int round = 0;

    while (round <= maxRounds) {
        // 调 LLM
        Result<ModelChatResponse> resp = modelGatewayFeignClient.chat(modelReq);
        ModelChatResponse mresp = resp.getData();
        accumulateUsage(result, mresp.getUsage());

        // 判断:有没有 tool_calls?
        boolean hasToolCalls = mresp.getToolCalls() != null && !mresp.getToolCalls().isEmpty();
        if (!hasToolCalls || round >= maxRounds) {
            result.answer = mresp.getContent();  // 没有工具调用,直接输出答案
            break;
        }

        // 执行工具调用
        for (ModelChatResponse.ToolCall tc : mresp.getToolCalls()) {
            ToolExecuteResponse toolResp = executeOneTool(tc, toolWhitelist, delegateAgentIds, ...);

            // 把 assistant 调用意图 + tool 结果回填 messages
            modelReq.getMessages().add(newMessage("assistant", "调用工具: " + tc.getFunction().getName()));
            modelReq.getMessages().add(newToolMessage(tc.getId(), toolResp.getResult()));
        }
        round++;
    }
    return result;
}

循环终止条件:要么 LLM 这次没返回 tool_calls(意味着它觉得自己能直接回答),要么 round 超过 maxToolCallRounds=5(防死循环)。

executeOneTool:三条分支

private ToolExecuteResponse executeOneTool(ToolCall toolCall, Set<Long> toolWhitelist, ...) {
    String funcName = toolCall.getFunction().getName();

    // ===== 分支 A:Agent 委派 =====
    if (funcName.startsWith("delegate_agent_")) {
        return executeAgentDelegate(toolCall, funcName, delegateAgentIds, ...);
    }

    // ===== 分支 B:普通工具 =====
    Long toolId = Long.parseLong(funcName);
    if (!toolWhitelist.contains(toolId)) {
        return buildFailedToolResp(toolId, "工具不在白名单");
    }
    ToolExecuteResponse tres = toolPluginFeignClient.execute(treq).getData();
    return tres;

    // ===== 分支 C:异常兜底 =====
    // 任何异常都返回结构化结果(toolId + success=false + errorMessage),不抛
}

executeAgentDelegate:本地 vs 远程

private ToolExecuteResponse executeAgentDelegate(...) {
    // 1. 提取目标 Agent ID + 白名单校验 + 深度检查
    Long targetAgentId = Long.parseLong(funcName.substring("delegate_agent_".length()));
    if (!delegateAgentIds.contains(targetAgentId)) { ... }  // 无权限
    if (delegateDepth >= maxDepth) { ... }  // 深度超限

    // 2. 判断:本地还是远程?
    AgentInfo targetAgent = null;
    boolean isLocal = true;
    try {
        targetAgent = agentService.loadAgent(targetAgentId);  // 本地数据库能查到
    } catch (Exception e) {
        isLocal = false;  // 查不到 → 远程 Agent(走 Python 编排层)
    }

    // 3. 本地委派:同 JVM 递归调 chat()
    if (isLocal) {
        ChatResponse delegateResp = chat(delegateReq);  // ⚠️ 递归调用!
        return ToolExecuteResponse(result=delegateResp.getAnswer());
    }

    // 4. 远程委派:走 Python 编排层 Feign
    AgentDelegateResponse remoteResp = pythonOrchestratorFeignClient
        .orchestrateAgentDelegate(delegateReq).getData();
    return ToolExecuteResponse(result=remoteResp.getResult());
}

本地委派是递归调用 chat()——这意味着 Agent A 委派给 Agent B,本质是 Agent A 的 LLM 决定"把这个 query 交给 Agent B 处理",然后 ChatService 直接调 chat(B的request),B 执行完返回答案,A 的 LLM 再决定要不要继续委派或者总结输出。

护栏:maxDelegateDepth = 3 防 A→B→A 递归死循环。每次递归 delegateDepth + 1,超过阈值返回 ToolExecuteResponse(success=false, result="委派深度超限,请直接回答"),上层 LLM 会收到这个结果并直接总结。

为什么 Java 本地委派有递归,Python 编排层也有递归?

本地委派是同 JVM 直接调 chat(),开销为 0,适合"Agent A 的子 Agent 也是本平台注册的 Agent"的场景。远程委派是跨 HTTP 调 Python,适合"Agent B 在 Python 编排层跑(可能用 LangGraph、可能是外部平台的 Agent)"的场景。

这就是 ch05 要讲的双栈架构设计——Java 做本地 Agent(零开销递归),Python 做远程 Agent(复杂编排),边界由"本地数据库能查到 agentId 吗"自动判定。

八、第 8-9 步:持久化 + 组装响应

// 8. 持久化
sessionContextService.appendMessage(tenantId, agent.getId(), sessionId,
    newMessage("user", request.getQuery()));
sessionContextService.appendMessage(tenantId, agent.getId(), sessionId,
    newMessage("assistant", loopResult.answer));

// 9. 响应
ChatResponse resp = new ChatResponse();
resp.setAnswer(loopResult.answer);
resp.setRagFragments(ragFragments);       // RAG 引用片段
resp.setToolCalls(loopResult.toolCalls);  // 工具调用记录
resp.setTokenUsage(loopResult.usage);     // Token 用量
resp.setCostMillis(System.currentTimeMillis() - start);
return resp;

持久化写入 Redis(双结构),不写 MySQL——会话数据是临时的(7 天 TTL 自动清理),高频读写用 Redis 比 MySQL 快一个数量级。如果需要长期会话存档,可以加一个异步消费 Redis 过期事件、写入 MySQL 归档表的消费者。

九、SSE 流式:打字机效果

聊天不止有非流式,还有 SSE 流式版本 chatStream():

public SseEmitter chatStream(ChatRequest request) {
    SseEmitter emitter = new SseEmitter(SSE_TIMEOUT_MS);  // 5 分钟超时
    Thread worker = new Thread(() -> {
        try {
            ChatResponse resp = chat(request);  // 还是走 chat(),只是输出方式变了
            String answer = resp.getAnswer();
            // 分片打字机推送(每次 4 字符,间隔 30ms)
            for (int i = 0; i < answer.length(); i += SSE_CHUNK_SIZE) {
                emitter.send(SseEmitter.event()
                    .name("message")
                    .data(answer.substring(i, Math.min(i + SSE_CHUNK_SIZE, answer.length()))));
                Thread.sleep(SSE_CHUNK_INTERVAL_MS);
            }
            // 推送完整结构(引用/工具/用量)
            emitter.send(SseEmitter.event().name("done").data(JSON.toJSONString(resp)));
            emitter.complete();
        } catch (Exception e) {
            emitter.completeWithError(e);
        }
    }, "aap-sse-" + request.getSessionId());
    worker.setDaemon(true);
    worker.start();
    return emitter;
}

注意:SSE 不边生成边推,而是先生成完再分片推——因为模型网关走的是非流式 Feign(modelReq.setStream(false))。真实生产中应该让模型网关也支持流式返回,但那需要改 Feign 为 WebClient + Flux,复杂度上升一个量级。当前实现是"打字机效果的体验优先",后端一次生成完,前端看起来像实时输出。

十、异常处理哲学

ChatServiceImpl 的异常处理有一套明确的策略:

外部依赖异常处理原因
RAG 召回try-catch → 降级为空列表RAG 是增强,缺失不阻断对话
Model-Gatewaytry-catch → throw BusinessExceptionLLM 是核心,调不通就是调不通
Tool-Plugintry-catch → 返回 ToolExecuteResponse(success=false)工具失败不能让整个对话挂掉,LLM 会收到失败结果并总结
Python 编排层try-catch → 返回 ToolExecuteResponse(success=false)同 Tool-Plugin
Agent 递归委派catch → 降级为远程委派本地查不到目标 Agent 时自动切远程

这和 AI 生成代码的风格完全不同——AI 生成的是"一路 Result.getData() 然后直接调 .xxx()",中间任何一步 null 就 NPE。我们是每一层都降级,保证对话链路尽可能走到终点。

十一、AI 生成 vs 人工硬焊总结

内容AI 生成质量我做了什么
chat() 方法骨架(9 步)★★★★☆基本框架正确,补了 delegateDepth 参数传递
Token 滑动窗口★☆☆☆☆AI 写的是"自后向前累加但方向反了"(最新消息被裁),重写为自后向前定位起点 + 至少保留最近一条
Token 估算★☆☆☆☆AI 用 tiktoken 重依赖,改成 CJK 按字、ASCII 按 4 字符的工程化近似
Redis 会话存储★★★☆☆Hash 写对了,但 ZSet 没写——补了 ZSet 做时间序 + TTL 刷新
Function Calling 工具循环★★☆☆☆AI 写了 while 框架,但没有 maxRounds 防死循环、没有 delegate_agent 递归、没有 本地 vs 远程路由
异常降级策略☆☆☆☆☆AI 完全没考虑降级——RAG/Tool/Python 任何一处 NPE 就炸。全部加了 try-catch + 降级返回
SSE 流式★★★☆☆AI 用了 WebClient 流式(但模型网关没流式),改成先生成完再分片推的简化实现

三个经验教训:

  1. Token 滑动窗口的方向不能反——AI 写的是"从前往后丢",会把最新消息裁掉,用户看到的是旧上下文回答。必须从后往前累加。
  2. 异常降级的边界要画清楚——核心依赖(LLM 调用)失败要 throw,增强依赖(RAG/工具/远程委派)失败要降级返回。这个边界 AI 给不出来,因为它不知道业务上什么是核心。
  3. Function Calling 要有护栏——maxToolCallRounds=5 + maxDelegateDepth=3 双保险,防死循环。AI 只会写 while 框架,不会主动加护栏参数。

下一章是核心卖点一——DAG 工作流引擎源码讲解。我们要逐行拆解 153 行 WorkflowEngine.java,对比 AI 生成的空壳和最终实现的差别。