1. 项目概述:构建基于LangGraph的智能问数系统
在数据驱动的商业环境中,业务人员每天需要处理大量数据查询需求。传统的数据查询方式通常要求用户掌握SQL等专业查询语言,或者依赖IT部门提供报表支持,这种模式存在响应慢、沟通成本高等问题。我们基于LangGraph框架开发的智能问数系统,旨在通过自然语言交互的方式,让业务人员可以直接用日常语言提问,系统自动理解并返回所需数据。
这个项目的核心挑战在于如何准确理解用户的自然语言问题,并将其转换为数据系统能够处理的元数据指令。就像一位经验丰富的数据分析师需要理解业务术语背后的真实含义一样,我们的系统需要具备类似的语义理解能力。
2. 系统架构设计思路
2.1 整体架构分层
我们的智能问数系统采用分层架构设计,每层都有明确的职责边界:
- 交互层:处理用户输入和结果展示
- 语义理解层(本文重点):将自然语言转换为结构化查询要素
- 查询生成层:根据理解结果构建数据查询
- 执行层:执行查询并返回结果
2.2 语义理解层的核心任务
语义理解层需要完成从自然语言到数据元数据的转换,主要包括四个关键步骤:
- 关键词抽取:提取问题中的核心业务术语
- 字段召回:识别问题涉及的数据库字段
- 指标召回:确定需要计算的业务指标
- 维度取值召回:获取具体的筛选条件值
这种分步处理的方式类似于人类理解问题的过程:先抓住关键词,再明确要查什么字段,计算什么指标,以及按什么条件筛选。
3. 关键词抽取节点实现
3.1 技术选型与实现
我们选择jieba分词库作为基础工具,主要考虑以下因素:
- 成熟稳定的中文分词能力
- 支持自定义词典和停用词表
- 提供关键词提取和词性标注功能
python复制import jieba.analyse
from langgraph.runtime import Runtime
async def extract_keywords(state: DataAgentState, runtime: Runtime[DataAgentContext]):
writer = runtime.stream_writer
writer({"type": "progress", "step": "抽取关键字", "status": "running"})
query = state["query"]
# 定义需要保留的词性
allow_pos = (
"n", "nr", "ns", "nt", "nz", # 名词类
"v", "vn", # 动词类
"a", "an", # 形容词类
"eng", "i", "l" # 英文、成语、习惯用语
)
# 提取关键词
keywords = jieba.analyse.extract_tags(query, allowPOS=allow_pos)
# 添加原始query作为兜底
keywords = list(set(keywords + [query]))
writer({"type": "progress", "step": "抽取关键字", "status": "success"})
return {"keywords": keywords}
3.2 关键设计决策
-
词性过滤策略:我们只保留名词、动词、形容词等有实际意义的词性,过滤掉助词、介词等对数据查询无意义的词汇。这种过滤可以显著减少后续处理的噪音。
-
原始Query保留:作为安全措施,我们将用户原始问题也加入关键词集合。这确保即使关键词提取不完全,系统也不会丢失核心查询意图。
-
流式处理反馈:通过流式接口实时返回处理状态,提升用户体验。用户可以看到系统正在"思考"的过程,而不是长时间等待。
提示:在实际应用中,建议根据业务领域调整allow_pos词性过滤规则。例如,金融领域可能需要特别关注数字和货币相关词汇。
4. 字段召回节点实现
4.1 召回流程设计
字段召回是将关键词映射到数据库实际字段的过程,我们采用多阶段召回策略:
- 关键词扩展:使用LLM对原始关键词进行同义词和关联词扩展
- 向量化处理:将扩展后的关键词转换为向量表示
- 向量检索:在Qdrant向量库中查找相似字段
- 结果去重:合并来自不同关键词的召回结果
python复制async def recall_column(state: DataAgentState, runtime: Runtime[DataAgentContext]):
writer = runtime.stream_writer
writer({"type": "progress", "step": "召回字段", "status": "running"})
query = state["query"]
keywords = state["keywords"]
embedding_client = runtime.context["embedding_client"]
column_qdrant_repository = runtime.context["column_qdrant_repository"]
try:
# LLM关键词扩展
prompt = PromptTemplate(
template=load_prompt("extend_keywords_for_column_recall"),
input_variables=["query"],
)
chain = prompt | llm | JsonOutputParser()
result = await chain.ainvoke({"query": query})
# 合并关键词
retrieved_columns_map = {}
keywords = list(set(keywords + result))
# 向量检索
for keyword in keywords:
embedding = await embedding_client.aembed_query(keyword)
payloads = await column_qdrant_repository.search(embedding)
for payload in payloads:
if payload.id not in retrieved_columns_map:
retrieved_columns_map[payload.id] = payload
writer({"type": "progress", "step": "召回字段", "status": "success"})
return {"retrieved_columns": list(retrieved_columns_map.values())}
except Exception as e:
writer({"type": "progress", "step": "召回字段", "status": "error"})
raise
4.2 向量检索实现细节
Qdrant向量数据库的检索实现考虑了以下因素:
python复制async def search(self, embedding: list[float], score_threshold=0.6, limit=5) -> list[ColumnInfo]:
result = await self.client.query_points(
collection_name=self.collection_name,
query=embedding,
score_threshold=score_threshold,
limit=limit
)
return [ColumnInfo(**point.payload) for point in result.points]
- 相似度阈值:设置0.6的阈值过滤低质量匹配
- 返回数量限制:每次查询最多返回5个最相关字段
- 结构化返回:将结果封装为ColumnInfo对象,便于后续处理
5. 指标召回节点实现
5.1 与字段召回的异同
指标召回与字段召回在架构上保持高度一致,主要区别在于:
- 使用不同的提示词模板(extend_keywords_for_metric_recall)
- 连接不同的向量库集合(指标而非字段)
- 返回MetricInfo而非ColumnInfo对象
这种设计确保了系统的一致性和可维护性,新增功能时只需最小化修改。
python复制async def recall_metric(state: DataAgentState, runtime: Runtime[DataAgentContext]):
# ...与字段召回类似的结构...
prompt = PromptTemplate(
template=load_prompt("extend_keywords_for_metric_recall"),
input_variables=["query"]
)
# ...后续处理流程相同...
5.2 统一架构的优势
- 代码复用:核心逻辑只需实现一次
- 维护简便:bug修复和性能优化可以统一应用
- 扩展容易:新增召回类型只需配置不同参数
- 认知一致:开发人员学习曲线平缓
6. 维度取值召回节点实现
6.1 技术选型考量
维度取值召回面临不同的技术挑战:
- 取值通常是短文本(如"华北"、"Q1")
- 需要支持模糊匹配(如"华"匹配"华北")
- 对错别字和同义词要有容忍度
- 查询延迟要求更高
基于这些特点,我们选择Elasticsearch而非向量数据库:
- 更擅长短文本检索
- 提供丰富的文本分析功能
- 支持模糊查询和同义词扩展
- 检索性能优异
python复制async def recall_value(state: DataAgentState, runtime: Runtime[DataAgentContext]):
writer = runtime.stream_writer
writer({"type": "progress", "step": "召回字段取值", "status": "running"})
query = state["query"]
keywords = state["keywords"]
value_es_repository = runtime.context["value_es_repository"]
try:
# LLM关键词扩展
prompt = PromptTemplate(
template=load_prompt("extend_keywords_for_value_recall"),
input_variables=["query"]
)
chain = prompt | llm | JsonOutputParser()
result = await chain.ainvoke({"query": query})
values_map = {}
keywords = list(set(keywords + result))
# ES全文检索
for keyword in keywords:
values = await value_es_repository.search(keyword)
for value in values:
if value.id not in values_map:
values_map[value.id] = value
writer({"type": "progress", "step": "召回字段取值", "status": "success"})
return {"retrieved_values": list(values_map.values())}
except Exception as e:
writer({"type": "progress", "step": "召回字段取值", "status": "error"})
raise
6.2 Elasticsearch检索配置
ES检索实现考虑了维度值的特点:
python复制async def search(self, keyword: str, score_threshold=0.6, limit=5) -> list[ValueInfo]:
result = await self.client.search(
index=self.index_name,
query={"match": {"value": keyword}},
min_score=score_threshold,
size=limit
)
return [ValueInfo(**hit['_source']) for hit in result['hits']['hits']]
- 使用match查询而非term查询,支持文本分析
- 设置最小分数阈值过滤低质量匹配
- 限制返回结果数量保证性能
- 将原始文档转换为ValueInfo对象
7. 工程实践与性能优化
7.1 异步处理设计
所有节点都采用异步实现,主要考虑:
- 高并发支持:避免IO等待阻塞线程
- 资源利用率:在等待外部服务响应时可以处理其他请求
- 响应速度:可以并行执行独立任务
关键实现要点:
- 使用async/await语法
- 选择支持异步的客户端库
- 合理控制并发度
7.2 流式进度反馈
通过流式接口实时返回处理状态:
python复制writer({"type": "progress", "step": "召回字段", "status": "running"})
# ...处理逻辑...
writer({"type": "progress", "step": "召回字段", "status": "success"})
状态机设计:
- running:开始处理
- success:成功完成
- error:处理失败
7.3 错误处理与日志
统一的错误处理模式:
- 捕获所有异常
- 更新状态为error
- 记录详细错误日志
- 重新抛出异常(由上层处理)
python复制try:
# 业务逻辑
except Exception as e:
writer({"type": "progress", "step": "召回字段", "status": "error"})
logger.error(f"召回字段失败: {str(e)}")
raise
日志记录要点:
- 包含足够上下文信息
- 区分不同日志级别
- 结构化日志便于分析
8. 实际应用中的经验总结
8.1 关键词抽取优化
在实际应用中,我们发现jieba的默认词典可能无法覆盖所有业务术语。通过以下措施显著提升了效果:
- 加载业务词典:将业务高频术语加入自定义词典
- 调整词频:对重要业务词提高词频权重
- 停用词优化:根据业务特点调整停用词表
python复制# 初始化时加载自定义词典
jieba.load_userdict("data/business_terms.txt")
8.2 向量检索调优
向量检索效果对系统准确性至关重要,我们通过以下方式优化:
- 嵌入模型选择:测试不同模型在业务领域的表现
- 分数阈值调整:通过实验确定最佳阈值
- 元数据过滤:结合业务规则缩小检索范围
- 混合检索:结合关键词和向量检索
8.3 性能监控指标
为确保系统稳定运行,我们监控以下关键指标:
- 各节点耗时:识别性能瓶颈
- 召回准确率:评估语义理解质量
- 错误率:及时发现系统问题
- 缓存命中率:评估缓存效果
提示:建议为高频查询建立缓存机制,可以显著提升响应速度。但要注意缓存失效策略,确保数据时效性。
9. 常见问题与解决方案
9.1 关键词抽取不准确
问题表现:
- 遗漏重要业务术语
- 包含无关词汇
解决方案:
- 扩充自定义词典
- 调整词性过滤规则
- 添加业务特定的后处理逻辑
9.2 字段召回结果不相关
问题表现:
- 返回字段与问题意图不符
- 漏掉明显相关字段
解决方案:
- 优化关键词扩展提示词
- 调整向量相似度阈值
- 检查字段元数据质量
9.3 系统响应慢
问题表现:
- 用户查询延迟高
- 系统吞吐量低
解决方案:
- 增加异步并发度
- 实现分级缓存
- 优化向量索引配置
- 考虑预计算常见查询
10. 扩展与演进方向
当前系统已经实现了基础语义理解能力,还可以从以下方向扩展:
- 多轮对话支持:处理需要澄清或补充信息的查询
- 查询意图分类:预先识别问题类型(趋势分析、对比等)
- 结果解释:提供查询结果的业务解释而不仅是原始数据
- 反馈学习:根据用户反馈持续优化模型
在实际部署中,我们发现将业务规则与机器学习相结合往往能取得最佳效果。纯粹的机器学习方法可能在边界情况下表现不稳定,而结合业务规则可以提供确定性的保障。
