diff --git a/.idea/.gitignore b/.idea/.gitignore deleted file mode 100644 index 13566b8..0000000 --- a/.idea/.gitignore +++ /dev/null @@ -1,8 +0,0 @@ -# Default ignored files -/shelf/ -/workspace.xml -# Editor-based HTTP Client requests -/httpRequests/ -# Datasource local storage ignored files -/dataSources/ -/dataSources.local.xml diff --git a/IMPACT_ANALYSIS.md b/IMPACT_ANALYSIS.md index d1866cb..74c01bc 100644 --- a/IMPACT_ANALYSIS.md +++ b/IMPACT_ANALYSIS.md @@ -551,3 +551,44 @@ | 项 | 说明 | |----|------| | `SSE_STREAM_CHUNK_CHARS` | 可选;每段 SSE 的 `content` 最大字符数,默认 `64`。 | + +--- + +# Impact Analysis Report — API 日志落盘与可观测性(追加) + +## 1. 改动概览 + +- **背景与目标**:终端日志偏少,难以排查 Text2SQL 全链路问题;需将更细粒度的日志写入仓库 `logs/`,且每次进程启动只保留一个 `.log` 文件。 +- **涉及模块**:`api_server.py`、`backend/utils/repo_logging.py`(新)、`backend/agents/orchestrator.py`、`backend/schema/indexer.py`、`backend/utils/dialog_classifier.py`。 +- **改动类型**:可观测性增强(行为对业务结果无影响)。 + +## 2. 方法级改动 + +| 位置 | 变更 | +|------|------| +| `configure_text2sql_api_logging` | 新建:清空 `logs/*.log`,创建 `logs/text2sql_api.log`(覆盖写);root 双 Handler(控制台简短格式 + 文件含 filename/lineno/funcName);关闭 Chroma/posthog 遥测 logger;`httpx`/`httpcore` 降为 WARNING。 | +| `api_server` 导入段 | 在 `from main import …` 之前调用上述配置;移除 `basicConfig`;`uvicorn.run(..., log_config=None)` 避免覆盖 root(勿用 `False`,否则会走 `fileConfig` 崩溃),并清空 uvicorn 自带 handler 改为 `propagate`。 | +| `_run_generate` / `nl_chat` / `_chat_stream_events` | 增加请求维度、对话上文长度、流式结束时的 valid/attempts/tables/SQL 摘要等 INFO 日志。 | +| `Text2SQLOrchestrator` | 向量粗筛输出表名+分数列表;选表输出 reasoning 预览;few-shot 注入 qid 与问题预览;生成 SQL 全文(超长截断至约 12k 字符);`_validate_sql` 返回前汇总 valid/errors/warnings/探针。 | +| `SchemaIndexer.search` | 检索结果由 DEBUG 改为 INFO,输出命中数与 top 表+分。 | +| `classify_dialog` 规则/hybrid 快路径 | DEBUG 改为 INFO,附用户输入预览;LLM 分支补充 reply 预览。 | + +## 3. 调用方与影响范围 + +- **调用方**:仅通过 `python api_server.py`(或等价导入 `api_server`)启动 API 时生效;`backend/main.py` CLI 仍使用自身 `basicConfig`,不落盘到本 `logs/text2sql_api.log`(未改 CLI)。 +- **破坏性变更**:否。 + +## 4. 风险与回滚 + +- **风险级别**:低。日志文件可能含用户查询片段与 SQL,需注意磁盘与隐私(内网演示场景可接受)。 +- **回滚**:删除 `repo_logging` 调用与相关增强日志,恢复 `basicConfig` 与默认 `uvicorn.run` 即可。 + +**回滚方式是否简单**:是。 + +## 5. 验证与测试 + +- 建议:`python -m py_compile api_server.py backend/utils/repo_logging.py`;启动 `python api_server.py` 后确认生成 `logs/text2sql_api.log` 且重启后仅保留该文件。 + +## 6. 配置变更 + +- 无新增环境变量;日志路径固定为仓库根下 `logs/text2sql_api.log`。 diff --git a/__pycache__/api_server.cpython-312.pyc b/__pycache__/api_server.cpython-312.pyc index 6600344..9ad4c41 100644 Binary files a/__pycache__/api_server.cpython-312.pyc and b/__pycache__/api_server.cpython-312.pyc differ diff --git a/api_server.py b/api_server.py index 5b6c08e..b174a07 100644 --- a/api_server.py +++ b/api_server.py @@ -28,6 +28,10 @@ _BACKEND_DIR = _REPO_DIR / "backend" if str(_BACKEND_DIR) not in sys.path: sys.path.insert(0, str(_BACKEND_DIR)) +from utils.repo_logging import configure_text2sql_api_logging + +configure_text2sql_api_logging(_REPO_DIR) + from main import setup_environment, load_schema, create_orchestrator, resolve_sql_dialect from agents.orchestrator import GenerationResult from nl_lite_store import lite_nl_store @@ -38,11 +42,6 @@ from utils.dialog_context import ( is_likely_follow_up, ) -logging.basicConfig( - level=logging.INFO, - format='%(asctime)s [%(levelname)s] %(name)s: %(message)s', - datefmt='%Y-%m-%d %H:%M:%S' -) logger = logging.getLogger(__name__) # per-request 覆盖 LLM client 时,使用锁避免并发串改 orchestrator.deepseek @@ -617,6 +616,15 @@ async def _run_generate( orch = get_orchestrator() dialect = resolve_sql_dialect(os.getenv("TEXT2SQL_DIALECT", "sqlserver")) dc = (dialog_context or "").strip() or None + q = question.strip() + logger.info( + "[GEN/API] dialect=%s top_k=%s dialog_context_chars=%s question_len=%s preview=%r", + dialect, + top_k, + len(dc) if dc else 0, + len(q), + q[:300] + ("…" if len(q) > 300 else ""), + ) async def _call_with_orch(o) -> GenerationResult: def _call() -> GenerationResult: @@ -643,6 +651,20 @@ async def _chat_stream_events(request: NLChatRequest) -> AsyncIterator[bytes]: text = request.message.strip() orch = get_orchestrator() dialog_block, last_data = await _load_session_text2sql_context(request, text) + logger.info( + "[API/stream] 开始: user_id=%r visitor_biz_id=%r session_id=%r service_code=%r model=%r " + "lang_code=%r msg_chars=%s preview=%r dialog_context_chars=%s last_turn_was_data_query=%s", + request.user_id, + request.visitor_biz_id, + request.session_id, + request.service_code, + request.model, + request.lang_code, + len(text), + text[:400] + ("…" if len(text) > 400 else ""), + len(dialog_block) if dialog_block else 0, + last_data, + ) lang = _normalize_lang_code(request.lang_code) async with _maybe_override_orch_llm(orch, request) as o: classified = await asyncio.to_thread( @@ -656,7 +678,11 @@ async def _chat_stream_events(request: NLChatRequest) -> AsyncIterator[bytes]: reply = (classified.reply_suggestion or "").strip() if not reply: reply = _localized_conversation_reply(lang) - logger.info("[API/stream] 对话意图: conversation(跳过 Text2SQL,与 CLI single_query 一致)") + logger.info( + "[API/stream] 对话意图 conversation(跳过 Text2SQL): reply_chars=%s reply_preview=%r", + len(reply), + reply[:500] + ("…" if len(reply) > 500 else ""), + ) yield _sse_data({"stage": "orchestrator", "stream_kind": "content", "content": "CHAT"}) async for pkt in _sse_stream_text_chunks("chat", reply): yield pkt @@ -669,10 +695,24 @@ async def _chat_stream_events(request: NLChatRequest) -> AsyncIterator[bytes]: try: result = await _run_generate(text, dialog_context=dialog_block or None, request=request) except Exception as e: - logger.error(f"[API/stream] 生成异常: {e}") + logger.error(f"[API/stream] 生成异常: {e}", exc_info=True) yield _sse_data({"code": 500, "msg": str(e), "data": None}) return + sql_out = (result.sql or "").strip() + logger.info( + "[API/stream] Text2SQL 完成: valid=%s attempts=%s tables_used=%s sql_chars=%s sql_head=%r", + result.valid, + result.attempts, + result.tables_used, + len(sql_out), + sql_out[:500] + ("…" if len(sql_out) > 500 else ""), + ) + if result.errors: + logger.warning("[API/stream] 错误列表: %s", result.errors) + if result.warnings: + logger.info("[API/stream] 警告列表: %s", result.warnings) + # What this SQL does / SQL 说明:按 lang_code 本地化 try: if isinstance(result.metadata, dict): @@ -747,9 +787,21 @@ async def nl_chat(request: NLChatRequest): lang = _normalize_lang_code(request.lang_code) raise HTTPException(status_code=400, detail=_localized_empty_input_reply(lang)) text = request.message.strip() - logger.info(f"[API] 问题: {text[:100]}...") orch = get_orchestrator() dialog_block, last_data = await _load_session_text2sql_context(request, text) + logger.info( + "[API/chat] 请求: user_id=%r session_id=%r service_code=%r model=%r lang=%r " + "msg_chars=%s preview=%r dialog_context_chars=%s last_turn_was_data_query=%s", + request.user_id, + request.session_id, + request.service_code, + request.model, + request.lang_code, + len(text), + text[:400] + ("…" if len(text) > 400 else ""), + len(dialog_block) if dialog_block else 0, + last_data, + ) lang = _normalize_lang_code(request.lang_code) async with _maybe_override_orch_llm(orch, request) as o: classified = await asyncio.to_thread( @@ -1004,5 +1056,7 @@ if __name__ == "__main__": host=host, port=port, reload=False, - log_level="info" + log_level="info", + # 必须为 None:False 仍会进入 uvicorn 的 fileConfig 分支并崩溃;None 才跳过覆盖 root 日志 + log_config=None, ) diff --git a/backend/__pycache__/main.cpython-312.pyc b/backend/__pycache__/main.cpython-312.pyc index d6f118d..f4d4c9a 100644 Binary files a/backend/__pycache__/main.cpython-312.pyc and b/backend/__pycache__/main.cpython-312.pyc differ diff --git a/backend/agents/__pycache__/orchestrator.cpython-312.pyc b/backend/agents/__pycache__/orchestrator.cpython-312.pyc index 276de29..40bd50d 100644 Binary files a/backend/agents/__pycache__/orchestrator.cpython-312.pyc and b/backend/agents/__pycache__/orchestrator.cpython-312.pyc differ diff --git a/backend/agents/orchestrator.py b/backend/agents/orchestrator.py index 692ab7a..935dd0e 100644 --- a/backend/agents/orchestrator.py +++ b/backend/agents/orchestrator.py @@ -182,7 +182,11 @@ class Text2SQLOrchestrator: indexer.ensure_index_for_schema(self.schema_manager) # 检索 - logger.info(f"[Orchestrator] 开始向量检索: query='{question[:50]}...'") + logger.info( + "[Orchestrator] 开始向量检索: query_chars=%s query_preview=%r", + len(question or ""), + (question or "")[:200] + ("…" if len(question or "") > 200 else ""), + ) results = indexer.search( query=question, top_k=top_k, @@ -190,7 +194,15 @@ class Text2SQLOrchestrator: ) candidate_tables = [r["table_name"] for r in results] - logger.debug(f"粗筛候选表:{candidate_tables[:10]}...(共{len(candidate_tables)}个)") + scored = [ + (r["table_name"], round(float(r.get("score", 0.0)), 4)) + for r in results[: min(25, len(results))] + ] + logger.info( + "[Orchestrator] 向量粗筛: 命中=%s 张(阈值内),表名+分: %s", + len(candidate_tables), + scored, + ) return candidate_tables except Exception as e: @@ -239,7 +251,13 @@ class Text2SQLOrchestrator: # 限制数量 relevant_tables = relevant_tables[:max_tables] - logger.info(f"LLM精筛选中表:{relevant_tables}") + rs = (reasoning or "").strip() + logger.info( + "LLM精筛选中表:%s | reasoning_chars=%s reasoning_preview=%r", + relevant_tables, + len(rs), + rs[:600] + ("…" if len(rs) > 600 else ""), + ) return relevant_tables, reasoning _BROKER_KEYWORDS_CN = ("对手方", "经纪商", "券商", "對手方") @@ -361,7 +379,12 @@ class Text2SQLOrchestrator: for i, ex in enumerate(examples) ]) schema_str = f"参考以下相似示例的SQL编写风格:\n\n{examples_prompt}\n\n【当前Schema】\n{schema_str}" - logger.debug(f"已注入 {len(examples)} 个few-shot示例: {[ex.qid for ex in examples]}") + logger.info( + "已注入 %s 个 few-shot 示例: qid=%s question_zh_preview=%r", + len(examples), + [ex.qid for ex in examples], + [((ex.question_zh or "")[:80] + "…") if len(ex.question_zh or "") > 80 else (ex.question_zh or "") for ex in examples], + ) except Exception as e: logger.warning(f"Few-shot检索失败: {e}") @@ -419,7 +442,9 @@ class Text2SQLOrchestrator: sql = normalize_sql_for_dialect(sql, dialect) - logger.debug(f"生成的SQL:{sql[:200]}...") + lim = 12000 + body = sql if len(sql) <= lim else sql[:lim] + "\n…(日志已截断)" + logger.info("生成的SQL(chars=%s):\n%s", len(sql), body) return sql def _validate_sql( @@ -562,6 +587,19 @@ class Text2SQLOrchestrator: empty_feedback = (empty_feedback or prefix) + follow is_valid = len(errors) == 0 + logger.info( + "[validate] 程序+探针+LLM 汇总: valid=%s err_count=%s warn_count=%s " + "db_execution_status=%s sql_chars=%s", + is_valid, + len(errors), + len(warnings), + db_execution_status, + len(sql or ""), + ) + if errors: + logger.info("[validate] errors 预览: %s", errors[:5]) + if warnings: + logger.info("[validate] warnings: %s", warnings[:5]) return is_valid, errors, warnings, db_execution_status, empty_feedback def generate( @@ -608,7 +646,12 @@ class Text2SQLOrchestrator: dc_raw = (dialog_context or "").strip() retrieval_question = self._merge_dialog_for_model(dc_raw, question, 4000) linker_question = self._merge_dialog_for_model(dc_raw, question, 6000) - logger.info(f"[GEN] 开始生成SQL:{question[:50]}...") + logger.info( + "[GEN] 开始生成SQL: question_chars=%s preview=%r dialog_context_chars=%s", + len(question or ""), + (question or "")[:300] + ("…" if len(question or "") > 300 else ""), + len(dc_raw) if dc_raw else 0, + ) attempt = 0 last_sql = None diff --git a/backend/llm/__pycache__/deepseek_client.cpython-312.pyc b/backend/llm/__pycache__/deepseek_client.cpython-312.pyc index 8abd54e..ed78f27 100644 Binary files a/backend/llm/__pycache__/deepseek_client.cpython-312.pyc and b/backend/llm/__pycache__/deepseek_client.cpython-312.pyc differ diff --git a/backend/schema/__pycache__/indexer.cpython-312.pyc b/backend/schema/__pycache__/indexer.cpython-312.pyc index 90af973..186a408 100644 Binary files a/backend/schema/__pycache__/indexer.cpython-312.pyc and b/backend/schema/__pycache__/indexer.cpython-312.pyc differ diff --git a/backend/schema/indexer.py b/backend/schema/indexer.py index 3613ef4..e880615 100644 --- a/backend/schema/indexer.py +++ b/backend/schema/indexer.py @@ -22,12 +22,8 @@ class SchemaIndexer: 功能: 1. 为表结构构建向量索引 2. 基于问题的表检索 -<<<<<<< HEAD 3. 默认使用 Chroma ``PersistentClient``,数据落在 ``persist_dir``(与 ``VECTOR_DB_PATH`` 一致); 仅当环境变量 ``SCHEMA_INDEXER_EPHEMERAL=true`` 或构造参数 ``use_ephemeral=True`` 时使用内存客户端。 -======= - 3. Chroma 使用持久化(PersistentClient),向量与元数据写入 ``persist_dir``,进程重启后可复用 ->>>>>>> 9369563942d8de442508613191670123d5ae9d08 """ def __init__( @@ -43,11 +39,7 @@ class SchemaIndexer: Args: embedder: Embedding模型实例 -<<<<<<< HEAD - persist_dir: Chroma 持久化目录(默认与编排器 ``vector_db_path`` / ``VECTOR_DB_PATH`` 一致) -======= persist_dir: Chroma 持久化根目录(磁盘路径) ->>>>>>> 9369563942d8de442508613191670123d5ae9d08 collection_name: 集合名称 use_ephemeral: 为 True 时使用内存 Chroma;为 None 时读环境变量 SCHEMA_INDEXER_EPHEMERAL """ @@ -56,17 +48,15 @@ class SchemaIndexer: self.persist_dir.mkdir(parents=True, exist_ok=True) self.collection_name = collection_name -<<<<<<< HEAD env_ephemeral = os.getenv("SCHEMA_INDEXER_EPHEMERAL", "").lower() in ( "1", "true", "yes", -======= + ) # 持久化 Chroma:数据落盘至 persist_dir self.client = chromadb.PersistentClient( path=str(self.persist_dir), settings=ChromaSettings(anonymized_telemetry=False), ->>>>>>> 9369563942d8de442508613191670123d5ae9d08 ) self._chroma_ephemeral = bool(use_ephemeral) if use_ephemeral is not None else env_ephemeral @@ -90,13 +80,8 @@ class SchemaIndexer: ) logger.info( -<<<<<<< HEAD f"[OK] 初始化SchemaIndexer({backend_desc}): collection={self.collection_name}, " f"count={self.collection.count()}" -======= - f"[OK] 初始化SchemaIndexer(Chroma持久化): collection={self.collection_name}, " - f"path={self.persist_dir}" ->>>>>>> 9369563942d8de442508613191670123d5ae9d08 ) def ensure_index_for_schema( @@ -255,7 +240,14 @@ class SchemaIndexer: "rank": idx + 1, }) - logger.debug(f"检索 '{query[:50]}...' -> 找到{len(formatted)}个相关表(阈值={score_threshold})") + top = [(x["table_name"], round(float(x.get("score", 0.0)), 4)) for x in formatted[:15]] + logger.info( + "Schema 向量检索: query_chars=%s 命中=%s(阈值=%s)top=%s", + len(query or ""), + len(formatted), + score_threshold, + top, + ) return formatted def search_by_table_names(self, table_names: List[str]) -> List[Dict]: @@ -332,9 +324,5 @@ class SchemaIndexer: "total_columns": total_columns, "avg_columns": total_columns / count if count > 0 else 0, "persist_dir": str(self.persist_dir), -<<<<<<< HEAD "chroma_mode": "memory" if self._chroma_ephemeral else "persistent", -======= - "chroma_mode": "persistent", ->>>>>>> 9369563942d8de442508613191670123d5ae9d08 } diff --git a/backend/utils/__pycache__/dialog_classifier.cpython-312.pyc b/backend/utils/__pycache__/dialog_classifier.cpython-312.pyc index 3df1ffe..7153eed 100644 Binary files a/backend/utils/__pycache__/dialog_classifier.cpython-312.pyc and b/backend/utils/__pycache__/dialog_classifier.cpython-312.pyc differ diff --git a/backend/utils/__pycache__/fewshot_selector.cpython-312.pyc b/backend/utils/__pycache__/fewshot_selector.cpython-312.pyc index a710b5e..99f9457 100644 Binary files a/backend/utils/__pycache__/fewshot_selector.cpython-312.pyc and b/backend/utils/__pycache__/fewshot_selector.cpython-312.pyc differ diff --git a/backend/utils/__pycache__/repo_logging.cpython-312.pyc b/backend/utils/__pycache__/repo_logging.cpython-312.pyc new file mode 100644 index 0000000..9adc801 Binary files /dev/null and b/backend/utils/__pycache__/repo_logging.cpython-312.pyc differ diff --git a/backend/utils/dialog_classifier.py b/backend/utils/dialog_classifier.py index e57bb4e..efa5b51 100644 --- a/backend/utils/dialog_classifier.py +++ b/backend/utils/dialog_classifier.py @@ -136,28 +136,28 @@ def _classify_dialog_rules( ) if _SQL_OR_QUERY_HINT_RE.search(t): - logger.debug("[dialog] intent=text2sql (query/business hint)") + logger.info("[dialog] intent=text2sql (query/business hint) preview=%r", user_text[:120]) return DialogClassifyResult(DialogIntent.TEXT2SQL, None) core = _strip_trailing_punct(t) if core.casefold() in _CHITCHAT_KEYS: - logger.debug("[dialog] intent=conversation (chitchat phrase)") + logger.info("[dialog] intent=conversation (chitchat phrase) core=%r", core) return DialogClassifyResult( DialogIntent.CONVERSATION, reply_suggestion=DEFAULT_CONVERSATION_REPLY ) if _META_QUESTION_RE.search(t): - logger.debug("[dialog] intent=conversation (meta question)") + logger.info("[dialog] intent=conversation (meta question) preview=%r", user_text[:120]) return DialogClassifyResult( DialogIntent.CONVERSATION, reply_suggestion=DEFAULT_CONVERSATION_REPLY, ) if last_turn_was_data_query: - logger.debug("[dialog] intent=text2sql (follow-up after data query)") + logger.info("[dialog] intent=text2sql (follow-up after data query) preview=%r", user_text[:120]) return DialogClassifyResult(DialogIntent.TEXT2SQL, None) - logger.debug("[dialog] intent=text2sql (default)") + logger.info("[dialog] intent=text2sql (default rules) preview=%r", user_text[:120]) return DialogClassifyResult(DialogIntent.TEXT2SQL, None) @@ -208,11 +208,15 @@ def _classify_dialog_llm( "若上一版 SQL 或结果不符合预期,请具体说明:希望增加/修改哪些条件、" "时间范围或统计维度,以便重新生成。" ) - logger.info("[dialog] intent=conversation (LLM)") + logger.info( + "[dialog] intent=conversation (LLM) reply_chars=%s reply_preview=%r", + len(reply), + reply[:300] + ("…" if len(reply) > 300 else ""), + ) return DialogClassifyResult(DialogIntent.CONVERSATION, reply_suggestion=reply) if intent_s == "text2sql": - logger.info("[dialog] intent=text2sql (LLM)") + logger.info("[dialog] intent=text2sql (LLM) user_preview=%r", user_text[:200]) return DialogClassifyResult(DialogIntent.TEXT2SQL, None) raise ValueError(f"unexpected intent field: {intent_s!r}") @@ -254,18 +258,18 @@ def classify_dialog( ) if _SQL_OR_QUERY_HINT_RE.search(t): - logger.debug("[dialog] intent=text2sql (query/business hint, hybrid fast)") + logger.info("[dialog] intent=text2sql (hybrid fast: query hint) preview=%r", user_text[:120]) return DialogClassifyResult(DialogIntent.TEXT2SQL, None) core = _strip_trailing_punct(t) if core.casefold() in _CHITCHAT_KEYS: - logger.debug("[dialog] intent=conversation (chitchat phrase, hybrid fast)") + logger.info("[dialog] intent=conversation (hybrid fast: chitchat) core=%r", core) return DialogClassifyResult( DialogIntent.CONVERSATION, reply_suggestion=DEFAULT_CONVERSATION_REPLY ) if _META_QUESTION_RE.search(t): - logger.debug("[dialog] intent=conversation (meta question, hybrid fast)") + logger.info("[dialog] intent=conversation (hybrid fast: meta) preview=%r", user_text[:120]) return DialogClassifyResult( DialogIntent.CONVERSATION, reply_suggestion=DEFAULT_CONVERSATION_REPLY, diff --git a/backend/utils/repo_logging.py b/backend/utils/repo_logging.py new file mode 100644 index 0000000..059f88b --- /dev/null +++ b/backend/utils/repo_logging.py @@ -0,0 +1,80 @@ +""" +Text2SQL API 仓库级日志:控制台 + 单一磁盘文件(每次进程启动清空 logs/*.log 后重写)。 + +由 ``api_server`` 在导入早期调用;与 ``uvicorn.run(..., log_config=False)`` 配合, +使 uvicorn / FastAPI 的日志经 root 统一落到文件。 +""" + +from __future__ import annotations + +import logging +from pathlib import Path + + +def configure_text2sql_api_logging(repo_root: Path) -> Path: + """ + 配置 root logger:控制台(INFO,简短格式)+ ``logs/text2sql_api.log``(INFO,含文件名/行号/函数)。 + + 每次调用会删除 ``logs`` 目录下所有 ``*.log``,再新建 ``text2sql_api.log``(覆盖写), + 保证一次运行仅保留一个日志文件。 + + Returns: + 主日志文件绝对路径。 + """ + log_dir = (repo_root / "logs").resolve() + log_dir.mkdir(parents=True, exist_ok=True) + for p in log_dir.glob("*.log"): + try: + p.unlink() + except OSError: + pass + + log_path = log_dir / "text2sql_api.log" + + root = logging.getLogger() + root.setLevel(logging.INFO) + for h in list(root.handlers): + root.removeHandler(h) + try: + h.close() + except Exception: + pass + + fmt_console = logging.Formatter( + "%(asctime)s [%(levelname)s] %(name)s: %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", + ) + fmt_file = logging.Formatter( + "%(asctime)s %(levelname)-7s [%(name)s] %(filename)s:%(lineno)d %(funcName)s() | %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", + ) + + fh = logging.FileHandler(log_path, mode="w", encoding="utf-8") + fh.setLevel(logging.INFO) + fh.setFormatter(fmt_file) + + ch = logging.StreamHandler() + ch.setLevel(logging.INFO) + ch.setFormatter(fmt_console) + + root.addHandler(fh) + root.addHandler(ch) + + # 依赖库降噪(避免 httpx 每条请求刷屏;Chroma 遥测与新版 posthog 不兼容时会打 ERROR) + logging.getLogger("httpx").setLevel(logging.WARNING) + logging.getLogger("httpcore").setLevel(logging.WARNING) + for _name in ( + "chromadb.telemetry", + "chromadb.telemetry.product.posthog", + "posthog", + ): + logging.getLogger(_name).disabled = True + + # Uvicorn 默认自带 handler;与 log_config=False 合用时改为只往 root 冒泡 + for name in ("uvicorn", "uvicorn.error", "uvicorn.access"): + lg = logging.getLogger(name) + lg.handlers.clear() + lg.propagate = True + + logging.getLogger(__name__).info("日志文件: %s", log_path) + return log_path diff --git a/data/embeddings/chroma_fewshot/9756e6ea-6d51-4e24-a490-4a58f0a585b9/length.bin b/data/embeddings/chroma_fewshot/9756e6ea-6d51-4e24-a490-4a58f0a585b9/length.bin index 2d95673..d33cd2c 100644 Binary files a/data/embeddings/chroma_fewshot/9756e6ea-6d51-4e24-a490-4a58f0a585b9/length.bin and b/data/embeddings/chroma_fewshot/9756e6ea-6d51-4e24-a490-4a58f0a585b9/length.bin differ diff --git a/logs/.gitkeep b/logs/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/logs/text2sql_api.log b/logs/text2sql_api.log new file mode 100644 index 0000000..d59031f --- /dev/null +++ b/logs/text2sql_api.log @@ -0,0 +1,67 @@ +2026-04-16 10:32:14 INFO [utils.repo_logging] repo_logging.py:79 configure_text2sql_api_logging() | 日志文件: C:\Users\24019\Desktop\backman-camel\logs\text2sql_api.log +2026-04-16 10:32:17 INFO [__main__] api_server.py:86 () | [OK] 已加载配置文件: C:\Users\24019\Desktop\backman-camel\.env +2026-04-16 10:32:17 INFO [__main__] api_server.py:1047 () | 启动服务: http://0.0.0.0:8041 +2026-04-16 10:32:17 INFO [__main__] api_server.py:1048 () | API文档: http://0.0.0.0:8041/docs +2026-04-16 10:32:17 INFO [uvicorn.error] server.py:92 _serve() | Started server process [6728] +2026-04-16 10:32:17 INFO [uvicorn.error] on.py:48 startup() | Waiting for application startup. +2026-04-16 10:32:17 INFO [__main__] api_server.py:741 lifespan() | ============================================================ +2026-04-16 10:32:17 INFO [__main__] api_server.py:742 lifespan() | Text2SQL API Server 启动中... +2026-04-16 10:32:17 INFO [__main__] api_server.py:743 lifespan() | ============================================================ +2026-04-16 10:32:17 INFO [main] main.py:126 setup_environment() | [OK] 环境检查通过 +2026-04-16 10:32:17 INFO [main] main.py:127 setup_environment() | - Schema: data\schemas\G3SB_MCDataDictionary_table_structure.json +2026-04-16 10:32:17 INFO [main] main.py:128 setup_environment() | - LLM: openai +2026-04-16 10:32:17 INFO [main] main.py:154 load_schema() | 加载Schema: ./data/schemas/G3SB_MCDataDictionary_table_structure.json +2026-04-16 10:32:17 INFO [main] main.py:156 load_schema() | 表注释(meta): ./data/schemas/G3SB_MCDataDictionary_table_meta.json +2026-04-16 10:32:18 INFO [schema.loader] loader.py:107 load_from_json() | [OK] 加载Schema完成(G3SB schemas): G3SB_MCDataDictionary_table_structure, 共2516张表 +2026-04-16 10:32:18 INFO [main] main.py:163 load_schema() | [OK] Schema加载完成: G3SB_MCDataDictionary_table_structure, 共2516张表, 59196个字段 +2026-04-16 10:32:21 INFO [llm.openai_client] openai_client.py:55 __init__() | [OK] OpenAIClient初始化: model=gpt-5.4, base_url=http://113.192.49.54:9080/v1 +2026-04-16 10:32:23 INFO [utils.embedding] embedding.py:144 __init__() | 使用 OpenAI Embedding API:model=text-embedding-ada-002,base_url=http://113.192.49.54:9080/v1,max_batch=100 +2026-04-16 10:32:25 INFO [utils.fewshot_chroma_store] fewshot_chroma_store.py:86 __init__() | [OK] FewShotChromaStore: 磁盘 data\embeddings\chroma_fewshot collection=fewshot_samples count=50 +2026-04-16 10:32:25 INFO [utils.fewshot_selector] fewshot_selector.py:128 _init_chroma_mode() | [Few-shot] 已从向量库加载(Chroma 50 条,data\embeddings\chroma_fewshot) +2026-04-16 10:32:25 INFO [utils.fewshot_selector] fewshot_selector.py:159 _init_chroma_mode() | [OK] Few-shot 使用 Chroma(50 条) +2026-04-16 10:32:25 INFO [agents.orchestrator] orchestrator.py:116 __init__() | Few-shot已启用: top_k=3, min_rating=7 +2026-04-16 10:32:25 INFO [agents.orchestrator] orchestrator.py:124 __init__() | [OK] Text2SQLOrchestrator初始化完成: max_retry=2, use_vector_search=True, fewshot=on, nl→zh_norm=on +2026-04-16 10:32:25 INFO [__main__] api_server.py:133 get_orchestrator() | [OK] Orchestrator 初始化完成 +2026-04-16 10:32:25 INFO [__main__] api_server.py:747 lifespan() | [OK] 服务已就绪 +2026-04-16 10:32:25 INFO [uvicorn.error] on.py:62 startup() | Application startup complete. +2026-04-16 10:32:25 INFO [uvicorn.error] server.py:224 _log_started_message() | Uvicorn running on http://0.0.0.0:8041 (Press CTRL+C to quit) +2026-04-16 10:32:42 INFO [uvicorn.access] httptools_impl.py:483 send() | 127.0.0.1:56978 - "POST /g3sb/api/nl/chat/stream HTTP/1.1" 200 +2026-04-16 10:32:42 INFO [__main__] api_server.py:654 _chat_stream_events() | [API/stream] 开始: user_id='anonymous' visitor_biz_id=None session_id=None service_code=None model='gpt-4o-mini' lang_code='auto' msg_chars=15 preview='所有客户账户之间的股票转移记录' dialog_context_chars=0 last_turn_was_data_query=False +2026-04-16 10:32:45 INFO [llm.openai_client] openai_client.py:55 __init__() | [OK] OpenAIClient初始化: model=gpt-4o-mini, base_url=http://113.192.49.54:9080/v1 +2026-04-16 10:32:45 INFO [utils.dialog_classifier] dialog_classifier.py:261 classify_dialog() | [dialog] intent=text2sql (hybrid fast: query hint) preview='所有客户账户之间的股票转移记录' +2026-04-16 10:32:45 INFO [__main__] api_server.py:620 _run_generate() | [GEN/API] dialect=tsql top_k=20 dialog_context_chars=0 question_len=15 preview='所有客户账户之间的股票转移记录' +2026-04-16 10:32:48 INFO [llm.openai_client] openai_client.py:55 __init__() | [OK] OpenAIClient初始化: model=gpt-4o-mini, base_url=http://113.192.49.54:9080/v1 +2026-04-16 10:32:51 INFO [agents.orchestrator] orchestrator.py:636 generate() | [GEN] 问句已归一中文:所有客户账户的股票转移记录 +2026-04-16 10:32:51 INFO [agents.orchestrator] orchestrator.py:649 generate() | [GEN] 开始生成SQL: question_chars=13 preview='所有客户账户的股票转移记录' dialog_context_chars=0 +2026-04-16 10:32:51 INFO [agents.orchestrator] orchestrator.py:664 generate() | 尝试 #1 +2026-04-16 10:32:51 INFO [schema.indexer] indexer.py:82 __init__() | [OK] 初始化SchemaIndexer(Chroma磁盘 path=data\embeddings\chroma): collection=schema_tables, count=2516 +2026-04-16 10:32:51 INFO [schema.indexer] indexer.py:106 ensure_index_for_schema() | Schema 向量索引已就绪(2516 张表),跳过向量化 +2026-04-16 10:32:51 INFO [agents.orchestrator] orchestrator.py:185 _coarse_filter() | [Orchestrator] 开始向量检索: query_chars=13 query_preview='所有客户账户的股票转移记录' +2026-04-16 10:32:53 INFO [utils.embedding] embedding.py:163 _set_dim_from_vector() | [OK] Embedding 向量维度:1536 +2026-04-16 10:32:53 INFO [schema.indexer] indexer.py:244 search() | Schema 向量检索: query_chars=13 命中=20(阈值=0.1)top=[('VSBHKRpt0430', 0.8264), ('VSBHKRpt0431', 0.8188), ('TSBTransferInstruction', 0.8147), ('VSBTransferInstruction', 0.8134), ('VSBHKRpt0397', 0.8126), ('VSBHKRpt0672', 0.8117), ('VSBHKRpt0090C', 0.81), ('VSBHKRpt0570', 0.8084), ('VSBHKRpt1030A', 0.8077), ('VSBRpt0999E', 0.8065), ('VSBRpt0999C', 0.8064), ('TSBAccountEntitlementRelease', 0.8063), ('VSBRpt0999A', 0.8059), ('TSBAccountInstrumentMovement', 0.8038), ('VSBTransferInstructionGenerationByAccountContract', 0.8033)] +2026-04-16 10:32:53 INFO [agents.orchestrator] orchestrator.py:201 _coarse_filter() | [Orchestrator] 向量粗筛: 命中=20 张(阈值内),表名+分: [('VSBHKRpt0430', 0.8264), ('VSBHKRpt0431', 0.8188), ('TSBTransferInstruction', 0.8147), ('VSBTransferInstruction', 0.8134), ('VSBHKRpt0397', 0.8126), ('VSBHKRpt0672', 0.8117), ('VSBHKRpt0090C', 0.81), ('VSBHKRpt0570', 0.8084), ('VSBHKRpt1030A', 0.8077), ('VSBRpt0999E', 0.8065), ('VSBRpt0999C', 0.8064), ('TSBAccountEntitlementRelease', 0.8063), ('VSBRpt0999A', 0.8059), ('TSBAccountInstrumentMovement', 0.8038), ('VSBTransferInstructionGenerationByAccountContract', 0.8033), ('VSBRpt0060', 0.8021), ('XCGatewayStockReconciliationReport', 0.802), ('TSBAccountEntitlementHold', 0.8013), ('WSBBatchLocationTransferDetail', 0.8012), ('VCAccountCashMovement', 0.8012)] +2026-04-16 10:32:57 INFO [agents.orchestrator] orchestrator.py:255 _llm_select_tables() | LLM精筛选中表:['VSBHKRpt0430', 'VSBHKRpt0431', 'TSBTransferInstruction', 'VSBTransferInstruction', 'TSBAccountInstrumentMovement'] | reasoning_chars=175 reasoning_preview='问题涉及客户账户的股票转移记录,VSBHKRpt0430和VSBHKRpt0431提供了客户股票转移的日报信息,TSBTransferInstruction和VSBTransferInstruction记录账户转移指令,TSBAccountInstrumentMovement则详细记录账户的股票移动交易。这些表共同涵盖了客户账户的股票转移相关信息。' +2026-04-16 10:32:57 INFO [agents.orchestrator] orchestrator.py:692 generate() | 选中表:['VSBHKRpt0430', 'VSBHKRpt0431', 'TSBTransferInstruction', 'VSBTransferInstruction', 'TSBAccountInstrumentMovement'],扩展后:['VSBHKRpt0430', 'VSBHKRpt0431', 'TSBTransferInstruction', 'TSBAccountInstrumentMovement', 'VSBTransferInstruction'] +2026-04-16 10:32:58 INFO [utils.fewshot_selector] fewshot_selector.py:271 _select_chroma() | Few-shot(Chroma)选择: 问题='所有客户账户的股票转移记录...' → 选中3个示例 (top_k=3, min_rating=7) +2026-04-16 10:32:58 INFO [agents.orchestrator] orchestrator.py:382 _generate_sql() | 已注入 3 个 few-shot 示例: qid=['Q2', 'Q40', 'Q22'] question_zh_preview=['列出今日所有客户账户之间的股票转移记录。', '列出今日执行的所有股票转移/移动。', '列出今日所有公司账户与客户账户之间的现金划转。'] +2026-04-16 10:33:07 INFO [agents.orchestrator] orchestrator.py:447 _generate_sql() | 生成的SQL(chars=476): +SELECT + m.MovementID, + m.AccountID AS FromAccountID, + m.TransferToAccountID, + m.InstrumentID, + i.Name AS InstrumentSymbol, + m.MovementType, + m.Quantity AS TransferQuantity, + m.ValueDate AS TransferDate +FROM TSBAccountInstrumentMovement m +LEFT JOIN MCInstrument i ON m.InstrumentID = i.InstrumentID +WHERE m.MovementType = 'T' + AND m.ValueDate >= CAST(GETDATE() AS DATE) + AND m.ValueDate < DATEADD(DAY, 1, CAST(GETDATE() AS DATE)) +ORDER BY m.ValueDate; +2026-04-16 10:33:07 INFO [db.engine] engine.py:45 get_engine() | SQLAlchemy engine initialized from database_url +2026-04-16 10:33:07 INFO [db.dbhub_tools] dbhub_tools.py:627 _execute_sql() | _execute_sql 执行语句数=1 readonly=True max_rows=1 +2026-04-16 10:33:10 INFO [agents.orchestrator] orchestrator.py:590 _validate_sql() | [validate] 程序+探针+LLM 汇总: valid=True err_count=0 warn_count=0 db_execution_status=0 sql_chars=476 +2026-04-16 10:33:10 INFO [agents.orchestrator] orchestrator.py:753 generate() | [OK] SQL生成与验证通过(1次尝试) +2026-04-16 10:33:10 INFO [__main__] api_server.py:703 _chat_stream_events() | [API/stream] Text2SQL 完成: valid=True attempts=1 tables_used=['VSBHKRpt0430', 'VSBHKRpt0431', 'TSBTransferInstruction', 'TSBAccountInstrumentMovement', 'VSBTransferInstruction'] sql_chars=476 sql_head="SELECT\n m.MovementID,\n m.AccountID AS FromAccountID,\n m.TransferToAccountID,\n m.InstrumentID,\n i.Name AS InstrumentSymbol,\n m.MovementType,\n m.Quantity AS TransferQuantity,\n m.ValueDate AS TransferDate\nFROM TSBAccountInstrumentMovement m\nLEFT JOIN MCInstrument i ON m.InstrumentID = i.InstrumentID\nWHERE m.MovementType = 'T'\n AND m.ValueDate >= CAST(GETDATE() AS DATE)\n AND m.ValueDate < DATEADD(DAY, 1, CAST(GETDATE() AS DATE))\nORDER BY m.ValueDate;" diff --git a/tech_architecture.md b/tech_architecture.md index d6a2126..db6e0b4 100644 --- a/tech_architecture.md +++ b/tech_architecture.md @@ -22,79 +22,123 @@ Text2SQL 系统由以下三个核心智能体组成: flowchart TD subgraph 输入层 A[自然语言问题] + Actx[可选会话上文 +dialog_context] end - + subgraph 核心处理层 - A --> B[意图分类 -Dialog Classifier] - B -->|非查询意图| C[返回对话回复] - - subgraph "智能体 1: Schema Linker" - D[表选择] - D1[向量检索粗筛 -SchemaIndexer] - D2[LLM 精筛 -DeepSeek] - D3[外键扩展 -关联表处理] + A --> B[意图分类 classify_dialog] + Actx --> B + B -->|conversation| C[直接返回引导性中文回复] + B -->|text2sql| N0[问句归一 normalize_nl_question +TRANSLATE_EN_TO_ZH=true 默认开启] + N0 --> D0[合并上文用于检索选表 +_merge_dialog_for_model] + + subgraph "Schema Linker 表选择模块" + D[表选择流程] + D1[向量粗筛 _coarse_filter +SchemaIndexer Chroma 或全表降级] + D2[LLM 精筛 _llm_select_tables +select_tables] + D2b[经纪商维度表优先 +_prioritize_broker_tables] + D3[外键双向扩展 +_expand_relations] + D4[生成紧凑 Schema 字符串 +to_compact_string] D --> D1 D1 --> D2 - D2 --> D3 + D2 --> D2b + D2b --> D3 + D3 --> D4 end - - B -->|查询意图| D - - subgraph "智能体 2: SQL Generator" - E[SQL 生成] - E1[Few-shot 示例增强] - E2[LLM 生成 SQL -DeepSeek] + + D0 --> D + + subgraph "SQL Generator SQL生成模块" + E[SQL 生成 _generate_sql] + E1[Few-shot 示例增强 FewShotSelector +Chroma 或 JSONL FEWSHOT_USE_CHROMA] + E2[LLM chat 生成 +temperature=0.0 top_p=1.0] + E3[sqlglot 方言归一 +normalize_sql_for_dialect] E --> E1 E1 --> E2 + E2 --> E3 end - - D3 --> E - - subgraph "智能体 3: Validator" - F[SQL 验证] - F1[程序验证 -语法检查] - F1b[数据库试执行 -行列探针 0/1/-1] - F2[LLM 语义验证 -未探针时] - F3[结果评估] + + D4 --> E + + subgraph "Validator SQL验证模块 _validate_sql" + F[程序校验] + F1[语法验证 sqlglot] + F1a[表列一致性验证 SchemaManager] + F1b[危险操作检查 check_dangerous_operations] + F1c[T-SQL 字面量规则检查 +check_no_cjk_in_sql_string_literals] + F1d[库探针 probe_sql_execution_status_ex +需 database_url 配置] + F2[LLM 语义验证 validate_sql +仅探针为 None 时调用] + F0a[探针 1 交付说明 +sql_probe_success_delivery_message] + F0b[探针 0 无数据说明 +empty_result_user_feedback] F --> F1 - F1 --> F1b - F1b --> F2 - F2 --> F3 + F1 --> F1a + F1a --> F1b + F1b --> F1c + F1c -->|无程序错误| F1d + F1c -->|有程序错误| I[未通过 max_retry 内重试] + F1d -->|探针 None| F2 + F1d -->|探针 1| F0a + F1d -->|探针 0| F0b + F1d -->|探针 -1| I + F2 -->|LLM 报错| I + F2 -->|通过| H[GenerationResult 有效 SQL] + F0a --> H + F0b --> H end - - E2 --> F - F3 -->|验证通过| H[返回有效 SQL] - F3 -->|验证失败| I[重试逻辑 -最多 max_retry 次] - I -->|重试| E + + E3 --> F + + I --> R[重试回合] + R --> R1[extract_tables 合并上次SQL引用的表] + R1 --> R2[外键再扩展 _expand_relations] + R2 --> R3[重新生成紧凑 Schema to_compact_string] + R3 -->|下一 attempt| E end - + subgraph 数据层 - J[Schema 管理 -SchemaManager] - K[向量数据库 -Chroma] - L[Few-shot 示例库 -JSONL] - M[DeepSeek API -大语言模型] - N[业务数据库 + J[SchemaManager +加载 Schema 元数据] + K[Chroma 持久化 +schema 向量索引 SchemaIndexer] + K2[Few-shot 向量库或 JSONL +FewShotSelector] + Emb[嵌入模型 +OpenAI 兼容 Embeddings API] + M[DeepSeek LLM +统一接口 deepseek.chat/select_tables +/validate_sql 等方法] + N[业务库 SQLAlchemy database_url] - J --> D + J --> D1 + J --> F1a + Emb --> D1 + Emb --> E1 K --> D1 - L --> E1 + K2 --> E1 + M --> N0 M --> D2 M --> E2 M --> F2 - N --> F1b + M --> F0a + M --> F0b + M --> B + N --> F1d end ``` @@ -103,66 +147,115 @@ database_url] ### 1. 意图分类 - **Dialog Classifier**:对用户输入的自然语言进行意图分类 -- **分类结果**:非查询意图直接返回对话回复,查询意图进入 SQL 生成流程 +- **分类模式**:支持 `rules`(仅规则)和 `hybrid`(混合)两种模式 + - `rules` 模式:仅使用关键词和短语规则,无 LLM 调用 + - `hybrid` 模式(默认):明显查数词/寒暄/元问题走规则快速返回;其余交 LLM 判断 +- **分类结果**:非查询意图(conversation)直接返回对话回复,查询意图(text2sql)进入 SQL 生成流程 +- **规则关键词**:包含查数词(查、查询、统计、汇总)、业务词(余额、交易、账户、持仓)、寒暄词(元问题、致谢等) -### 2. Schema Linker 算法 +### 2. Schema Linker 表选择模块 -- **向量检索粗筛**:使用 `SchemaIndexer` 对用户问题进行向量检索,获取相关表的候选列表 - - 利用预训练的嵌入模型将表名和描述转换为向量 +- **问句归一化(可选)**:`normalize_nl_question_for_text2sql` + - 将用户中英文问题归一为标准中文(默认开启,`TRANSLATE_EN_TO_ZH=true`) + - 使同一语义的中英文表述对齐,保证 SQL 一致性 + - 使用 LLM temperature=0 贪心解码 + +- **对话上文合并**:`_merge_dialog_for_model` + - 拼接上文与当前问句,控制总长度(默认 4000 字符用于检索,6000 字符用于选表) + - 优先保留当前问句完整 + +- **向量检索粗筛**:`_coarse_filter` → `SchemaIndexer` + - 使用嵌入模型将表结构转换为向量 - 使用 Chroma DB 进行相似度搜索 - - 返回 top-k 个最相关的表 + - 返回 top-k 个最相关的候选表 + - 可配置是否启用向量检索(`use_vector_search`) -* **LLM 精筛**:调用 DeepSeek 模型对候选表进行精筛,选择最相关的表 - - 构造包含表名和描述的提示 - - 让 LLM 基于用户问题选择最相关的表 -* **外键扩展**:自动添加与选中表相关的关联表,确保查询完整性 - - 分析表之间的外键关系 - - 自动添加被引用和引用当前表的关联表 +- **LLM 精筛**:`_llm_select_tables` → `select_tables` + - 构造包含候选表名和注释的提示 + - 调用 DeepSeek 模型选择最相关的表 + - 返回相关表列表和选择理由 -### 3. SQL Generator 算法 +- **经纪商维度表优先(可选)**:`_prioritize_broker_tables` + - 问题涉及对手方/经纪商时,优先纳入 `TSBBrokerContract` 与 `MCBroker` + - 避免仅选中报表视图却无 BrokerID -- **Few-shot 示例增强**:根据用户问题从示例库中选择相似的示例,增强 SQL 生成质量 - - 计算用户问题与示例库中问题的相似度 - - 选择 top-k 个最相似的高质量示例 +- **外键双向扩展**:`_expand_relations` + - 自动添加被引用的表(外键指向的表) + - 自动添加引用当前表的表(反向外键) + - 确保 JOIN 完整性 + +- **生成紧凑 Schema 字符串**:`to_compact_string` + - 根据选中的表生成紧凑的 Schema 描述 + - 包含表名、注释、字段、类型、主键、外键等信息 + +### 3. SQL Generator SQL生成模块 + +- **Few-shot 示例增强**:`FewShotSelector` + - 支持 Chroma 向量库或 JSONL 两种存储方式 + - 根据用户问题计算语义相似度 + - 选择 top-k 个最相似的高质量示例(支持评分、标签、难度过滤) - 将示例注入到提示中,指导 SQL 生成 -- **LLM 生成 SQL**:调用 DeepSeek 模型,根据选中的表结构和示例生成 SQL 语句 - - 构造包含表结构、用户问题和示例的提示 - - 指导 LLM 生成符合特定 SQL 方言的语句 - - 清理和规范化生成的 SQL + - 可配置 `FEWSHOT_USE_CHROMA=true/false` -### 4. Validator 算法 +- **LLM 生成 SQL**:`_generate_sql` → `deepseek.chat` + - 构造包含 Schema、用户问题、示例(可选)、对话上文的提示 + - 使用 `temperature=0.0, top_p=1.0` 贪心解码 + - T-SQL 特殊要求:使用方括号 `[]`、禁止反引号、`+` 拼接字符串、今日用 `CAST(GETDATE() AS DATE)` -- **程序验证**:检查 SQL 语法是否正确,表和列是否存在于 Schema 中,是否包含危险操作 - - 使用 sqlglot 进行语法检查 - - 验证表和列是否存在于 Schema 中 - - 检查是否包含危险操作(如 DROP、DELETE 等) -- **数据库试执行(探针)**:在程序验证全部通过后,于已配置的业务库上对 SQL 做只读试执行,结果用单一状态码表示,**不把具体行列数据交给 Validator LLM**: - - **1**:执行成功,且**有返回行或有返回列**(列名或数据行至少其一非空)→ **跳过** Validator 的语义审核 LLM,**直接将 SQL 交付用户** - - **0**:执行成功,但**既无列也无行** → 仍交付 SQL,**跳过** Validator 语义审核 LLM,另调 LLM 生成简短中文说明(可能原因 + 请用户补充条件),随响应一并返回 - - **-1**:执行失败 → 视为本次 SQL 不可用,**不交付**,进入重试;**不调用** Validator 语义审核 LLM - - 未配置 `database_url` 时跳过探针(可选告警),此时走 **LLM 语义验证** 作为兜底 -- **LLM 语义验证**:在未命中上述探针结果(即未配置库、未执行探针)时调用 DeepSeek,验证 SQL 语义是否符合用户意图 +- **SQL 方言归一**:`normalize_sql_for_dialect` + - 使用 sqlglot 清理和规范化生成的 SQL + - 提取 markdown 代码块中的 SQL + +### 4. Validator SQL验证模块 + +- **程序校验**(确定性规则,按顺序执行): + - **语法验证**:`validate_sql_syntax`(sqlglot) + - **Schema 一致性验证**:`validate_schema_consistency`(检查表和列是否存在于 Schema 中) + - **危险操作检查**:`check_dangerous_operations`(DROP、DELETE、UPDATE、INSERT、ALTER、TRUNCATE 等) + - **T-SQL 字面量规则检查**:`check_no_cjk_in_sql_string_literals`(禁止在字符串字面量中出现中文) + +- **数据库试执行(探针)**:`probe_sql_execution_status_ex` + - 仅在程序校验全部通过后执行 + - 需要配置 `database_url` 环境变量 + - 返回状态码: + - **1**:执行成功,且**有返回行** → **跳过 LLM 语义验证**,直接将 SQL 交付用户,可调用 `sql_probe_success_delivery_message` 生成交付说明 + - **0**:执行成功,但**数据行数为 0** → 跳过 LLM 语义验证,调用 `empty_result_user_feedback` 生成补充说明和追问引导 + - **-1**:执行失败 → 进入重试流程,不交付 SQL + - **None**:未配置 `database_url`,跳过探针 + +- **LLM 语义验证**:`validate_sql`(仅探针为 None 时调用) - 构造包含 SQL、表结构和用户问题的提示 - 让 LLM 评估 SQL 是否正确回答了用户问题 - 收集错误和警告信息 -- **结果评估**:验证通过则返回有效 SQL,失败则进入重试逻辑 - - 最多重试 max\_retry 次 - - 每次重试使用相同的表结构,但重新生成 SQL + - 如果程序校验已通过 Schema 一致性,则过滤掉 LLM 常误报的 `unknown_table` 和 `unknown_column` 错误 + +- **重试逻辑**: + - 最多重试 `max_retry` 次(默认 2) + - 每次重试: + - 合并上次失败 SQL 中引用的表(`extract_tables_from_sql`) + - 外键再扩展 + - 重新生成紧凑 Schema + - 将上次校验错误反馈给 LLM ### 5. 数据层支持 -- **Schema 管理**:管理数据库表结构和元数据 - - 从 JSON 文件加载表结构 +- **Schema 管理**:`SchemaManager` 管理数据库表结构和元数据 + - 从 JSON/G3SB 文件加载表结构 - 提供表和列的查询接口 - 分析表之间的外键关系 -- **向量数据库**:存储表和列的向量表示,用于快速检索 - - 使用 Chroma DB 存储向量 + - 生成紧凑的 Schema 描述字符串 +- **向量数据库**:`SchemaIndexer` + Chroma DB 存储和检索表向量 - 支持增量更新和查询 -- **Few-shot 示例库**:存储高质量的自然语言到 SQL 的示例 - - 从 JSONL 文件加载示例 - - 支持基于相似度的示例检索 -- **DeepSeek API**:提供大语言模型能力,用于表选择、SQL 生成和验证 - - 调用 DeepSeek 聊天模型 + - 可配置向量数据库路径 +- **Few-shot 示例库**:`FewShotSelector` 管理高质量示例 + - 支持 Chroma 向量库或 JSONL 文件两种存储方式 + - 基于语义相似度检索示例 + - 支持评分、标签、难度过滤 +- **DeepSeek API**:统一接口提供大语言模型能力 + - `deepseek.chat` 通用对话 + - `select_tables` 表选择 + - `validate_sql` SQL 验证 + - `normalize_nl_question` 问句归一 - 支持不同的模型参数配置 ## 技术栈 @@ -177,12 +270,16 @@ database_url] ## 算法特点 -1. **多智能体协作**:Schema Linker、SQL Generator、Validator 三个智能体协同工作,各负责专门任务,形成完整的处理流程 -2. **向量检索增强**:Schema Linker 使用向量数据库快速筛选相关表,提高表选择效率 -3. **Few-shot 学习**:SQL Generator 利用示例库增强 SQL 生成质量,学习最佳实践 -4. **多层验证**:程序校验 + 库上探针;探针为 1/0/-1 时不再走 Validator 语义 LLM(1 直接交付、0 附带无数据说明、-1 重试),未配置库时仍以 LLM 语义验证兜底 -5. **自动关联表扩展**:Schema Linker 通过外键关系自动扩展相关表,提高查询完整性 -6. **可配置性**:支持多种配置参数,如温度、最大重试次数、向量检索开关等,适应不同场景需求 +1. **模块化编排架构**:Text2SQLOrchestrator 统一编排表选择、SQL 生成、SQL 验证三大模块,各模块协同工作形成完整处理流程 +2. **问句归一化**:支持中英文问句归一为标准中文(默认开启),保证同一语义的查询生成一致的 SQL +3. **向量检索增强**:Schema Linker 使用 Chroma 向量数据库快速筛选相关表,提高表选择效率,可配置启用/禁用 +4. **经纪商维度表优先**:针对证券/期货业务场景,涉及对手方/经纪商的问题自动优先纳入经纪商维度表 +5. **Few-shot 学习**:SQL Generator 利用示例库增强 SQL 生成质量,支持 Chroma 向量库或 JSONL 存储方式,可配置启用/禁用 +6. **多层验证机制**:程序校验(语法 + Schema 一致性 + 危险操作 + T-SQL 字面量规则)+ 库上探针 + LLM 语义验证 +7. **智能探针分流**:探针返回 1(有数据)直接交付、0(无数据)附带说明引导用户补充条件、-1(失败)自动重试,None 时走 LLM 语义验证兜底 +8. **自动关联表扩展**:通过外键关系双向扩展相关表,确保 JOIN 完整性 +9. **对话上下文支持**:支持多轮对话的上下文合并,处理指代消解和续问 +10. **可配置性**:支持多种配置参数,如最大重试次数、向量检索开关、Few-shot 配置、问句归一开关等,适应不同场景需求 ## Few-shot 学习详细说明 @@ -192,28 +289,54 @@ Few-shot 学习是一种机器学习方法,指通过少量示例来指导模 ### 在 Text2SQL 系统中的应用 -在 Text2SQL 系统中,Few-shot 学习主要应用于 SQL Generator 智能体,具体流程如下: +在 Text2SQL 系统中,Few-shot 学习应用于 SQL 生成阶段,通过 `FewShotSelector` 组件实现,具体流程如下: -1. **示例库构建**:系统维护一个存储高质量自然语言到 SQL 示例的库,从 JSONL 文件加载。这些示例包含各种类型的 SQL 查询场景,如简单查询、复杂连接、聚合操作等。 -2. **相似度匹配**:当用户提出自然语言问题时,系统会计算该问题与示例库中问题的语义相似度。 -3. **示例选择**:基于相似度排序,选择最相关的 top-k 个高质量示例。 -4. **提示增强**:将这些示例注入到给大语言模型的提示中,指导模型生成更准确的 SQL 语句。 -5. **生成指导**:模型参考示例的结构和风格,结合用户问题和表结构,生成符合要求的 SQL 语句。 +1. **示例库构建**:系统维护一个存储高质量自然语言到 SQL 示例的库,支持两种存储方式: + - **Chroma 向量库**:`FEWSHOT_USE_CHROMA=true` 时使用,提供高效的相似度检索 + - **JSONL 文件**:传统方式,从 `all_samples.jsonl` 加载示例 + - 每条示例包含:问题(中文/英文)、SQL 语句、评分、标签、难度等级等元数据 + +2. **相似度匹配**:当用户提出自然语言问题时,使用嵌入模型将问题转换为向量,计算向量相似度 + +3. **示例选择**:基于相似度排序,支持多维度过滤: + - `top_k`:选择最相关的 top-k 个示例(默认 3) + - `min_rating`:最低评分过滤(默认 7) + - `required_tags`:必须包含的标签 + - `max_difficulty`:最大难度过滤 + +4. **提示增强**:将选中的示例注入到给大语言模型的提示中,格式为: + ``` + 参考以下相似示例的SQL编写风格: + + 示例 1: + 问题:xxx + SQL: + xxx + + 【当前Schema】 + xxx + ``` + +5. **生成指导**:模型参考示例的结构和风格,结合用户问题和表结构,生成符合要求的 SQL 语句 ### 技术实现 -- **示例存储**:使用 JSONL 格式存储示例,每条示例包含自然语言问题和对应的 SQL 语句。 -- **相似度计算**:利用预训练的嵌入模型将问题转换为向量,计算向量相似度。 -- **示例选择**:根据相似度得分选择最相关的示例。 -- **提示构造**:将选中的示例与用户问题、表结构一起构造提示,确保模型能理解任务要求。 +- **示例存储**:支持 Chroma 向量库持久化或 JSONL 文件加载两种方式 +- **向量索引**: + - Chroma 模式:使用 Chroma DB 存储和检索向量 + - JSONL 模式:内存 numpy 数组 + 可选 .npy 缓存 +- **相似度计算**:利用预训练的嵌入模型将问题转换为向量,使用余弦相似度计算 +- **示例选择**:根据相似度得分和元数据过滤选择最相关的示例 +- **提示构造**:将选中的示例与用户问题、表结构一起构造提示,确保模型能理解任务要求 ### 优势 -1. **减少数据需求**:不需要大量的标注数据,仅需少量高质量示例。 -2. **提高生成质量**:通过示例指导,模型能生成更符合特定场景的 SQL 语句。 -3. **学习最佳实践**:示例库可以包含领域专家编写的高质量 SQL,使模型学习到最佳实践。 -4. **适应不同场景**:通过扩展示例库,可以适应不同领域和复杂度的 SQL 生成需求。 -5. **灵活性**:可以根据具体应用场景调整示例库,提高系统的适应性。 +1. **减少数据需求**:不需要大量的标注数据,仅需少量高质量示例 +2. **提高生成质量**:通过示例指导,模型能生成更符合特定场景的 SQL 语句 +3. **学习最佳实践**:示例库可以包含领域专家编写的高质量 SQL,使模型学习到最佳实践 +4. **适应不同场景**:通过扩展示例库,可以适应不同领域和复杂度的 SQL 生成需求 +5. **灵活性**:支持 Chroma 和 JSONL 两种存储方式,可根据场景选择 +6. **高质量保障**:通过评分、标签、难度等元数据过滤,确保选用高质量示例 ### 应用效果 @@ -223,20 +346,29 @@ Few-shot 学习是一种机器学习方法,指通过少量示例来指导模 - 生成更符合用户意图的 SQL 语句 - 减少生成错误 - 提高系统的泛化能力 +- 通过示例学习特定业务场景的 SQL 编写风格 ## 性能优化 -1. **向量索引预构建**:提前构建向量索引,加速首次查询 -2. **缓存机制**:缓存常见查询的结果,提高响应速度 -3. **并行处理**:对多个查询进行并行处理,提高系统吞吐量 -4. **模型调优**:调整 LLM 参数,平衡生成质量和速度 +1. **向量索引预构建**:`build_vector_index` 方法提前构建向量索引,加速首次查询 +2. **懒加载机制**:向量索引在首次使用时才加载,避免启动开销 +3. **Few-shot 缓存**:向量缓存到 `.npy` 文件,换模型自动重建 +4. **贪心解码**:SQL 生成使用 `temperature=0.0` 贪心解码,提高一致性和速度 +5. **库探针优化**:仅取 1 行数据(`max_rows=1`)判断探针状态,最小化数据库负载 +6. **SQL 执行优化**: + - 聚合查询优先直接执行而非 COUNT 包裹 + - 自动去掉派生表内无意义的 ORDER BY + - 支持多种 SQL 方言的语法处理 ## 扩展性 -1. **支持多种数据库**:通过配置支持不同的 SQL 方言 -2. **可插拔的 LLM**:支持替换不同的大语言模型 -3. **自定义示例库**:可根据特定领域扩展示例库 +1. **支持多种数据库**:通过 sqlglot 支持不同的 SQL 方言(T-SQL、MySQL、PostgreSQL 等) +2. **可插拔的 LLM**:统一通过 `DeepSeekClient` 接口调用,支持替换不同的大语言模型 +3. **可配置的嵌入模型**:通过环境变量支持 OpenAI、ModelScope、DashScope 等多种嵌入 API +4. **自定义示例库**:可根据特定领域扩展示例库,支持 Chroma 或 JSONL 格式 +5. **Schema 灵活加载**:支持从 JSON 文件或 G3SB 格式加载 Schema +6. **意图分类可扩展**:支持 rules 和 hybrid 两种模式,可扩展更多分类策略 ## 总结 -Text2SQL 算法通过多智能体协作、向量检索和大语言模型技术,实现了从自然语言到 SQL 的高效准确转换。算法流程清晰,逻辑完善,具有良好的可扩展性和可配置性,能够满足不同场景下的 SQL 生成需求。 +Text2SQL 算法通过模块化编排架构、向量检索、问句归一化和 Few-shot 学习技术,实现了从自然语言到 SQL 的高效准确转换。系统采用多层验证机制(程序校验 + 库上探针 + LLM 语义验证),配合智能探针分流策略,确保生成的 SQL 既正确又符合用户意图。算法流程清晰,逻辑完善,具有良好的可扩展性和可配置性,能够满足证券/期货等金融业务场景下的 SQL 生成需求。