1. 项目概述:智能问答系统的路由决策挑战
在企业级智能问答场景中,最令人头疼的问题莫过于:当用户抛出一个问题时,系统该如何判断应该去文档库检索还是查询数据库?这个看似简单的决策背后,实际上涉及到自然语言理解、任务分类和工具调用的复杂过程。
传统解决方案通常采用硬编码规则或简单关键词匹配,但这种方法存在明显缺陷:
- 规则维护成本高:每次新增文档类型或数据库表都需要更新规则
- 泛化能力差:无法处理同义表达和复杂语义
- 容错性低:一旦判断错误,整个回答就会偏离预期
我们开发的RAG-SQL Router系统,正是为了解决这些痛点。其核心创新点在于:
- 动态路由机制:基于问题语义而非关键词进行工具选择
- 混合执行能力:可同时调用多个工具处理复合问题
- 质量保障层:通过Cleanlab Codex对输出结果进行可信度评估
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 系统架构深度解析
2.1 核心组件交互流程
系统采用事件驱动的Workflow设计模式,主要组件及其交互关系如下:
mermaid复制graph TD
A[用户问题] --> B(Router Agent)
B --> C{RAG工具}
B --> D{SQL工具}
C --> E[文档检索结果]
D --> F[数据库查询结果]
E --> G[结果整合]
F --> G
G --> H[Cleanlab验证]
H --> I[最终响应]
2.2 路由智能体实现细节
路由决策的核心在于LLM对工具描述的准确理解。我们采用函数调用(function calling)方式,为每个工具定义清晰的元数据:
python复制tools_def = [
{
"type": "function",
"function": {
"name": "sql_query_engine",
"description": (
"用于查询结构化数据,适合回答需要精确数值的问题。"
"示例问题:'Q3销售额最高的产品是什么?'"
"注意:仅当问题明确涉及数据库中的表字段时使用"
),
"parameters": {...}
}
},
# RAG工具定义类似
]
关键设计要点:
- 描述差异化:确保两种工具的描述没有语义重叠
- 示例引导:包含典型问题示例帮助LLM理解
- 限制条件:明确说明使用边界
2.3 混合查询执行流程
对于需要同时使用两种工具的复合问题,系统执行以下步骤:
- 问题分解:识别问题中的独立子问题
- 并行执行:异步调用多个工具
- 结果聚合:按逻辑关系组合不同来源的结果
典型场景示例:
用户问题:"东京的人口是多少?文档中关于东京的介绍有哪些?"
处理流程:
python复制async def handle_compound_question(question):
# 步骤1:问题分解
sub_questions = await llm.decompose(question)
# 步骤2:并行查询
tasks = [
router.run(q) for q in sub_questions
]
results = await asyncio.gather(*tasks)
# 步骤3:结果整合
final_answer = await llm.summarize(results)
return final_answer
3. 核心实现详解
3.1 环境配置与依赖管理
推荐使用Poetry进行依赖管理,比直接使用pip更可靠:
python复制# pyproject.toml
[tool.poetry.dependencies]
python = "^3.9"
llama-index = "^0.10.0"
sqlalchemy = "^2.0.0"
cleanlab-codex = "^1.2.0"
[build-system]
requires = ["poetry-core>=1.0.0"]
build-backend = "poetry.core.masonry.api"
关键依赖说明:
llama-index:提供核心的RAG和SQL查询能力sqlalchemy:数据库ORM层,支持多种数据库后端cleanlab-codex:回答质量验证工具
3.2 数据库连接最佳实践
生产环境推荐使用连接池和管理器:
python复制from sqlalchemy import create_engine
from sqlalchemy.pool import QueuePool
engine = create_engine(
"postgresql://user:pass@host:5432/db",
poolclass=QueuePool,
pool_size=10,
max_overflow=5,
pool_pre_ping=True,
pool_recycle=3600,
connect_args={
"connect_timeout": 5,
"application_name": "rag-sql-router"
}
)
连接池参数说明:
pool_size:保持的连接数max_overflow:允许临时超过的连接数pool_pre_ping:自动检测连接有效性pool_recycle:连接自动回收时间(秒)
3.3 RAG索引优化技巧
对于文档处理,建议采用分层索引策略:
python复制from llama_index.core import VectorStoreIndex
from llama_index.core.node_parser import HierarchicalNodeParser
node_parser = HierarchicalNodeParser(
chunk_sizes=[2048, 512, 128], # 三级分块
chunk_overlap=0.2
)
documents = load_documents() # 加载原始文档
nodes = node_parser.get_nodes_from_documents(documents)
# 创建多粒度索引
base_index = VectorStoreIndex(nodes[2]) # 最细粒度
mid_index = VectorStoreIndex(nodes[1])
top_index = VectorStoreIndex(nodes[0])
这种分层结构可以:
- 提高检索精度:细粒度匹配具体内容
- 保持上下文:粗粒度提供背景信息
- 优化性能:减少不必要的细粒度检索
4. 生产环境部署方案
4.1 性能优化策略
查询缓存实现
python复制from functools import lru_cache
import hashlib
@lru_cache(maxsize=1000)
def cached_query(query: str, db_schema: str) -> str:
"""带数据库模式感知的缓存"""
pass
def get_query_hash(query: str, db_schema: str) -> str:
"""考虑数据库结构的哈希生成"""
return hashlib.sha256(
f"{query}||{db_schema}".encode()
).hexdigest()
缓存策略特点:
- 模式感知:当数据库结构变化时自动失效
- 多级缓存:内存缓存+Redis持久化缓存
- 语义哈希:相同语义不同表达的问题能命中缓存
负载均衡设计
python复制from concurrent.futures import ThreadPoolExecutor
from collections import defaultdict
class LoadBalancer:
def __init__(self):
self.executors = {
"rag": ThreadPoolExecutor(max_workers=8),
"sql": ThreadPoolExecutor(max_workers=4)
}
self.metrics = defaultdict(int)
def dispatch(self, task_type: str, fn, *args):
self.metrics[task_type] += 1
return self.executors[task_type].submit(fn, *args)
4.2 监控与告警系统
推荐使用Prometheus+Grafana搭建监控看板,关键指标包括:
python复制from prometheus_client import Counter, Gauge
# 定义指标
QUERY_COUNT = Counter(
'router_queries_total',
'Total query count',
['tool_type']
)
LATENCY = Gauge(
'router_latency_seconds',
'Query processing latency',
['tool_type']
)
ERROR_COUNT = Counter(
'router_errors_total',
'Total error count',
['error_type']
)
# 在关键位置埋点
def process_query(query):
start_time = time.time()
try:
result = router.run(query)
QUERY_COUNT.labels(tool_type=result.tool).inc()
return result
except Exception as e:
ERROR_COUNT.labels(error_type=type(e).__name__).inc()
raise
finally:
LATENCY.labels(tool_type=result.tool).set(
time.time() - start_time
)
4.3 安全防护措施
SQL注入防护
python复制from sqlalchemy import text
from sqlalchemy.exc import SQLAlchemyError
def safe_query(engine, query: str, params: dict):
"""参数化查询防止注入"""
try:
with engine.connect() as conn:
stmt = text(query)
result = conn.execute(stmt, params)
return result.fetchall()
except SQLAlchemyError as e:
logger.error(f"SQL error: {e}")
raise
访问控制方案
python复制from functools import wraps
from flask import request, abort
def role_required(role):
"""基于角色的访问控制"""
def decorator(f):
@wraps(f)
def wrapper(*args, **kwargs):
user_role = get_current_user_role()
if user_role != role:
abort(403)
return f(*args, **kwargs)
return wrapper
return decorator
@app.route('/api/query')
@role_required('data_analyst')
def query_api():
"""需要data_analyst角色的API"""
pass
5. 典型问题排查指南
5.1 路由决策异常
症状:明显应该用SQL查询的问题却走了RAG路径
排查步骤:
- 检查工具描述是否足够差异化
- 验证LLM是否收到了完整的工具定义
- 分析问题是否包含歧义表达
修复方案:
python复制# 改进后的SQL工具描述
sql_tool = QueryEngineTool.from_defaults(
description=(
"专门用于查询结构化数据库中的精确数值和统计信息。"
"适用于以下场景:"
"- 需要具体数字的问题(如'销售额'、'用户数')"
"- 涉及比较、排序的问题(如'最高的'、'增长最快的')"
"- 明确提及数据库字段的问题"
"示例:'2023年Q4的营收增长率是多少?'"
)
)
5.2 混合查询结果混乱
症状:多个工具的结果组合后逻辑不通顺
解决方案:
python复制def integrate_results(question, *results):
"""智能结果整合"""
integration_prompt = f"""
请将以下查询结果整合成连贯的回答:
原始问题:{question}
查询结果:{results}
"""
return llm.generate(integration_prompt)
5.3 性能瓶颈分析
常见瓶颈点:
- RAG检索耗时过长
- SQL查询没有走索引
- LLM生成响应时间不稳定
优化工具:
python复制# 使用cProfile进行性能分析
import cProfile
profiler = cProfile.Profile()
profiler.enable()
# 执行待测代码
router.run("测试问题")
profiler.disable()
profiler.print_stats(sort='cumtime')
6. 扩展与演进方向
6.1 多模态支持
扩展系统处理图像和表格数据的能力:
python复制class MultiModalRouter:
def __init__(self):
self.tools = {
"text": TextQueryTool(),
"image": ImageAnalysisTool(),
"table": TableExtractionTool()
}
def detect_modality(self, query):
"""检测问题涉及的数据模态"""
return llm.classify(query, ["text", "image", "table"])
6.2 持续学习机制
实现系统的自我优化:
python复制class FeedbackLearner:
def __init__(self):
self.feedback_db = FeedbackDatabase()
def process_feedback(self, query, response, user_rating):
"""处理用户反馈"""
if user_rating < 3: # 负面反馈
self.analyze_failure(query, response)
self.retrain_if_needed()
def retrain_if_needed(self):
"""达到阈值后触发重新训练"""
if self.feedback_db.low_rating_count() > 100:
self.retrain_router()
6.3 边缘计算部署
适用于数据敏感场景的本地化方案:
python复制class EdgeDeployment:
def __init__(self, model_path):
self.model = load_compressed_model(model_path)
def query(self, question):
"""本地执行查询"""
return self.model.generate(question)
这套系统的真正价值在于它解决了AI落地过程中的"最后一公里"问题——不是简单地堆砌技术组件,而是让这些组件能够根据实际场景智能协作。随着业务需求的变化,系统可以通过新增工具类型(如API调用工具、计算工具等)来不断扩展能力边界,而无需推翻重来。
