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

对比项SimpleLLMServiceRoutingLLMService
有多少模型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用户问题:"报销"可能召回无关文档("报销流程"、"报销限额"都召回)
用 dynamicIntentsLLM 优化后的 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)
输入自然语言promptquery + 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 ?

对比PGVectorMilvus
本质PostgreSQL扩展插件独立向量数据库
部署依赖PostgreSQL(简单)独立部署(复杂)
性能百万级向量够用十亿级向量专精
功能基础向量检索高级索引、分布式、多模态

为什么选PGVector

  • ✅ 项目是简化版,不需要Milvus的复杂性
  • ✅ 已有PostgreSQL,直接装插件即可
  • ✅ 学习成本低,面试容易讲清楚

9.为什么不用 ES ?

ES(Elasticsearch) 是分布式搜索引擎,擅长全文检索

与本项目的关系: 当前关键词检索: → SQL LIKE '%关键词%'(性能差,无排序)

ES关键词检索: → 倒排索引 + BM25算法(性能强,结果准确)

为什么不用?

  • ⚠️ 简化版项目,不想引入ES的复杂性
  • ✅ 用SQL LIKE作为简化版关键词检索(够用)

10.为什么不用 Qdrant ?

Qdrant 也是向量数据库,类似 Milvus

向量数据库语言特点
PGVectorSQL轻量,依赖PostgreSQL
MilvusGo功能强大,适合大规模
QdrantRust性能极致,云原生

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

对比传统RAGGraphRAG
检索方式向量相似度图谱关系推理
优势召回快理解实体关系
案例"报销流程"召回文档"报销流程需要哪些部门审批"推理出人事→财务链路

为什么不用?

  • ⚠️ 简化版项目,不需要复杂推理
  • ⚠️ 构建知识图谱成本高

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) — 先到先得。后来的立即返回"该对话正在处理中"。不存在数据竞争

更多推荐