4、基于 RAG 的数据治理智能问答系统——流式对话接口设计
2026/7/22 13:25:41 网站建设 项目流程

概述

InfoHelper 是一个面向数据治理领域的智能问答系统,基于RAG(检索增强生成)架构,结合Elasticsearch 向量检索Neo4j 图数据库MySQL 关系型数据库三种数据源,实现对数据标准、数据质量、数据资产、数据集市、元数据、数据血缘、数据安全等场景的智能问答。

核心入口/chat/send是一个SSE(Server-Sent Events)流式接口,支持逐 token 推送 LLM 生成内容,并在各个环节插入进度消息,降低用户等待焦虑。


接口定义

基本信息

属性
方法POST
路径/chat/send
Content-Typeapplication/x-www-form-urlencoded
响应类型text/event-stream(SSE)
传输模式流式(Flux)

请求参数

参数类型必填说明示例
userIdString用户唯一标识,用于会话归属、权限判断"test_user_001"
contentString用户问题文本,支持中文"人员相关的数据有哪些"
conversationIdString会话 ID,不传则自动创建,格式为UUID+userId"1494ea83d8254fffbf63a9964baa791ctest_user_001"

响应格式

响应是一个 SSE 事件流,每一条消息以data:开头,包含以下三类事件:

事件类型格式说明
进度通知[PROGRESS]:xxx当前处理阶段提示
LLM Token纯文本逐 token 流式输出的回答内容
结束标记[DONE]:conversationId流结束,返回会话 ID 供后续对话使用

真实响应示例

以下是一次实际请求userId=test_user_001&content=人员相关的数据有哪些的 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. 消息持久化 │ │ 保存用户消息 (type=USER, content=原始问题) │ │ 创建空助手消息 (type=ASSISTANT, content=null) 占位 │ └────────┬────────────────────────────────────────────┘ │ ~0ms ▼ ┌─────────────────────────────────────────────────────┐ │ 3. 意图识别 (qwen3.7-plus, ~11.5s) │ │ [PROGRESS]:正在识别您的意图... │ │ 输入: "人员相关的数据有哪些" │ │ 输出: intent=数据集市查询, data_domain=人员 │ │ → 清除意图缓存,避免污染后续对话 │ └────────┬────────────────────────────────────────────┘ │ related=true ▼ ┌─────────────────────────────────────────────────────┐ │ 4. RAG 检索增强生成 │ │ │ │ a. 问题改写 (qwen3.7-plus, ~9.3s) │ │ [PROGRESS]:正在优化您的问题... │ │ 输入: "人员相关的数据有哪些" │ │ 输出: "人员相关数据资产查询" │ │ │ │ b. 查询路由 (qwen3.7-plus, ~19.8s) │ │ [PROGRESS]:正在路由您的问题... │ │ 输入: "人员相关数据资产查询" │ │ 输出: strategy=relational_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 数据需要怎么脱敏”

核心区分:带"我的" → 数据资产查询;不带"我的"且问"有什么" → 数据集市查询。

阶段 3:RAG 检索增强生成

如果问题与数据治理相关(related=true),进入 RAG 管道。核心组件及职责:

环节组件输入 → 输出耗时参考
问题改写InfoHelperQueryTransformer“人员相关的数据有哪些” → “人员相关数据资产查询”~9s
查询路由InfoHelperQueryRouter改写后问题 →strategy=relational_db~20s
SQL 检索InfoHelperSqlDatabaseContentRetriever自然语言 → SQL → 查询结果~1s
ES 检索InfoHelperElasticsearchContentRetriever向量/全文检索 → 文档片段~1s
重排序BgeScoringModel(Onnx)检索结果 → 语义排序 Top-K本地推理
内容聚合InfoHelperReRankingContentAggregator去重、截断、注入 Prompt~0s
生成qwen3.7-plus(streaming)Prompt + 检索结果 → 逐 token 回答~8s

查询路由策略(QueryRouter LLM 智能判断走哪路):

LLM 返回 strategy实际走的检索器典型场景
relational_dbSQL(MySQL)结构化查询、按数据域/主题查元数据
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自然语言 → SQL
Text-to-Cypherqwen3.7-plus自然语言 → Cypher
RAG 回答生成qwen3.7-plus结合检索结果流式生成,temperature=0.2
标题生成qwen3.6-flash轻量快速,异步非阻塞
文本向量化text-embedding-v4ES KNN 向量检索
重排序bge-reranker-v2-m3(Onnx)本地推理,语义精排

实际执行时序

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 查询路由完成 → strategy=relational_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(()->{StringaiTitle=titleSummaryService.generateTitle(content);updateConversationTitle(conversationId,aiTitle);});

实际效果:临时标题 “人员相关的数据有哪些” → 语义标题 “人员相关数据有哪些”。

4. 意图缓存隔离

意图识别完成后立即清除 LLM 的历史记忆缓存,避免意图识别的对话上下文污染后续 RAG 对话:

databaseChatMemoryStore.evictCache(finalConversationId);

5. 消息占位-回填机制

流式返回时无法预知最终内容,先保存空消息占位获取messageId,流结束后再回填完整内容:

INSERT chat_message (content=null) → 获取 messageId → LLM 流式生成 → UPDATE chat_message SET content=完整回答

前端集成示例

asyncfunctionchat(userId,content,conversationId){constresponse=awaitfetch('/chat/send',{method:'POST',headers:{'Content-Type':'application/x-www-form-urlencoded'},body:newURLSearchParams({userId,content,conversationId})});constreader=response.body.getReader();constdecoder=newTextDecoder();letbuffer='';while(true){const{done,value}=awaitreader.read();if(done)break;buffer+=decoder.decode(value,{stream:true});constlines=buffer.split('\n');buffer=lines.pop();for(constlineoflines){if(line.startsWith('data:[PROGRESS]')){showProgress(line.replace('data:',''));}elseif(line.startsWith('data:[DONE]')){constconvId=line.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,便于性能优化定位

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询