Agentic RAG 项目梳理
my-ragent 项目梳理
一、全链路
1.用户输入,前端 -> 后端
@GetMapping(value = "/rag/chat", produces = "text/event-stream;charset=UTF-8")
public SseEmitter chat(@RequestParam String question,
@RequestParam(required = false) String conversationId,
@RequestParam(required = false, defaultValue = "false") Boolean deepThinking) {
//Spring的SSE通道
SseEmitter emitter = new SseEmitter(ragDefaultProperties.getSseTimeoutMs());
//传入emitter
ragChatService.streamChat(question, conversationId, deepThinking, emitter);
//Spring接管,持续推送
return emitter;
}
参数:用户问题、对话ID(非必须)、是否开启深度思考(非必须)
rag 默认配置中的参数:
- collectionName: 默认向量集合名称,用于指定在向量数据库中存储向量数据的默认集合名称
- dimension: 指定向量的维数,要与所使用的 Embedding 模型的输出维数保持一致。例如 BGE 输出 1024 维向量,PG 建表时定义 embedding vector(1024) 。如果你存 256 维的向量到这个字段,PG 直接拒绝写入。向量维度不是元数据,是索引结构的物理约束
- metricType: 向量相似度度量类型,常见取值有:余弦相似度、欧氏距离、内积
- sseTimeouts: SSE 超时时间,防止 SSE 连接泄露,超时后自动关闭连接。连接泄露指客户端断开了但代码里没 complete(),线程和连接仍然占用,积压多了 Tomcat 线程池耗尽,新请求进不来。超时配置就是为了兜底自动断开这种死连接
通过 rag 默认配置参数创建了 SSE 发射器,并传入调用 Service 层的函数中,最终给前端返回发射器对象,前端将 SSE 流式输出的内容显示出来
2.Service层
@Override
public void streamChat(String question, String conversationId, Boolean deepThinking, SseEmitter emitter) {
//生成ConversationId
String actualConversationId = StrUtil.isBlank(conversationId) ? IdUtil.getSnowflakeNextIdStr() : conversationId;
//生成taskId
String taskId = IdUtil.getSnowflakeNextIdStr();
log.info("开始流式对话,会话ID:{},任务ID:{}", actualConversationId, taskId);
boolean thinkingEnabled = Boolean.TRUE.equals(deepThinking);
//创建SSE事件处理器
StreamCallBack callback = callbackFactory.createChatEventHandler(emitter, actualConversationId, taskId);
//构造StreamChatContext
//关键点:callback 被放进 Context,Pipeline 每一步产生的内容都通过 callback.onContent() 推送到前端
StreamChatContext ctx = StreamChatContext.builder()
.question(question)
.conversationId(actualConversationId)
.taskId(taskId)
.deepThinking(thinkingEnabled)
.userId(UserContext.getUserId())
.callback(callback)
.build();
RLock lock = redisson.getLock("ragent:lock:" + actualConversationId);
boolean locked = false;
try {
locked = lock.tryLock(0, 5, TimeUnit.MINUTES);
if (!locked) {
callback.onContent("该对话正在处理中,请稍后再试");
callback.onComplete();
return;
}
//进入Pipeline
chatPipeline.execute(ctx);
} catch (Exception e) {
log.error("流式对话处理异常,会话ID:{},任务ID:{}", actualConversationId, taskId, e);
callback.onError(e);
} finally {
if (locked && lock.isHeldByCurrentThread()) {
try {
lock.unlock();
} catch (Exception ignored) {}
}
}
}
1.如果 Controller 中没有传对话ID,那么在 Service 中使用雪花算法生成
2.使用雪花算法生成任务ID
3.使用工厂模式创建 SSE 事件处理器
除了传入的 SSE 发射器、对话ID、任务ID外,通过工厂模式还传入了:
- modelProperties: 包括提供商配置(URL、密钥、端点映射)、聊天模型组配置(默认模型、深度思考模型、候选模型列表)、向量嵌入模型组配置、重排序模型组配置、模型选择策略配置(失败阈值、熔断器打开持续时间)、流式响应配置(消息分块大小)
问题1:怎么理解熔断器?
问题2:消息分块大小真的用上了吗?
- memoryService: 管理对话记忆,下文详细说
- conversationGroupService: 对话组服务接口,提供对话信息、摘要和对话信息的查询功能 例如获取指定对话中最新的用户消息列表、统计用户在指定对话中的消息数量等
- taskManager: 流式任务管理器
熔断器
熔断器是防止级联失败的保护机制,在项目中:
ModelRoutingExecutor.java - 熔断器核心逻辑
if (!healthStore.allowCall(target.id())) continue; // 熔断器打开,跳过这个模型
具体流程:
- DeepSeek 连续失败2次(failure-threshold: 2)
- 熔断器打开,标记 DeepSeek 为"不可用"
- 接下来的30秒(open-duration-ms: 30000),所有请求自动跳过 DeepSeek
- 直接 fallback 到 SiliconFlow Qwen
- 30秒后熔断器半开,尝试一次 DeepSeek,成功则恢复,失败则继续熔断
为什么需要熔断器?
- ❌ 没有:每次请求都等DeepSeek超时(可能几十秒)→ 用户体验差
- ✅ 有:立即切换,用户无感知 → 系统高可用
消息分块
在 StreamChatEventHandler.java 中,消息分块派上了用场:
private void sendChunked(String type, String content) {
int length = content.length();
int idx = 0;
int count = 0;
StringBuilder buffer = new StringBuilder();
while (idx < length) {
int codePoint = content.codePointAt(idx);
buffer.appendCodePoint(codePoint);
idx += Character.charCount(codePoint);
count++;
if (count >= messageChunkSize) {
sender.sendEvent(SSEEventType.MESSAGE.value(), new MessageDelta(type, buffer.toString()));
buffer.setLength(0);
count = 0;
}
}
if (!buffer.isEmpty()) {
sender.sendEvent(SSEEventType.MESSAGE.value(), new MessageDelta(type, buffer.toString()));
}
}
每达到 messageChunkSize 个字符,就发送一次。避免每个 token 都发送一次 HTTP 请求,造成性能浪费
4.锁的作用:防止同一用户对着一个 conversation 连发两条消息导致同一段对话出现两路并发 SSE 推送。锁保护的是一个 conversation 同时最多一个请求在处理
分布式场景下,一个用户可能对应多个JVM实例!
典型部署架构: 用户 → Nginx负载均衡 → [JVM实例1, JVM实例2, JVM实例3] ↓ Redis(共享)
为什么需要分布式锁?
| 场景 | 无锁问题 | 有锁效果 |
|---|---|---|
| 用户连发两条消息 | JVM1处理消息1,JVM2处理消息2 → 对话记录错乱 | 只有一个JVM能获得锁,另一个直接返回 |
| 多设备登录 | 手机+电脑同时发消息 → 数据竞争 | 保护同一 conversation 同一时刻只有一个请求处理 |
5.StreamChatContent: 流式对话上下文,可理解为保存必要信息的容器。除了不可变的参数外,还有经过 Pipeline 时填充的中间状态,以及各阶段耗时+总耗时
3.进入 Pipeline
1.记录消融实验配置
所谓消融实验,就是用配置文件控制流水线是否经过某个组件的处理,例如关闭问题重写功能、意图解析功能,以便于通过对比获取数据指标
2.加载历史记忆
List<ChatMessage> history = memoryService.loadAndAppend(
ctx.getConversationId(), ctx.getUserId(),
ChatMessage.user(ctx.getQuestion())
);
ctx.setHistory(history);
调用上文中提及的 memoryService 服务类提供的 loadAndAppend 接口
传参:
- 对话ID
- 用户ID
- 用户问题(封装成 ChatMessage 类)
然后将新问题加入到对话历史当中
3.问题拆分 + 重写
if (!ablation.isQueryRewriteEnabled()) {
// 关闭改写:直接使用原问题,不拆分
ctx.setRewriteResult(new RewriteResult(ctx.getQuestion(), List.of()));
ctx.setRewriteMs(System.currentTimeMillis() - t);
return;
}
RewriteResult result = queryRewriteService.rewriteWithSplit(
ctx.getQuestion(), ctx.getHistory()
);
ctx.setRewriteResult(result);
调用 queryRewriteService 中的 rewirteWithSplit 接口
传参:
- 用户问题
- 对话历史
然后将拆分+重写的结果存入 context 的 RewriteResult 中
4.意图解析,采用树状结构意图节点
if (!ablation.isIntentEnabled()) {
// 关闭意图识别:以重写后的问题作为唯一意图,走 KB 检索
SubQuestionIntent defaultIntent = new SubQuestionIntent(
ctx.getRewriteResult().rewrittenQuestion(), List.of());
ctx.setSubIntents(List.of(defaultIntent));
ctx.setIntentMs(System.currentTimeMillis() - t);
return;
}
List<SubQuestionIntent> subIntents = intentResolver.resolve(ctx.getRewriteResult());
ctx.setSubIntents(subIntents);
调用 IntentResolver 的 resolve 接口
传参:
- 重写后的问题
然后将解析出的意图存入 context 的 SubIntents 中
5.歧义引导:意图模糊时,反问用户
if(handleGuideAmbiguity(ctx)) {
ctx.setTotalMs(System.currentTimeMillis() - start);
logTiming(ctx, sw);
return;
}
handleGuideAmbiguity:
String prompt = guidanceService.check(ctx.getQuestion(), ctx.getSubIntents());
if (prompt == null) {
return false; // 意图明确,继续
}
log.info("意图模糊,反问用户: {}", prompt);
ctx.getCallback().onContent(prompt);
ctx.getCallback().onComplete();
return true; // 短路,不检索
调用 GuidanceService 的 check 接口
传参:
- 用户问题
- 子意图
6.判断是否为纯系统问题
if(handleSystemOnlyQuestion(ctx)) {
ctx.setTotalMs(System.currentTimeMillis() - start);
logTiming(ctx, sw);
return;
}
handleSystemOnlyQuestion:
List<SubQuestionIntent> subIntents = ctx.getSubIntents();
if (subIntents != null && !subIntents.isEmpty()) {
boolean hasKbOrMcp = subIntents.stream()
.flatMap(si -> si.nodeScores().stream())
.anyMatch(ns -> ns.getNode() != null &&
(ns.getNode().getKind() == IntentKind.KB ||
ns.getNode().getKind() == IntentKind.MCP));
if (hasKbOrMcp) {
return false; // 有 KB 或 MCP 意图,走检索/工具调用
}
}
log.info("无 KB/MCP 意图匹配,走纯 LLM 回答");
RagResponse(ctx, RetrievalContext.empty());
如果节点类型不是 KB 或 MCP,就直接调用 RagResponse 方法,调用 LLM 回答
走 RagResponse 时,传参:
- context 上下文
- 构造一个空的 Retrive 上下文作为参数传入
7.进入 Agent 循环,代替原有 Pipeline 的 retrieve + response
agentCore.execute(ctx, ctx.getSubIntents(), ctx.getCallback());
传参:
- context 上下文
- 子意图
- SSE 事件处理器
检索
RetrievalContext result = retrievalEngine.retrieve(ctx.getSubIntents(), 5);
调用 retrievalEngine 的 retrieve 接口
二、Pipeline 中各部分的具体实现
1.加载对话记忆:loadAndAppend
default List<ChatMessage> loadAndAppend(String conversationId, String userId, ChatMessage message) {
List<ChatMessage> history = load(conversationId, userId);
append(conversationId, userId, message);
return history;
}
默认接口实现中,实际上就是先 load 后 append,分别看一下 load 和 append 的具体实现:
load
@Override
public List<ChatMessage> load(String conversationId, String userId) {
if (StrUtil.isBlank(conversationId) || StrUtil.isBlank(userId)) return List.of();
try {
// 1. 检查是否有摘要
var summaryDo = groupService.findLatestSummary(conversationId, userId);
String summaryText = (summaryDo != null) ? summaryDo.getContent() : null;
// 2. 取最近消息
List<ConversationMessageVO> messages = messageService.listMessages(
conversationId, userId, DEFAULT_HISTORY_LIMIT, ConversationMessageOrder.ASC);
if (messages == null || messages.isEmpty()) return List.of();
List<ChatMessage> result = new ArrayList<>();
// 3. 有摘要时:作为 SYSTEM 消息放在最前面(节省 token)
if (StrUtil.isNotBlank(summaryText)) {
result.add(ChatMessage.system("以下是之前的对话摘要:\n" + summaryText));
}
// 4. 追加消息
for (ConversationMessageVO vo : messages) {
ChatMessage msg = new ChatMessage(
ChatMessage.Role.fromString(vo.getRole()), vo.getContent());
if (vo.getThinkingContent() != null) msg.setThinkingContent(vo.getThinkingContent());
result.add(msg);
}
return Collections.unmodifiableList(result);
} catch (Exception e) {
log.error("加载对话记忆失败 - conversationId: {}", conversationId, e);
return List.of();
}
}
1.调用 ConversationGroupService 中的 findLatestSummary 接口,检查是否有摘要。如果有,将其作为 SYSTEM 消息
@Override
public ConversationSummaryDO findLatestSummary(String conversationId, String userId) {
if (StrUtil.isBlank(conversationId) || StrUtil.isBlank(userId)) {
return null;
}
return summaryMapper.selectOne(
Wrappers.lambdaQuery(ConversationSummaryDO.class)
.eq(ConversationSummaryDO::getConversationId, conversationId)
.eq(ConversationSummaryDO::getUserId, userId)
.eq(ConversationSummaryDO::getDeleted, 0)
.orderByDesc(ConversationSummaryDO::getId)
.last("limit 1")
);
}
2.调用 ConversationMessageService 中的 listMessages 接口,获取前 DEFAULT_HISTORY_LIMIT 条消息,此处默认是20
TODO: 可以考虑改用 token-aware 的滑动窗口
| 方法 | 当前实现(固定条数) | Token-Aware(动态) |
|---|---|---|
| 策略 | 固定保留最近20条消息 | 根据 token 数动态裁剪 |
| 问题 | 20条短消息浪费空间,20条长消息超 token 限制 | 自动适配不同长度消息 |
| 实现 | DEFAULT_HISTORY_LIMIT = 20 | 估算每条消息 token数,累加上限(如4000 tokens) |
3.将 VO 转换为纯 ChatMessage,加入到 result 中,返回
append
@Override
public String append(String conversationId, String userId, ChatMessage message) {
if (StrUtil.isBlank(conversationId) || StrUtil.isBlank(userId) || message == null) return null;
ConversationMessageBO bo = ConversationMessageBO.builder()
.conversationId(conversationId).userId(userId)
.role(message.getRole().name().toLowerCase())
.content(message.getContent())
.thinkingContent(message.getThinkingContent())
.thinkingDuration(message.getThinkingDuration())
.build();
try {
String msgId = messageService.addMessage(bo);
compressIfNeeded(conversationId, userId, message);
return msgId;
} catch (Exception e) {
log.error("追加消息失败 - conversationId: {}", conversationId, e);
return null;
}
}
1.根据传入参数构建 BO
2.调用 ConversationMessageService 中的 addMessage 接口,实际就是 INSERT
@Override
public String addMessage(ConversationMessageBO conversationMessage) {
ConversationMessageDO messageDO = BeanUtil.toBean(conversationMessage, ConversationMessageDO.class);
conversationMessageMapper.insert(messageDO);
return messageDO.getId();
}
3.若消息数量超过阈值,压缩前几条消息为摘要
compress
private void compressIfNeeded(String conversationId, String userId, ChatMessage message) {
try {
long count = groupService.countUserMessages(conversationId, userId);
if (count < COMPRESS_THRESHOLD) return;
// 取前 COMPRESS_COUNT 条消息 + 新消息 → 压缩
List<ConversationMessageVO> recent = messageService.listMessages(
conversationId, userId, COMPRESS_COUNT, ConversationMessageOrder.ASC);
if (recent == null || recent.isEmpty()) return;
StringBuilder sb = new StringBuilder();
for (ConversationMessageVO m : recent) {
sb.append(m.getRole()).append(": ").append(m.getContent()).append("\n");
}
sb.append("assistant: ").append(message.getContent());
String prompt = "请将以下对话压缩成一段 100 字以内的摘要,保留关键信息和上下文:\n\n" + sb;
String summary = llmService.chat(prompt);
if (StrUtil.isBlank(summary)) return;
messageService.addMessageSummary(ConversationSummaryBO.builder()
.conversationId(conversationId).userId(userId)
.content(summary).build());
log.info("会话 {} 已压缩摘要: {}", conversationId, StrUtil.maxLength(summary, 80));
} catch (Exception e) {
log.warn("会话摘要压缩失败: {}", e.getMessage());
}
}
取前 COMPRESS_COUNT 条消息(此处默认为10条)与新消息 -> 拼接为字符串 -> 调用 LLM 进行压缩 -> 加入摘要
chat 的具体实现:
SimpleLLMService:
@Override
public String chat(ChatRequest request) {
ModelTarget target = selectTarget(false);
ChatClient client = resolveClient(target);
if (client == null) {
throw new IllegalStateException("无可用的 chat 模型");
}
return client.chat(request, target);
}
client 中封装了不同厂商的信息,但chat方法底层都调用 OpenAI 提供的统一接口,以 DeepSeekChatClient 为例:
public class DeepSeekChatClient extends AbstractOpenAIStyleChatClient {
public DeepSeekChatClient(OkHttpClient syncHttpClient,
OkHttpClient streamingHttpClient,
Executor modelStreamExecutor) {
super(syncHttpClient, streamingHttpClient, modelStreamExecutor);
}
@Override
public String provider() {
return "deepseek";
}
@Override
public String chat(ChatRequest request, ModelTarget target) {
return doChat(request, target);
}
@Override
public StreamCancellationHandle streamChat(ChatRequest request, StreamCallBack callback, ModelTarget target) {
return doStreamChat(request, callback, target);
}
}
doChat 与 doStreamChat 的具体实现在 AbstractOpenAIStyleChatClient 类中,以 doChat 为例:
protected String doChat(ChatRequest request, ModelTarget target) {
AIModelProperties.ProviderConfig provider = HttpResponseHandler.requireProvider(target, provider());
if (requiresApiKey()) {
HttpResponseHandler.requireApiKey(provider, provider());
}
JsonObject reqBody = buildRequestBody(request, target, false);
Request requestHttp = newAuthorizedRequest(provider, target)
.post(RequestBody.create(reqBody.toString(), HttpMediaTypes.JSON))
.build();
JsonObject respJson;
try (Response response = syncHttpClient.newCall(requestHttp).execute()) {
if (!response.isSuccessful()) {
String body = HttpResponseHandler.readBody(response.body());
log.warn("{} 同步请求失败: status={}, body={}", provider(), response.code(), body);
throw new ModelClientException(
provider() + " 同步请求失败: HTTP " + response.code(),
ModelClientErrorType.fromHttpStatus(response.code()),
response.code()
);
}
respJson = HttpResponseHandler.parseJson(response.body(), provider());
} catch (IOException e) {
throw new ModelClientException(
provider() + " 同步请求失败: " + e.getMessage(),
ModelClientErrorType.NETWORK_ERROR, null, e);
}
return extractChatContent(respJson);
}
RoutingLLMService:
@Override
public String chat(ChatRequest request) {
return executor.executeWithFallback(
executor.selectCandidates(),
(client, target) -> client.chat(request, target)
);
}
RoutingLLMService vs SimpleLLMService
| 对比项 | SimpleLLMService | RoutingLLMService |
|---|---|---|
| 有多少模型 | 1 个(默认没选到就报错) | N 个候选(优先级 + 熔断) |
| 模型挂了咋办 | 直接抛异常 | 自动切下一个模型 |
| 选型逻辑 | selectTarget() 选一个 | selectCandidates() 排序所有候选 → 逐个试 |
executeWithFallback: DeepSeek 挂了 → 自动试 SiliconFlow Qwen → 挂了试 Ollama → 都挂了就报错
总结:加载最近20条消息 -> 追加新消息 -> 若超10条调用 LLM 进行压缩
2.问题拆分+重写
如果消融实验开关没开,直接使用原问题
正常启用问题重写:
@Override
@RagTraceNode(name = "query-rewrite-and-split", type = "REWRITE")
public RewriteResult rewriteWithSplit(String userQuestion, List<ChatMessage> history) {
if (!ragConfigProperties.getQueryRewriteEnabled()) {
String normalized = queryTermMappingService.normalize(userQuestion);
List<String> subs = ruleBasedSplit(normalized);
return new RewriteResult(normalized, subs);
}
String normalizedQuestion = queryTermMappingService.normalize(userQuestion);
return callLLMRewriteAndSplit(normalizedQuestion, userQuestion, history);
}
normalize: 从 Redis/MySQL 中加载映射规则并进行归一化
逻辑:
- 如果未开启问题重写功能,先归一化,再通过 ruleBasedSplit 按照分隔符拆分问题,最后返回
- 如果开启,先归一化,然后通过 LLM 进行重写,流程为:载入重写提示词 -> 传入提示词、归一化后的问题、对话历史以构建一个 ChatMessage 对象作为发给 LLM 的请求 -> LLM 补全指代、优化表达、将多问句拆分为子问题 -> 把 LLM 返回的 JSON 文本改造成结构化的 RewriteResult 并做好异常兜底(用原问题作单子问题继续往下走)
3.意图解析
如果消融实验开关没开,以重写后的问题作为唯一意图,走 KB 检索
正常启用意图解析:
@RagTraceNode(name = "intent-resolve", type = "INTENT")
public List<SubQuestionIntent> resolve(RewriteResult rewriteResult) {
//判断重写后的问题里的子问题是否为空:非空则直接取,空则取该问题本身
List<String> subQuestions = CollUtil.isNotEmpty(rewriteResult.subQuestions())
? rewriteResult.subQuestions()
: List.of(rewriteResult.rewrittenQuestion());
//并行做意图分类
List<CompletableFuture<SubQuestionIntent>> tasks = subQuestions.stream()
.map(q -> CompletableFuture.supplyAsync( //把每个子问题包装成 CompletableFuture,返回一个“未来会给你结果”的凭据
() -> {
//在线程池里执行
try {
return new SubQuestionIntent(q, classifyIntents(q));
} catch (Exception e) {
log.error("子问题意图分类失败,降级为空意图,question:{}", q, e);
return new SubQuestionIntent(q, List.of());
}
},
intentClassifyExecutor //线程池
))
.toList();
//此时任务已经提交,多个线程在跑
//等待结果,逐个join
//join:拿着凭据,阻塞等待,返回结果
List<SubQuestionIntent> subIntents = tasks.stream()
.map(CompletableFuture::join)
.toList();
//限制总意图数量
return capTotalIntents(subIntents);
}
classifyIntents 调用了 IntentClassifier 中的 classifyTargets,其实现类为 DefaultIntentClassifier
整个类做了三件事:加载意图树 → 分类打分 → 解析结果。按方法拆解:
数据加载
loadIntentTreeData() — 从 Redis 加载意图树,Redis 为空则从数据库加载并回填缓存
返回一个 IntentTreeData record,包含:全量节点列表、叶子节点列表、id → node 的 Map
loadIntentTreeFromDB() — 从 t_intent_node 表查出所有未删除启用的节点(扁平数据),用 intentCode/parentCode 把平表组装成树:
DB 平表:
| code | parent | name |
| hr | null | 人事 |
| sal | hr | 薪资 |
组装后:
hr → children: [sal]
flatten() — 栈式遍历,把树压平成列表。
fillFullPath() — 递归填充 fullPath,让每个节点都存储从根节点到自己的全路径
效果:"集团信息化 > 人事 > 薪资查询",拼进 LLM prompt 时让模型理解分类的层级结构。
核心分类
classifyTargets(String question)
输入用户问题,输出每个叶子节点的匹配分数。流程:
① 加载所有叶子节点(从 Redis/DB)
② buildPrompt:把叶子节点的 id、path、description、examples 拼成 prompt
③ 调 LLM 打分
④ 解析 JSON → List<NodeScore>
⑤ 按 score 降序排序
LLM 返回的 JSON 格式:
[
{"id": "sal", "score": 0.95},
{"id": "leave", "score": 0.12}
]
topKAboveThreshold() — 滤掉低分,只取前 N 个。相当于一个便捷过滤方法
Prompt 构造
buildPrompt(List<IntentNode> leafNodes)
把所有叶子节点拼成一个列表字符串:
- id=sal
path=集团信息化 > 人事 > 薪资查询
description=查询员工薪资、奖金相关信息
type=KB
examples=我这个月工资多少 / 年终奖什么时候发
- id=weather
path=外部服务 > 天气查询
description=查询指定城市的天气
type=MCP
toolId=mcp_tool_001
然后塞进 prompt 模板的 {intent_list} 占位符:
return promptTemplateLoader.render(
INTENT_CLASSIFIER_PROMPT_PATH,
Map.of("intent_list", sb.toString())
);
类型标识会传给 LLM 以让它知道:打高分的如果是 MCP 类型,后续就走工具调用而不是知识库检索
注册表查询
getNodeById(String id) — 实现 IntentNodeRegistry 接口,根据 intentCode 查节点。LLM 返回的 JSON 里只有 id(等于 intentCode),拿到完整节点后才能知道它的 collectionName 或 mcpToolId,供后续检索/工具调用使用
一句话总结
loadIntentTreeFromDB 从数据库拼树,classifyTargets 把叶子节点当选项让 LLM 给每个打分,topKAboveThreshold 取高分结果,后续检索步骤拿到命中的叶子节点,就知道该去哪个集合搜文档
4.歧义引导
/**
* 检查是否需要反问用户
*
* @return 反问文本,null 表示不需要反问
*/
public String check(String question, List<SubQuestionIntent> subIntents) {
// 多个子问题 → 不反问
if (CollUtil.isEmpty(subIntents) || subIntents.size() != 1) {
return null;
}
List<NodeScore> scores = subIntents.get(0).nodeScores();
if (CollUtil.isEmpty(scores) || scores.size() < 2) {
return null;
}
NodeScore top = scores.get(0);
NodeScore second = scores.get(1);
// 最高分太低 → 意图不清
if (top.getScore() < MIN_SCORE) {
return buildPrompt("未能确定您的问题属于哪个类别,请补充更多信息。");
}
// 两名太接近 → 反问
double ratio = second.getScore() / top.getScore();
if (ratio > AMBIGUITY_RATIO) {
String topName = top.getNode().getName();
String secondName = second.getNode().getName();
String msg = String.format("您是想咨询【%s】还是【%s】相关的问题?请补充说明以便我更准确地为您解答。",
topName, secondName);
return buildPrompt(msg);
}
return null; // 意图明确
}
反问情况:
- 最高分太低
- 最高分与次高分太接近
不反问情况:
- 多个子问题
- 意图不多于一个
- 非反问情况
5.ReAct
public void execute(StreamChatContext ctx, List<SubQuestionIntent> subIntents, StreamCallBack callback) {
List<ChatMessage> history = new ArrayList<>();
history.add(ChatMessage.system(buildSystemPrompt()));
history.add(ChatMessage.user(ctx.getQuestion()));
for (int round = 1; round <= MAX_ROUNDS; round++) {
log.info("[Agent] 第 {} 轮...", round);
String decision = llmService.chat(ChatRequest.builder().messages(history).temperature(0.1).build());
log.info("[Agent] LLM 决策: {}", StrUtil.maxLength(decision, 200));
Action action = parse(decision);
if (action == null) {
// LLM 未输出合法 JSON → 清理非法字符后输出
String cleaned = decision.replaceAll("[{}]", "").replaceAll("\"[^\"]*\"", "").trim();
if (cleaned.length() < 5) cleaned = "抱歉,我暂时无法回答这个问题。";
callback.onContent(cleaned);
callback.onComplete();
return;
}
if ("answer".equals(action.type)) {
callback.onContent(action.content);
callback.onComplete();
return;
}
// 执行动作 + 返回包含质量信息的观察结果
String observation = executeAction(action, ctx.getQuestion(), subIntents);
history.add(ChatMessage.assistant(decision));
history.add(ChatMessage.system("【观察结果】\n" + observation));
log.info("[Agent] 第 {} 轮 → {}('{}') → {}", round, action.type, action.content,
StrUtil.maxLength(observation, 120));
}
callback.onContent("抱歉,我暂时无法完成这个任务(超出最大思考轮数)。");
callback.onComplete();
}
流程:
1.构建系统提示词并加入到 history 中
private String buildSystemPrompt() {
String tools = toolRegistry.all().isEmpty() ? "" :
toolRegistry.all().entrySet().stream()
.map(e -> " - " + e.getValue().name() + ": " + e.getValue().description())
.collect(Collectors.joining("\n"));
return """
你是企业内部知识助手。你要通过"思考→行动→观察"的循环来收集信息,然后回答用户问题。
可用的行动:
- search_kb: 搜索企业内部知识库。输入搜索关键词,返回文档片段及其相关度分数。
%s
工作流程:
第1步:思考用户需要什么信息,输出搜索行动
第2步:观察搜索结果(注意检查相关度分数,<0.7说明可能没搜到最相关的内容)
第3步:如果相关度偏低或信息不全 → 换更精准的关键词重新搜索
第4步:如果信息足够 → 输出最终回答
优先级规则:
- 工具返回的实时数据优先于知识库内容
- 知识库内容与工具返回矛盾时,以工具返回为准
- 如果工具参数不足(如需要员工工号但用户没提供),直接反问用户补充,不要反复调工具
格式要求(每次只输出一个JSON,不要有其他文字):
{"tool": "search_kb", "query": "搜索关键词"}
{"answer": "你的最终回答"}
""".formatted(tools);
}
把已注册工具及其描述、硬编码的系统提示词转为字符串
2.将用户问题加入到 history 中
3.LLM 进行决策
4.解析决策,封装成 Action 对象,Action中包含:
- type: 决策类型
- content: 决策内容
5.异常处理:如果 LLM 没有输出合法 JSON,清理非法字符后输出“抱歉,我暂时无法回答这个问题”
6.如果 type 是 answer,说明 LLM 已经完成了工作,直接调用 callback.onContent() 推送内容到前端
7.调用 executeAction
private String executeAction(Action action, String question, List<SubQuestionIntent> subIntents) {
if ("search_kb".equals(action.type)) {
String query = action.content != null && !action.content.isBlank() ? action.content : question;
// 用 LLM 指定的 query 创建临时意图,确保检索使用 LLM 的关键词而非固定 query
List<SubQuestionIntent> dynamicIntents = subIntents.stream()
.map(si -> new SubQuestionIntent(query, si.nodeScores()))
.toList();
RetrievalContext rc = retrievalEngine.retrieve(dynamicIntents, 5);
List<RetrievedChunk> chunks = new ArrayList<>();
if (rc.getIntentChunks() != null) {
for (var entry : rc.getIntentChunks().entrySet()) {
chunks.addAll(entry.getValue());
}
}
if (chunks.isEmpty()) return "检索结果为空,请尝试其他关键词。";
float maxScore = chunks.stream().map(RetrievedChunk::getScore).max(Float::compare).orElse(0f);
chunks = rerankService.rerank(query, chunks);
StringBuilder sb = new StringBuilder();
sb.append(String.format("共命中 %d 篇文档,最高相关度 %.2f。\n", chunks.size(), maxScore));
sb.append("---\n");
for (int i = 0; i < chunks.size(); i++) {
sb.append(String.format("文档%d(相关度%.2f): %s\n",
i + 1, chunks.get(i).getScore(), chunks.get(i).getText()));
}
return sb.toString();
}
Tool tool = toolRegistry.get(action.type);
if (tool != null) {
Map<String, Object> params = mcpExtractor.extract(question, tool);
try {
String result = tool.execute(params);
return "工具调用成功:\n" + result;
} catch (Exception e) {
return "工具调用失败:" + e.getMessage();
}
}
return "未知工具:" + action.type;
}
传入参数:
- LLM 输出的 Action(type + content)
- 用户问题
- 子意图列表
流程:
如果动作是 search_kb:
(a) 构建 query, 如果 action.content 有内容就使用 action.content, 因为这是 LLM 输出的关键词,优先级更高,否则使用用户问题本身。
TODO: 这里是有缺陷的!!!
- ❌ LLM可能输出多余内容
- ❌ query字段可能包含完整句子而非关键词
- ❌ 工程上没有二次清洗机制
改进方案:再次调用 LLM,提取关键词
(b) 创建一个新的子问题意图列表,其中问题使用 (a) 中 LLM 构建的 query, 让 LLM 自主决定检索关键词,而不是用用户原始问题
| 方案 | 检索关键词 | 问题 |
|---|---|---|
| 不用 dynamicIntents | 用户问题:"报销" | 可能召回无关文档("报销流程"、"报销限额"都召回) |
| 用 dynamicIntents | LLM 优化后的 query:"公司差旅报销标准" | 召回更精准 |
(c) 调用 retrieve 进行检索,传入子问题意图列表
(d) 提取出检索出的文档,存入 chunks 中
RetrieveContext 中的 mcpContext 冗余,作为保留字段。因为一开始只计划做 Pipeline, ReAct 是后面新增的功能
(e) 对 chunks 进行精排(rerank)
(f) 遍历 chunks 中的元素,将其转换为“文档编号+相关度+文档内容”的字符串,返回该字符串
如果动作是调用 mcp 工具:
(a) 在已注册工具的列表中找到对应工具
(b) 如果存在这样的工具,通过 mcpExtractor.extract 解析出参数,执行并返回工具调用后的结果(String)
(c) 如果不存在这样的工具,返回提示字符串
不管是检索还是调用工具,最终返回的都是检索/工具调用后的结果,只不过是以字符串形式表示的结果。因为这个结果要作为 observation 交回给 LLM 进行判断,让 LLM 判断相关度是否足够、信息是否齐全,以此来决定下一步是要再次检索/调用工具,还是直接回答。这体现了 ReAct 的 Thought -> Action -> Observation 的循环
8.将 LLM 开始时的决策以及执行检索/工具调用后返回的结果(observation)加入到 history 当中,作为下一轮 ReAct 的信息传给 LLM
9.异常处理:如果超过了最大轮数,调用 callback.onContent() 向前端推送对应字符串
6.检索
public RetrievalContext retrieve(List<SubQuestionIntent> subIntents, int topK) {
if (CollUtil.isEmpty(subIntents)) return RetrievalContext.empty();
int k = topK > 0 ? topK : 5;
Map<String, List<RetrievedChunk>> intentChunks = new HashMap<>();
List<String> contexts = new ArrayList<>();
List<String> mcpResults = new ArrayList<>();
for (SubQuestionIntent intent : subIntents) {
String q = intent.subQuestion();
List<NodeScore> kbIntents = NodeScoreFilters.kb(intent.nodeScores());
List<NodeScore> mcpIntents = NodeScoreFilters.mcp(intent.nodeScores());
// ==== MCP 工具调用 ====
boolean hasMcp = false;
for (NodeScore ns : mcpIntents) {
hasMcp = true;
String toolId = ns.getNode().getMcpToolId();
Tool tool = toolRegistry.get(toolId);
if (tool == null) {
log.warn("MCP 工具不存在: {}", toolId);
continue;
}
Map<String, Object> params = mcpExtractor.extract(q, tool);
String result = tool.execute(params);
mcpResults.add(result);
log.info("MCP 工具调用: {} → {}", toolId, StrUtil.maxLength(result, 200));
}
// 只有 MCP 意图,没有 KB 意图 → 跳过检索
if (hasMcp && kbIntents.isEmpty()) {
continue;
}
// ==== 并行调用所有检索通道 ====
List<CompletableFuture<List<RetrievedChunk>>> futures = channels.stream()
.map(ch -> CompletableFuture.supplyAsync(
() -> ch.search(q, k), retrieveExecutor))
.toList();
List<List<RetrievedChunk>> allResults = futures.stream()
.map(CompletableFuture::join)
.toList();
// ==== RRF 融合 ====
List<RetrievedChunk> fused = rrfFusion(allResults);
// ==== 去重 ====
fused = deduplicate(fused);
// ==== Rerank 精排 ====
fused = rerankService.rerank(q, fused);
// ==== Parent-Child 块聚合 ====
fused = enrichWithParent(fused);
if (CollUtil.isEmpty(fused)) {
log.info("子问题未检索到相关文档:{}", q);
continue;
}
fused = applyPerQuestionBudget(fused);
log.info("子问题检索到 {} 个文档片段({} 通道, 预算裁剪后):{}",
fused.size(), channels.size(), q);
String key = CollUtil.isNotEmpty(kbIntents) ? kbIntents.get(0).getNode().getId() : "default";
intentChunks.put(key, fused);
contexts.add(formatContext(q, fused));
}
contexts = applyTotalBudget(contexts);
return RetrievalContext.builder()
.kbContext(String.join("\n\n---\n\n", contexts))
.mcpContext(String.join("\n", mcpResults))
.intentChunks(intentChunks).build();
}
对每个子问题:
1.过滤出 MCP 工具调用和 KB 检索两种意图
2.对于 MCP 工具调用,尝试从已注册工具中找到其需要的工具,调用 LLM 提取参数,调用工具执行
3.并行调用所有通道
当前只有两个通道实现类:关键词检索与向量检索。调用其 search 接口进行检索,传入参数为:
- q:子问题
- k:topK
4.进行 RRF 融合
传入参数:search 返回的 <List<List<RetrieveChunk>>,即多通道检索后返回的检索块序列
private List<RetrievedChunk> rrfFusion(List<List<RetrievedChunk>> allResults) {
Map<String, RetrievedChunk> merged = new LinkedHashMap<>();
for (List<RetrievedChunk> results : allResults) {
for (int rank = 0; rank < results.size(); rank++) {
String key = textKey(results.get(rank).getText());
double score = 1.0 / (RRF_K + rank + 1);
if (merged.containsKey(key)) {
merged.get(key).setScore(merged.get(key).getScore() + (float) score);
} else {
RetrievedChunk c = results.get(rank);
c.setScore((float) score);
merged.put(key, c);
}
}
}
List<RetrievedChunk> result = new ArrayList<>(merged.values());
result.sort((a, b) -> Float.compare(b.getScore(), a.getScore()));
return result;
}
流程:
1.拿到每一个 chunk 命中的文本(results.get(rank).getText()),如果长度大于40就只取前40个元素作为 key
只要40个的原因:提高去重效率,避免长文本key冲突
如果两个文档开头40个字符相同,可能会被误判为重复,这是工程权衡:
- ✅ 优点:哈希表key更短,内存占用小,比较速度快
- ❌ 缺点:有极小概率误判(两个不同文档前40字符相同)
改进方案: 用内容哈希(MD5/SHA256)作为key,避免误判
2.计算其得分
拿到的 List 已经做了排序,allResults 来源于并行调用所有检索通道时 searchChannel 的 search 接口
SearchChannel 有两个实现类重写了 search 方法,分别是 KeywordSearchChannel 和 VectorSearchChannel
KeywordSearchChannel 中,找到至少命中一个 token 的文档后,按照命中数进行了排序:
// 按命中数排序
return hitCount.entrySet().stream()
.sorted(Map.Entry.<Integer, Integer>comparingByValue().reversed())
.limit(topK)
.map(e -> {
int id = e.getKey();
int hits = e.getValue();
double score = Math.min(1.0, (double) hits / queryTokens.size() + hits * 0.05);
return RetrievedChunk.builder()
.text(documents.get(id))
.score((float) score)
.build();
})
.collect(Collectors.toList());
在 VectorSearchChannel 中,调用了 VectorStoreService 里的 search 方法。VectorStoreService 接口有两个实现类,InMemoryVectorStoreService 和 PgVectorStoreService
在 PgVectorStoreService 中,按照 embedding 后的结果已经做了排序:
List<Row> rows = vectorJdbc.query(
"SELECT id, content, 1 - (embedding <=> ?::vector) AS score FROM t_knowledge_chunk_vector " +
"WHERE enabled = 1 AND deleted = 0 AND embedding IS NOT NULL " +
"ORDER BY embedding <=> ?::vector LIMIT ?",
ps -> {
ps.setString(1, vecStr);
ps.setString(2, vecStr);
ps.setInt(3, topK);
},
(rs, rn) -> new Row(rs.getLong("id"), rs.getString("content"), rs.getDouble("score"))
);
InMemoryVectorStoreService 中,对关键词匹配得分也做了降序排列的处理:
// 3. 按分数降序,取 topK
scored.sort((a, b) -> Double.compare(b.score, a.score));
List<ScoredChunk> top = scored.subList(0, Math.min(topK, scored.size()));
3.如果 merged 里已经有某个 key ,说明这个 chunk 多次命中,给它加分
4.如果 chunk 第一次命中,就设置其 score 字段的值,将其放入哈希表中,value 为 RetrieveChunk 对象本身,key 为其命中文本片段
5.取 merged 的 value 部分作为 result ,降序排列,返回
5.去重
传入参数:进行 RRF 融合后返回的按照 score 排列的 RetrieveChunk 列表
private List<RetrievedChunk> deduplicate(List<RetrievedChunk> chunks) {
if (chunks.size() <= 1) return chunks;
List<RetrievedChunk> result = new ArrayList<>();
result.add(chunks.get(0));
for (int i = 1; i < chunks.size(); i++) {
boolean dup = false;
for (RetrievedChunk kept : result) {
if (textOverlap(chunks.get(i).getText(), kept.getText()) > 0.7) {
dup = true;
break;
}
}
if (!dup) result.add(chunks.get(i));
}
return result;
}
流程:
1.把第一个 Chunk 放入 result 中
2.遍历后续的 chunk
3.对于每一个遍历到的 chunk ,将其与已经在 result 中的 chunk 作比较(执行 textOverlap )
4.如果分数大于0.7,说明 result 中已经有重复度较高的 chunk,需要去掉,不再与已存于 result 中的其它 chunk 比较
5.去掉,实际上就是不放入 result 中
6.如果与 result 中的 chunk 没有明显重复,就将其放入 result 中
7.最终,返回 result
6.Rerank 精排
传入参数:
- 子意图列表
- 去重后的 RetrieveChunk 列表
public List<RetrievedChunk> rerank(String query, List<RetrievedChunk> candidates) {
if (candidates.isEmpty()) return candidates;
JsonObject body = new JsonObject();
body.addProperty("model", MODEL);
body.addProperty("query", query.substring(0, Math.min(query.length(), 256)));
JsonArray docs = new JsonArray();
for (RetrievedChunk c : candidates) {
String text = c.getText();
if (text != null && text.length() > 400) {
text = text.substring(0, 400);
}
docs.add(text != null ? text : "");
}
body.add("documents", docs);
try {
Request req = new Request.Builder()
.url(RERANK_URL)
.header("Authorization", "Bearer " + apiKey)
.post(RequestBody.create(body.toString(), MediaType.parse("application/json")))
.build();
try (Response resp = client.newCall(req).execute()) {
if (!resp.isSuccessful()) {
return candidates;
}
JsonObject result = JsonParser.parseString(resp.body().string()).getAsJsonObject();
JsonArray results = result.getAsJsonArray("results");
List<RetrievedChunk> reranked = new ArrayList<>();
for (int i = 0; i < results.size(); i++) {
JsonObject r = results.get(i).getAsJsonObject();
int idx = r.get("index").getAsInt();
double score = r.get("relevance_score").getAsDouble();
RetrievedChunk c = candidates.get(idx);
c.setScore((float) score);
reranked.add(c);
}
reranked.sort((a, b) -> Float.compare(b.getScore(), a.getScore()));
log.debug("Rerank: {} 个候选重新打分, 最高分={}", reranked.size(),
reranked.isEmpty() ? 0 : String.format("%.2f", reranked.get(0).getScore()));
return reranked;
}
} catch (IOException e) {
log.warn("Rerank 调用失败,使用原始排序: {}", e.getMessage());
return candidates;
}
}
流程:
1.构建 JSON 字符串:
- 模型
- 问题,不长于256
- 文档,每个不长于400
2.根据 JSON 字符串构建 HTTP 请求
3.发送 HTTP 请求(clinet.newCall(req).execute),将响应结果转为 JSON 对象并组装到 JsonArray 中
4.从传入的 candidates 中,通过 idx 找到对应的 chunk,更新其分数,加入 reranked 列表中
5.对 reranked 进行排序,返回 reranked 列表
BGE-Reranker-v2-m3 是专门用于重排序的模型,不是通用LLM
工作原理: 输入:query + [doc1, doc2, doc3, ...] 输出:[{index: 0, relevance_score: 0.95}, {index: 1, relevance_score: 0.12}, ...]
与通用LLM的区别:
| 对比 | 通用LLM(GPT/DeepSeek) | Reranker模型(BGE) |
|---|---|---|
| 输入 | 自然语言prompt | query + documents |
| 输出 | 自由文本 | 结构化分数 |
| 任务 | 生成/理解 | 相似度打分 |
| 速度 | 慢(秒级) | 快(毫秒级) |
不需要 Prompt 或者 schema 告诉它输出格式:
- Reranker 模型是专门训练的,输入格式固定
- 不需要告诉它"判断相关性",它天生就是做这个的
7.Parent-Child 块聚合
传入参数:Rerank 后的 RetrieveChunk 列表
private List<RetrievedChunk> enrichWithParent(List<RetrievedChunk> chunks) {
return chunks.stream().map(chunk -> {
try {
Long parentId = vectorJdbc.queryForObject(
"SELECT parent_id FROM t_knowledge_chunk_vector WHERE id = ?::bigint",
Long.class, chunk.getId());
if (parentId == null || parentId == Long.parseLong(chunk.getId())) return chunk;
String parentText = vectorJdbc.queryForObject(
"SELECT content FROM t_knowledge_chunk_vector WHERE id = ?", String.class, parentId);
if (parentText != null && parentText.length() > chunk.getText().length() + 50) {
chunk.setText(parentText); // 用 Parent 完整上下文替换
}
} catch (Exception ignored) {}
return chunk;
}).toList();
}
作用:检索小文档块,返回完整父文档
场景: 数据库存储: Parent 文档(1000字):完整报销流程规定 └─ Child 文档1(200字):报销流程第1步 └─ Child 文档2(200字):报销流程第2步 └─ Child 文档3(200字):报销流程第3步
检索流程: 用户问:"报销流程是什么?" → 向量检索召回 Child 文档1(200字,片段信息) → enrichWithParent 查询数据库:这个 Child 的 parent_id 是谁? → 获取 Parent 文档(1000字,完整信息) → 用 Parent 替换 Child,返回给 LLM
好处: LLM获得完整上下文,而不是片段信息
三、为什么——技术选型
1.为什么是 Pipeline + ReAct?
回顾一下流程:加载记忆 -> 重写问题 -> 意图解析 -> 歧义引导 -> agent循环
如果只有 ReAct, 显然前面四个流程会消耗较多 token,而且难以保证 LLM 会按照这个流程做。因此对于有固定流程的部分,需要使用 Pipeline
如果只有 Pipeline, 不能处理复杂推理
混合架构——SYSTEM 问题一秒过,KB/MCP 才走 Agent,针对不同复杂度梯度分流
2.为什么是 SSE 而不是 WebSocket?
SSE 轻量、HTTP 天然支持、浏览器原生 EventSource API。WebSocket 需要 upgrade 协议、额外维护连接状态
RAG 是单向流式推送(服务端推给前端),不是双向实时通信,SSE 刚好匹配
3.为什么是 Parent-Child 而不是固定 Chunk?具体是如何实现的?
核心矛盾:检索粒度 vs LLM上下文长度
| 方案 | 检索粒度 | LLM 上下文 | 问题 |
|---|---|---|---|
| 固定 Chunk | 统一500字 | 500字 | ❌ 重要文档被切断,信息不完整 |
| Parent-Child | 检索 Child(200字) | 返回 Parent(1000字) | ✅ 检索精准 + 信息完整 |
具体例子: 用户问:"报销需要什么材料?"
固定Chunk方案: → 召回:"...需要发票、..."(只有片段,没说清楚) → LLM回答:需要发票(信息不完整)
Parent-Child方案: → 召回Child:"需要发票..." → 替换Parent:"报销需要准备:发票原件、审批单、银行卡号..." → LLM回答:需要发票原件、审批单、银行卡号(信息完整)
4.为什么是 Hybrid Retrieval(混合检索)而不是只 Embedding?
| 检索方式 | 优点 | 缺点 | 失败案例 |
|---|---|---|---|
| 向量检索 | 语义理解强 | 对专有名词差 | "GPT-4"搜不到"GPT4"(符号敏感) |
| 关键词检索 | 精确匹配强 | 无法理解语义 | "请假"搜不到"休假申请"(同义词) |
RRF融合: → 合并两个结果,去重,按综合分数排序 → 召回更全更准
5.为什么意图解析采用树状结构存储意图?
让 LLM 理解业务分层,提高分类准确性
| 方案 | LLM看到的Prompt | 问题 |
|---|---|---|
| 扁平列表 | - 薪资查询 - 请假申请 - 天气查询 | ❌ LLM不知道这些意图的关系,容易混淆"薪资查询"和"工资条查看" |
| 树状结构 | - 集团信息化 > 人事 > 薪资查询 - 集团信息化 > 人事 > 请假申请 - 外部服务 > 天气查询 | ✅ LLM理解层级关系,分类更准 |
用户问:"我这个月工资多少?"
扁平列表: → LLM困惑:薪资查询?工资条查看?薪资计算? → 可能分类错误
树状结构: → LLM看到"人事 > 薪资查询",理解这是人事相关 → 分类准确率提升
6.为什么不用 LangGraph ?
其实这个项目使用 LangGraph 是不错的,因为有固定流程,每个流程的状态都会交给下一个流程,类似于 LangGraph 中的 state,本项目中我想自己实现一下 Pipeline,因此没有使用 LangGraph/LangChain 这种框架
7.为什么不用 Redis Stream ?
Redis Stream 是Redis 5.0引入的消息队列功能
与本项目的关系: 当前方案: 用户请求 → Controller → Pipeline → SSE推送(同步等待)
Redis Stream方案: 用户请求 → Controller → Redis Stream消息 → 消费者异步处理 → SSE推送
为什么不用?
- ✅ 当前场景:单机部署,SSE足够,不需要MQ
- ⚠️ 如果需要:异步处理、消息持久化、消费者组,才考虑Redis Stream
8.为什么不用 Milvus ?
| 对比 | PGVector | Milvus |
|---|---|---|
| 本质 | PostgreSQL扩展插件 | 独立向量数据库 |
| 部署 | 依赖PostgreSQL(简单) | 独立部署(复杂) |
| 性能 | 百万级向量够用 | 十亿级向量专精 |
| 功能 | 基础向量检索 | 高级索引、分布式、多模态 |
为什么选PGVector
- ✅ 项目是简化版,不需要Milvus的复杂性
- ✅ 已有PostgreSQL,直接装插件即可
- ✅ 学习成本低,面试容易讲清楚
9.为什么不用 ES ?
ES(Elasticsearch) 是分布式搜索引擎,擅长全文检索
与本项目的关系: 当前关键词检索: → SQL LIKE '%关键词%'(性能差,无排序)
ES关键词检索: → 倒排索引 + BM25算法(性能强,结果准确)
为什么不用?
- ⚠️ 简化版项目,不想引入ES的复杂性
- ✅ 用SQL LIKE作为简化版关键词检索(够用)
10.为什么不用 Qdrant ?
Qdrant 也是向量数据库,类似 Milvus
| 向量数据库 | 语言 | 特点 |
|---|---|---|
| PGVector | SQL | 轻量,依赖PostgreSQL |
| Milvus | Go | 功能强大,适合大规模 |
| Qdrant | Rust | 性能极致,云原生 |
11.你预留了支持 MQ 的接口,为什么考虑使用 MQ ?解决了什么问题?为什么使用 RocketMQ ?
MQ 解决的问题:
| 问题 | 无MQ | 有MQ |
|---|---|---|
| 高并发 | 1000个请求同时来 → 系统崩溃 | 请求进入队列 → 慢慢处理 |
| 异步解耦 | 用户等LLM响应(30秒) | 用户提交 → 立即返回 → 后台处理 |
主流 MQ 对比:
| MQ | 特点 | 适用场景 |
|---|---|---|
| Kafka | 高吞吐(百万/秒) | 大数据、日志采集 |
| RocketMQ | 事务消息、顺序保证 | 金融、订单 |
| RabbitMQ | 轻量、易用 | 中小规模、快速开发 |
本项目为什么预留MQ接口?
- ⚠️ 当前不需要(用户量小)
- ✅ 为未来扩展预留(用户量增长后可平滑升级)
12.为什么使用 Redisson ?
Redisson 是实现分布式锁的 Java 客户端库,底层基于 Redis 的 SETNX + Lua 脚本
分布式锁的必要性: 时间线: T1: 用户A发消息1 → Nginx转发到JVM1 T2: 用户A发消息2 → Nginx转发到JVM2 T3: JVM1和JVM2同时操作同一个conversation → 数据错乱!
解决方案: JVM1获得锁 → JVM2等待 → JVM1处理完释放锁 → JVM2获得锁 → 数据一致
13.为什么使用 PGVector ?
| 组件 | 作用 | 在链路中的位置 |
|---|---|---|
| PGVector | 向量数据库(存储+检索) | 检索阶段:存储文档向量,检索相似文档 |
| BGE Embedding | 向量化模型(文本→向量) | 入库阶段:文档转为向量;检索阶段:问题转为向量 |
完整流程: 入库阶段: 文档 → BGE Embedding → 向量 → PGVector 存储
检索阶段: 用户问题 → BGE Embedding → 问题向量 → PGVector 检索相似文档
14.为什么使用 PostgreSQL ?
原因:PGVector插件只有PostgreSQL支持
PostgreSQL + PGVector插件 = 向量数据库
如果用MySQL:
- ❌ 没有向量存储能力
- ❌ 需要额外部署Milvus/Qdrant
15.为什么不用 GraphRAG ?
GraphRAG = 知识图谱 + RAG
| 对比 | 传统RAG | GraphRAG |
|---|---|---|
| 检索方式 | 向量相似度 | 图谱关系推理 |
| 优势 | 召回快 | 理解实体关系 |
| 案例 | "报销流程"召回文档 | "报销流程需要哪些部门审批"推理出人事→财务链路 |
为什么不用?
- ⚠️ 简化版项目,不需要复杂推理
- ⚠️ 构建知识图谱成本高
16.为什么不用多 Agent?
| 框架 | 特点 | 适用场景 |
|---|---|---|
| CrewAI | 角色扮演,每个 Agent 一个角色 | 复杂任务拆解(研究员+写作员+编辑) |
| AutoGen | 微软开源,对话式协作 | 需要频繁交互的场景 |
| LangGraph | 状态机,有向图 | 固定流程的多 Agent |
为什么本项目不用?
- ✅ 任务流程固定(Pipeline 足够)
- ⚠️ 多Agent成本高(多次 LLM 调用)
- ⚠️ 协作复杂(需要设计 Agent 间通信)
四、为什么——各模块的作用
1.为什么要重写问题?
解决代词指代不明的问题,例如用户问“那个是什么?”,根据聊天记录将“那个”重写为具体名词
2.为什么要解析意图?
便于 ReAct 时 LLM 更好判断意图,节省 LLM 决策成本。同时提高检索召回率
3.为什么粗排后还要 Rerank ?
| 阶段 | 方法 | 速度 | 准确度 | 作用 |
|---|---|---|---|---|
| 粗排 | RRF 融合 + 向量检索 | 快(毫秒级) | 中等 | 从海量文档中快速筛选 Top-k |
| 精排 | BGE Rerank模型 | 慢(秒级) | 高 | 对Top-k重新打分,选出最相关的 |
10万篇文档 → 粗排(RRF)→ 召回前100篇 → 精排→ 选出前5篇最相关的 (快,降低计算量) (准,提高相关性)
4.Trace 模块的作用?
作用:全链路性能追踪 + 问题排查
| 场景 | 无Trace | 有Trace |
|---|---|---|
| 性能分析 | 不知道哪个阶段慢 | 发现"意图解析耗时800ms" |
| 问题排查 | 只能看 log,信息散乱 | 数据库完整记录每一步耗时、错误 |
| 面试讲解 | 只能说"系统运行正常" | 展示数据:检索平均120ms, Rerank 80ms |
例如:
数据库记录:
SELECT * FROM t_rag_trace_node WHERE trace_id = 'xxx';
结果:
- 意图解析:200ms
- 检索:120ms
- Agent决策:800ms
五、异常处理
1.哪些功能用到了 LLM ?如果 LLM 超时怎么办?
用到 LLM 的功能:重写问题、意图解析、Agent决策、解析工具参数、压缩摘要
如果 LLM 超时:ModelRoutingExecutor 逐个试候选模型,DeepSeek 超时 → 试 SiliconFlow → 试 Ollama
若全挂 → 抛异常 → Controller Advice 兜底(GlobalExceptionHandler)
2.Embedding 模型挂了怎么办?
BgeEmbeddingService 返回 null → PgVectorStoreService 跳过无 embedding 的文档 → 不影响已索引的文档。新增文档入库时会失败但不会让整个系统崩
3.PGVector 查询失败怎么办?
PgVectorStoreService.search() catch 异常返回空列表 → 关键词通道仍有结果 → 降级为关键词-only 检索
4.Tool 调用异常怎么办?
AgentCore.executeAction() catch 异常返回 "工具调用失败" → LLM 看到降级信息 → 自动切到纯 KB 策略
5.MCP 不可用怎么办?
ToolRegistry 拿不到 Tool → "Tool 不存在" → LLM 自动跳过工具,走纯检索
6.SSE 中途断开怎么办?
Tomcat 检测到客户端断开(IOException)→ SseEmitter 自动 complete。服务端不会阻塞。前端用 EventSource.onerror 重连
7.用户连续发送两个请求怎么办?
RLock tryLock(conversationId) — 先到先得。后来的立即返回"该对话正在处理中"。不存在数据竞争
更多推荐
所有评论(0)