4、基于 RAG 的数据治理智能问答系统——流式对话接口设计

4、基于 RAG 的数据治理智能问答系统——流式对话接口设计 概述InfoHelper 是一个面向数据治理领域的智能问答系统基于RAG检索增强生成架构结合Elasticsearch 向量检索、Neo4j 图数据库和MySQL 关系型数据库三种数据源实现对数据标准、数据质量、数据资产、数据集市、元数据、数据血缘、数据安全等场景的智能问答。核心入口/chat/send是一个SSEServer-Sent Events流式接口支持逐 token 推送 LLM 生成内容并在各个环节插入进度消息降低用户等待焦虑。接口定义基本信息属性值方法POST路径/chat/sendContent-Typeapplication/x-www-form-urlencoded响应类型text/event-streamSSE传输模式流式Flux请求参数参数类型必填说明示例userIdString是用户唯一标识用于会话归属、权限判断test_user_001contentString是用户问题文本支持中文人员相关的数据有哪些conversationIdString否会话 ID不传则自动创建格式为UUIDuserId1494ea83d8254fffbf63a9964baa791ctest_user_001响应格式响应是一个 SSE 事件流每一条消息以data:开头包含以下三类事件事件类型格式说明进度通知[PROGRESS]:xxx当前处理阶段提示LLM Token纯文本逐 token 流式输出的回答内容结束标记[DONE]:conversationId流结束返回会话 ID 供后续对话使用真实响应示例以下是一次实际请求userIdtest_user_001content人员相关的数据有哪些的 SSE 流输出data:[PROGRESS]:正在识别您的意图... data:[PROGRESS]:正在优化您的问题... data:[PROGRESS]:正在路由您的问题... data:[PROGRESS]:正在检索知识库内容... data:[PROGRESS]:正在排序筛选结果... data:[PROGRESS]:正在生成回答... data:在数据集市中与人员相关的数据主要分布在**教职工管理**、**学生管理**以及**一卡通用户信息**等主题域中。以下是具体的数据实体及其核心字段信息 data:### 1. 教职工基本数据 data:属于**教职工管理主题域BKJG**主要存储教职工的基础身份信息。 data:* **数据实体**RS_StaffInfo data:* **数据大类**教职工基本数据类 data:* **核心字段** data: * StfNo职工号 data: * Name姓名 data: * GenderCode性别码 data:### 2. 学生基本数据 data:* **核心字段** data: * StuNo学号 data: * Name姓名 data: * GenderCode性别码 data: * BirthDate出生日期 data: * MajorCode专业代码 data: * IDCard证件号码 data:* **责任单位**教务处 data:### 3. 卡用户信息数据 data:属于**用户基本信息数据类**关联一卡通或校园卡系统。 data:* **数据实体**YKT_CardUserInfo data:* **核心字段** data: * UserId用户ID data: * Name姓名 data: * CardNo卡号 data: * Balance余额 data:### 4. 项目经费相关人员 data:用于标识项目负责人。 data:* **数据实体**CW_Project data:* **数据大类**项目经费数据类 data:* **核心字段** data: * ChargerId负责人ID关联具体人员 data:[DONE]:1494ea83d8254fffbf63a9964baa791ctest_user_001处理流程整个对话链路可拆分为5 个阶段用户输入: 人员相关的数据有哪些 │ ▼ ┌─────────────────────────────────────────────────────┐ │ 1. 会话管理 │ │ 创建 conversationId UUID userId │ │ 临时标题 人员相关的数据有哪些 │ │ └─ 虚拟线程异步: qwen3.6-flash 生成标题 → 回写DB │ └────────┬────────────────────────────────────────────┘ │ ~0ms ▼ ┌─────────────────────────────────────────────────────┐ │ 2. 消息持久化 │ │ 保存用户消息 (typeUSER, content原始问题) │ │ 创建空助手消息 (typeASSISTANT, contentnull) 占位 │ └────────┬────────────────────────────────────────────┘ │ ~0ms ▼ ┌─────────────────────────────────────────────────────┐ │ 3. 意图识别 (qwen3.7-plus, ~11.5s) │ │ [PROGRESS]:正在识别您的意图... │ │ 输入: 人员相关的数据有哪些 │ │ 输出: intent数据集市查询, data_domain人员 │ │ → 清除意图缓存避免污染后续对话 │ └────────┬────────────────────────────────────────────┘ │ relatedtrue ▼ ┌─────────────────────────────────────────────────────┐ │ 4. RAG 检索增强生成 │ │ │ │ a. 问题改写 (qwen3.7-plus, ~9.3s) │ │ [PROGRESS]:正在优化您的问题... │ │ 输入: 人员相关的数据有哪些 │ │ 输出: 人员相关数据资产查询 │ │ │ │ b. 查询路由 (qwen3.7-plus, ~19.8s) │ │ [PROGRESS]:正在路由您的问题... │ │ 输入: 人员相关数据资产查询 │ │ 输出: strategyrelational_db │ │ │ │ c. 多源检索 (仅激活 relational_db 一路) │ │ [PROGRESS]:正在检索知识库内容... │ │ → Text-to-SQL: 查询 data_metadata 表 │ │ → ES 检索: 补充知识库文档 │ │ │ │ d. 重排序/融合 │ │ [PROGRESS]:正在排序筛选结果... │ │ → RRF 融合 去重取 Top-K │ │ │ │ e. 内容注入 Prompt 加载业务 prompt │ │ → 注入>阶段 1会话管理如果conversationId为空生成格式为UUID userId的会话 ID以用户问题截取前 20 个字符作为临时标题创建会话通过Java 虚拟线程异步调用qwen3.6-flash生成语义化标题回写到数据库不阻塞主流程用户无感知阶段 2意图识别调用qwen3.7-plus加载intent-recognition-new-prompt.txt约 3300 token 的 system prompt判断用户问题是否属于数据治理领域并提取关键实体。实际执行结果输入 “人员相关的数据有哪些”{reasoning:1.相关性涉及数据查询与数据治理相关。2.场景定位用户想了解有哪些与人员相关的数据处于浏览数据集市阶段。3.辨析属于按数据域/主题搜索公开数据目录而非查询自己已申请的数据资产。- 数据集市查询,related:true,intent:数据集市查询,entities:{database_name:null,table_name:null,column_name:null,data_domain:人员,data_owner:null,data_level:null,application_id:null,system_name:null}}支持 8 种意图意图数据来源适用场景示例数据标准咨询知识库了解规范定义“字段命名规范是什么”数据质量检查知识库检查/监控数据质量“数据完整性怎么检查”数据资产查询my_data_table我的查自己已申请的数据“我申请了哪些数据表”数据集市查询data_metadata公开浏览数据目录“人员相关的数据有哪些”数据申请审批data_application申请/查看审批状态“申请单 APP20240101 通过了吗”元数据管理知识库管理元数据“怎么修改表字段注释”数据血缘分析Neo4j追踪数据来源流向“uid 从哪个上游表来的”数据安全合规知识库脱敏/加密/合规“L3 数据需要怎么脱敏”核心区分带我的 → 数据资产查询不带我的且问有什么 → 数据集市查询。阶段 3RAG 检索增强生成如果问题与数据治理相关relatedtrue进入 RAG 管道。核心组件及职责环节组件输入 → 输出耗时参考问题改写InfoHelperQueryTransformer“人员相关的数据有哪些” → “人员相关数据资产查询”~9s查询路由InfoHelperQueryRouter改写后问题 →strategyrelational_db~20sSQL 检索InfoHelperSqlDatabaseContentRetriever自然语言 → SQL → 查询结果~1sES 检索InfoHelperElasticsearchContentRetriever向量/全文检索 → 文档片段~1s重排序BgeScoringModelOnnx检索结果 → 语义排序 Top-K本地推理内容聚合InfoHelperReRankingContentAggregator去重、截断、注入 Prompt~0s生成qwen3.7-plusstreamingPrompt 检索结果 → 逐 token 回答~8s查询路由策略QueryRouter LLM 智能判断走哪路LLM 返回 strategy实际走的检索器典型场景relational_dbSQLMySQL结构化查询、按数据域/主题查元数据graph_dbNeo4j血缘追踪、上下游依赖knowledge_baseES 向量/全文标准咨询、文档语义搜索其他/异常三路并行兜底模糊问题、无法判断注本例中 “人员相关的数据有哪些” 被路由到relational_db因为 QueryRouter 判断这是结构化元数据查询按数据域条件检索。阶段 4权限控制通过RoleEnum控制数据访问范围角色枚举值权限范围普通用户USER只能看公开文档已申请数据用户APPLICANT可见自己申请的数据管理员ADMIN全部可见模型架构用途模型说明意图识别qwen3.7-plus结构化 JSON 输出8 种意图分类问题改写qwen3.7-plus口语化 → 标准检索查询查询路由qwen3.7-plus判断走哪路检索器 置信度Text-to-SQLqwen3.7-plus自然语言 → SQLText-to-Cypherqwen3.7-plus自然语言 → CypherRAG 回答生成qwen3.7-plus结合检索结果流式生成temperature0.2标题生成qwen3.6-flash轻量快速异步非阻塞文本向量化text-embedding-v4ES KNN 向量检索重排序bge-reranker-v2-m3Onnx本地推理语义精排实际执行时序以content人员相关的数据有哪些为例完整链路耗时12:26:15 会话创建 (INSERT chat_conversation) 12:26:16 消息持久化 (INSERT chat_message) 12:26:16 异步标题生成开始 (qwen3.6-flash) 12:26:16 意图识别开始 (qwen3.7-plus) 12:26:28 意图识别完成 → intent数据集市查询, data_domain人员 (~11.5s) 12:26:28 问题改写开始 (qwen3.7-plus) 12:26:37 问题改写完成 → 人员相关数据资产查询 (~9.3s) 12:26:37 查询路由开始 (qwen3.7-plus) 12:26:57 查询路由完成 → strategyrelational_db (~19.8s) 12:27:46 SQL检索 ES检索 排序 注入Prompt 完成 12:27:47 LLM流式生成开始 (qwen3.7-plus, streaming) 12:27:55 流式生成完成回填回答 (UPDATE chat_message) ───────────────────────────────────────── 总耗时: ~100s其中 LLM 调用占 ~48s网络/检索占 ~52s技术要点1. 流式推送 进度通知使用Flux.concatWith()将 RAG 各环节的进度消息与 LLM token 流串联到同一个 Flux 中确保进度消息先于 token 到达前端returnFlux.just([PROGRESS]:正在识别您的意图...).concatWith(/* 意图识别 分支路由 */).concatWith(Mono.just([DONE]:finalConversationId));2. 多数据源指路不盲搜四路检索器全部构建但QueryRouter 只激活 1 路常规情况避免无效检索。只有路由失败时才走三路并行兜底。3. 异步标题生成会话创建后不阻塞主流程通过 Java 虚拟线程异步调用qwen3.6-flash生成标题Thread.ofVirtual().name(title-summary-conversationId).start(()-{StringaiTitletitleSummaryService.generateTitle(content);updateConversationTitle(conversationId,aiTitle);});实际效果临时标题 “人员相关的数据有哪些” → 语义标题 “人员相关数据有哪些”。4. 意图缓存隔离意图识别完成后立即清除 LLM 的历史记忆缓存避免意图识别的对话上下文污染后续 RAG 对话databaseChatMemoryStore.evictCache(finalConversationId);5. 消息占位-回填机制流式返回时无法预知最终内容先保存空消息占位获取messageId流结束后再回填完整内容INSERT chat_message (contentnull) → 获取 messageId → LLM 流式生成 → UPDATE chat_message SET content完整回答前端集成示例asyncfunctionchat(userId,content,conversationId){constresponseawaitfetch(/chat/send,{method:POST,headers:{Content-Type:application/x-www-form-urlencoded},body:newURLSearchParams({userId,content,conversationId})});constreaderresponse.body.getReader();constdecodernewTextDecoder();letbuffer;while(true){const{done,value}awaitreader.read();if(done)break;bufferdecoder.decode(value,{stream:true});constlinesbuffer.split(\n);bufferlines.pop();for(constlineoflines){if(line.startsWith(data:[PROGRESS])){showProgress(line.replace(data:,));}elseif(line.startsWith(data:[DONE])){constconvIdline.split(:)[2];saveConversationId(convId);}elseif(line.startsWith(data:)){appendToken(line.replace(data:,));}}}}总结/chat/send是 InfoHelper 的唯一对话入口封装了意图识别、多源 RAG 检索、流式生成、会话管理、权限控制等完整链路。通过 SSE 进度通知机制用户可以在每个处理阶段看到实时反馈。整个架构的设计亮点意图识别做第一道分拣无关问题走通用对话兜底降低 RAG 成本指路不盲搜QueryRouter 智能判断走哪路检索器只有失败时才三路并行兜底BGE Reranker 本地重排序Onnx 模型本地推理无需网络调用提升结果精度进度消息与 LLM 输出同流推送前端只需一个 SSE 连接降低复杂度全链路耗时透明意图识别 ~11s、改写 ~9s、路由 ~20s、生成 ~8s便于性能优化定位