从零搭建 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-Gateway | try-catch → throw BusinessException | LLM 是核心,调不通就是调不通 |
| Tool-Plugin | try-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 流式(但模型网关没流式),改成先生成完再分片推的简化实现 |
三个经验教训:
- Token 滑动窗口的方向不能反——AI 写的是"从前往后丢",会把最新消息裁掉,用户看到的是旧上下文回答。必须从后往前累加。
- 异常降级的边界要画清楚——核心依赖(LLM 调用)失败要 throw,增强依赖(RAG/工具/远程委派)失败要降级返回。这个边界 AI 给不出来,因为它不知道业务上什么是核心。
- Function Calling 要有护栏——
maxToolCallRounds=5+maxDelegateDepth=3双保险,防死循环。AI 只会写 while 框架,不会主动加护栏参数。
下一章是核心卖点一——DAG 工作流引擎源码讲解。我们要逐行拆解 153 行 WorkflowEngine.java,对比 AI 生成的空壳和最终实现的差别。
