diff --git a/IMPACT_ANALYSIS.md b/IMPACT_ANALYSIS.md
index 74c01bc..83b1027 100644
--- a/IMPACT_ANALYSIS.md
+++ b/IMPACT_ANALYSIS.md
@@ -592,3 +592,83 @@
## 6. 配置变更
- 无新增环境变量;日志路径固定为仓库根下 `logs/text2sql_api.log`。
+
+---
+
+# Impact Analysis Report — Chroma Few-shot 黄金 SQL 条件适配(追加)
+
+## 1. 改动概览
+
+- **背景与目标**:`data/embeddings/chroma_fewshot` 中存储的问答-SQL 视为已校验正确答案;当向量检索与当前问题足够相似时,不再仅作「风格示例」,而是以该 SQL 为骨架,仅让 LLM 调整 WHERE/HAVING 等与用户问题相关的**特殊条件**。
+- **涉及模块**:`backend/config/prompts.py`(`GOLDEN_SQL_ADAPT_*`)、`backend/utils/fewshot_selector.py`(`select_best_with_score`、`is_chroma_backend`)、`backend/agents/orchestrator.py`(`_generate_sql_golden_adapt`、`_generate_sql` 分支)、`api_server.py`(响应中可选 `fewshot_golden`)。
+- **改动类型**:生成策略增强(默认开启,可用环境变量关闭或调阈值)。
+
+## 2. 方法级改动
+
+| 位置 | 变更 |
+|------|------|
+| `FewShotSelector` | 新增 `_passes_filters`、`select_best_with_score`;`is_chroma_backend` 属性。 |
+| `Text2SQLOrchestrator._generate_sql` | 在 Chroma 且满足阈值时调用 `_generate_sql_golden_adapt` 并设置 `_last_fewshot_golden`;否则沿用原 few-shot 注入 + 全量生成。 |
+| `_generate_sql_golden_adapt` | 使用 `GOLDEN_SQL_ADAPT_SYSTEM/USER`,强调保留 JOIN/SELECT 结构、仅改条件。 |
+| `generate` | `metadata` 增加 `fewshot_golden_reuse`、`fewshot_golden_qid`、`fewshot_golden_score`;失败路径亦回传(若曾走黄金分支)。 |
+| `_nl_dict_from_generation` | 若存在黄金复用,增加 `fewshot_golden` 字段供前端展示/调试。 |
+
+## 3. 破坏性变更
+
+- **否**。未命中阈值时行为与原先「多示例风格参考」一致。
+
+## 4. 配置变更
+
+| 环境变量 | 含义 | 默认 |
+|----------|------|------|
+| `FEWSHOT_GOLDEN_REUSE` | 是否启用黄金 SQL 条件适配 | `true` |
+| `FEWSHOT_GOLDEN_ONLY_CHROMA` | 是否仅在 Chroma 后端启用(与 `chroma_fewshot` 策略一致) | `true` |
+| `FEWSHOT_GOLDEN_MIN_SCORE` | 最低相似度(1-距离,约等于余弦相似度) | `0.88` |
+| `FEWSHOT_GOLDEN_SQL_PROMPT_MAX` | 写入提示词的标准答案 SQL 最大字符数 | `16000` |
+
+## 5. 风险与回滚
+
+- **风险**:阈值过低可能把不太相似的问题强行套在同一 SQL 上;过高则很少触发黄金分支。
+- **回滚**:设 `FEWSHOT_GOLDEN_REUSE=false` 或回退相关提交。
+
+**回滚方式是否简单**:是。
+
+---
+
+# Impact Analysis Report — 流式 SSE sql_gen 分片与节流
+
+## 1. 改动概览
+
+- **背景与目标**:前端约定每条 SSE 为 `data: {"stage":"sql_gen","stream_kind":"content","content":"..."}`;需在服务端控制拆成「逐字符多条」或「与 LLM delta 一致」,并支持可选分片间隔。
+- **涉及模块**:`api_server.py`(`/g3sb/api/nl/chat/stream`、`_iter_sql_gen_content_pieces`、`_sse_stream_text_chunks`)。
+- **改动类型**:行为调整(默认分片粒度默认更贴近前端的 `char`;可配置回 `delta`)。
+
+## 2. 方法级改动
+
+| 位置 | 变更 |
+|------|------|
+| `NLChatRequest` | 新增可选 `sql_stream_granularity`(`sqlStreamGranularity`);明确 `streaming_throttle` 为相邻 content 间隔毫秒。 |
+| `_iter_sql_gen_content_pieces` | 支持传入 `mode`;未设置时 `SSE_SQL_GEN_SPLIT` 默认 `char`。 |
+| `_chat_stream_events` | sql_gen 循环按粒度拆分后对每条 `data` 可选 `asyncio.sleep(throttle)`;寒暄分支 `_sse_stream_text_chunks` 支持相同节流。 |
+
+## 3. 调用方与影响范围
+
+- **调用方**:仅 SSE 流式客户端;非流式 `/g3sb/api/nl/chat` 不变。
+- **破坏性变更**:否。未传新字段时:粒度由 `SSE_SQL_GEN_SPLIT` 决定(默认 `char`,事件条数多于旧版 `delta`);若需旧行为可设 `SSE_SQL_GEN_SPLIT=delta` 或请求体 `sqlStreamGranularity: "delta"`。
+
+## 4. 配置变更
+
+| 环境变量 | 含义 | 默认 |
+|----------|------|------|
+| `SSE_SQL_GEN_SPLIT` | `char`:逐 Unicode 标量多条 SSE;`delta`:与 LLM 增量一致 | `char`(未设置 env 时由代码默认) |
+
+## 5. 风险与回滚
+
+- **风险级别**:低。`char` 模式下 SSE 条数增加,带宽与前端拼接次数上升。
+- **回滚**:设 `SSE_SQL_GEN_SPLIT=delta` 或请求传 `sqlStreamGranularity: "delta"`。
+
+**回滚方式是否简单**:是。
+
+## 6. 验证与测试
+
+- 已执行:`python -m py_compile api_server.py`。
diff --git a/__pycache__/api_server.cpython-312.pyc b/__pycache__/api_server.cpython-312.pyc
index 9ad4c41..937744a 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 b174a07..031f1a4 100644
--- a/api_server.py
+++ b/api_server.py
@@ -9,7 +9,6 @@ from __future__ import annotations
import os
import sys
import json
-import html as html_lib
import asyncio
import logging
from pathlib import Path
@@ -151,7 +150,15 @@ class NLChatRequest(BaseModel):
session_id: Optional[str] = Field(None, description="会话ID")
visitor_biz_id: Optional[str] = Field(None, description="访客业务ID")
user_id: Optional[str] = Field(None, description="用户ID")
- streaming_throttle: Optional[int] = Field(None, description="流式节流参数")
+ streaming_throttle: Optional[int] = Field(
+ None,
+ description="流式节流:相邻 content 分片之间的间隔毫秒数(0/None 表示不延迟)",
+ )
+ sql_stream_granularity: Optional[str] = Field(
+ None,
+ validation_alias=AliasChoices("sql_stream_granularity", "sqlStreamGranularity"),
+ description="SQL 生成流式分片:delta | char;不传则使用环境变量 SSE_SQL_GEN_SPLIT(默认 char)",
+ )
class IntentPayload(BaseModel):
@@ -247,6 +254,23 @@ def _sse_data(obj: Dict[str, Any]) -> bytes:
return f"data: {json.dumps(obj, ensure_ascii=False)}\n\n".encode("utf-8")
+def _iter_sql_gen_content_pieces(text: str, mode: Optional[str] = None) -> List[str]:
+ """
+ 将 LLM 流式片段再拆成前端期望的多条 SSE(与 chatStore onDelta 累加一致)。
+ 每条均为:{"stage": "sql_gen", "stream_kind": "content", "content": "..."}
+
+ mode / 环境变量 SSE_SQL_GEN_SPLIT:
+ - char(请求默认):按 Unicode 标量逐字符发送(与「用户」「问题」逐条 data 一致)
+ - delta:与上游 LLM 每次 delta 一致(块更大、事件更少)
+ """
+ if not text:
+ return []
+ raw = (mode or os.getenv("SSE_SQL_GEN_SPLIT", "char") or "char").strip().lower()
+ if raw in ("delta", "none", "0", "false"):
+ return [text]
+ return list(text)
+
+
# 与前端 chatStore onDelta 一致:多次 { stage, stream_kind: content, content } 累加
_DEFAULT_SSE_CHUNK_CHARS = int(os.getenv("SSE_STREAM_CHUNK_CHARS", "64"))
@@ -256,16 +280,21 @@ async def _sse_stream_text_chunks(
content: str,
*,
chunk_size: Optional[int] = None,
+ throttle_ms: Optional[int] = None,
) -> AsyncIterator[bytes]:
"""将长文本拆成多段 SSE,便于浏览器逐段渲染(流式)。"""
if not content:
return
size = max(8, chunk_size or _DEFAULT_SSE_CHUNK_CHARS)
+ delay = (throttle_ms or 0) / 1000.0 if throttle_ms and throttle_ms > 0 else 0.0
for i in range(0, len(content), size):
yield _sse_data(
{"stage": stage, "stream_kind": "content", "content": content[i : i + size]}
)
- await asyncio.sleep(0)
+ if delay:
+ await asyncio.sleep(delay)
+ else:
+ await asyncio.sleep(0)
def _conversation_nl_dict(reply: str) -> Dict[str, Any]:
@@ -584,29 +613,14 @@ def _nl_dict_from_generation(result: GenerationResult) -> Dict[str, Any]:
qz = result.metadata.get("question_zh_normalized")
if qo and qz:
payload["query_normalization"] = {"original": qo, "zh": qz}
+ if result.metadata.get("fewshot_golden_reuse"):
+ payload["fewshot_golden"] = {
+ "qid": result.metadata.get("fewshot_golden_qid"),
+ "score": result.metadata.get("fewshot_golden_score"),
+ }
return payload
-def _sql_gen_stream_html(result: GenerationResult) -> str:
- sql = (result.sql or "").strip()
- inner = json.dumps({"sql": sql}, ensure_ascii=False)
- parts = [f"{inner}"]
- explain = ""
- if result.metadata.get("sql_delivery_message"):
- parts.append(
- f'
{html_lib.escape(str(result.metadata["sql_delivery_message"]))}
'
- )
- if result.metadata.get("db_empty_feedback"):
- explain = str(result.metadata["db_empty_feedback"])
- elif result.warnings:
- explain = str(result.warnings[0])
- elif result.errors:
- explain = "; ".join(str(e) for e in result.errors[:5])
- if explain.strip():
- parts.append(f'{html_lib.escape(explain)}
')
- return "".join(parts)
-
-
async def _run_generate(
question: str,
top_k: int = 20,
@@ -684,7 +698,9 @@ async def _chat_stream_events(request: NLChatRequest) -> AsyncIterator[bytes]:
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):
+ async for pkt in _sse_stream_text_chunks(
+ "chat", reply, throttle_ms=request.streaming_throttle
+ ):
yield pkt
data_dict = _conversation_nl_dict(reply)
yield _sse_data({"code": 200, "msg": "success", "data": data_dict})
@@ -692,8 +708,65 @@ async def _chat_stream_events(request: NLChatRequest) -> AsyncIterator[bytes]:
return
yield _sse_data({"stage": "orchestrator", "stream_kind": "content", "content": "DATA_QUERY"})
+ dialect = resolve_sql_dialect(os.getenv("TEXT2SQL_DIALECT", "sqlserver"))
+ dc_ctx = (dialog_block or "").strip() or None
+ top_k = 20
+ loop = asyncio.get_running_loop()
+ logger.info(
+ "[GEN/API/stream] dialect=%s top_k=%s dialog_context_chars=%s question_len=%s preview=%r",
+ dialect,
+ top_k,
+ len(dc_ctx) if dc_ctx else 0,
+ len(text),
+ text[:300] + ("…" if len(text) > 300 else ""),
+ )
try:
- result = await _run_generate(text, dialog_context=dialog_block or None, request=request)
+ async with _maybe_override_orch_llm(orch, request) as o:
+ chunk_queue: asyncio.Queue[Optional[str]] = asyncio.Queue()
+ holder: Dict[str, Any] = {}
+
+ def _run_generate_sync() -> None:
+ try:
+ holder["result"] = o.generate(
+ question=text.strip(),
+ dialect=dialect,
+ top_k_candidates=top_k,
+ dialog_context=dc_ctx,
+ sql_stream_callback=lambda c: loop.call_soon_threadsafe(
+ chunk_queue.put_nowait, c
+ ),
+ )
+ except Exception as e:
+ holder["error"] = e
+ finally:
+ loop.call_soon_threadsafe(chunk_queue.put_nowait, None)
+
+ gen_task = asyncio.create_task(asyncio.to_thread(_run_generate_sync))
+ throttle_ms = request.streaming_throttle or 0
+ sql_chunk_delay = throttle_ms / 1000.0 if throttle_ms > 0 else 0.0
+ gran = (
+ (request.sql_stream_granularity or "").strip()
+ or os.getenv("SSE_SQL_GEN_SPLIT", "char")
+ or "char"
+ ).lower()
+ while True:
+ piece = await chunk_queue.get()
+ if piece is None:
+ break
+ for frag in _iter_sql_gen_content_pieces(piece, mode=gran):
+ yield _sse_data(
+ {
+ "stage": "sql_gen",
+ "stream_kind": "content",
+ "content": frag,
+ }
+ )
+ if sql_chunk_delay:
+ await asyncio.sleep(sql_chunk_delay)
+ await gen_task
+ if holder.get("error"):
+ raise holder["error"]
+ result = holder["result"]
except Exception as e:
logger.error(f"[API/stream] 生成异常: {e}", exc_info=True)
yield _sse_data({"code": 500, "msg": str(e), "data": None})
@@ -727,9 +800,6 @@ async def _chat_stream_events(request: NLChatRequest) -> AsyncIterator[bytes]:
except Exception:
pass
- html = _sql_gen_stream_html(result)
- async for pkt in _sse_stream_text_chunks("sql_gen", html):
- yield pkt
data_dict = _nl_dict_from_generation(result)
msg = "success" if result.valid else "partial"
yield _sse_data({"code": 200, "msg": msg, "data": data_dict})
@@ -856,7 +926,10 @@ async def nl_chat(request: NLChatRequest):
async def nl_chat_stream(request: NLChatRequest):
"""
自然语言对话流式接口(SSE)。
- 事件体为 JSON:分片 delta `{stage, stream_kind, content}` 或结束包 `{code, msg, data}`。
+ 事件体为 JSON:分片 `data: {"stage","stream_kind","content"}`(例:sql_gen 时
+ `stream_kind` 为 `content`)或结束包 `{code, msg, data}`。
+ 可选:`sql_stream_granularity` / `sqlStreamGranularity`(delta|char,未传则 `SSE_SQL_GEN_SPLIT`,默认 char)、
+ `streaming_throttle`(相邻 content 分片间隔毫秒)。
"""
return StreamingResponse(
_chat_stream_events(request),
diff --git a/backend/agents/__pycache__/orchestrator.cpython-312.pyc b/backend/agents/__pycache__/orchestrator.cpython-312.pyc
index 40bd50d..081ed23 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 935dd0e..99eef6d 100644
--- a/backend/agents/orchestrator.py
+++ b/backend/agents/orchestrator.py
@@ -5,13 +5,13 @@ Text2SQL 多智能体编排器
import logging
import os # 新增
-from typing import Dict, List, Optional, Tuple
+from typing import Callable, Dict, List, Optional, Tuple
from dataclasses import dataclass, field
from schema.manager import SchemaManager
from schema.indexer import SchemaIndexer
from llm.deepseek_client import DeepSeekClient, DeepSeekConfig
-from utils.fewshot_selector import FewShotSelector # 新增
+from utils.fewshot_selector import ExperienceSample, FewShotSelector # 新增
logger = logging.getLogger(__name__)
@@ -121,6 +121,9 @@ class Text2SQLOrchestrator:
logger.warning(f"Few-shot加载失败: {e},将使用标准生成")
self.fewshot_enabled = False
+ # 若本轮走「Chroma 黄金 SQL 条件适配」,在 metadata 中回传 qid/分数
+ self._last_fewshot_golden: Optional[Tuple[str, float]] = None
+
logger.info(
f"[OK] Text2SQLOrchestrator初始化完成: "
f"max_retry={max_retry}, use_vector_search={use_vector_search}"
@@ -335,6 +338,118 @@ class Text2SQLOrchestrator:
return expanded
+ def _sql_chat_completion_text(
+ self,
+ messages: List[Dict[str, str]],
+ *,
+ sql_stream_callback: Optional[Callable[[str], None]] = None,
+ ) -> str:
+ """
+ SQL 生成相关的一次 LLM 调用:可选流式回调(将原始 completion 文本分片回传)。
+ 无回调或非流式失败时与非流式 chat 行为一致。
+ """
+ ds = self.deepseek
+ if sql_stream_callback is not None and hasattr(ds, "chat_stream"):
+ try:
+ parts: List[str] = []
+ for piece in ds.chat_stream(messages, temperature=0.0, top_p=1.0):
+ if piece:
+ parts.append(piece)
+ sql_stream_callback(piece)
+ return "".join(parts).strip()
+ except Exception as e:
+ logger.warning("[GEN] chat_stream 失败,回退非流式: %s", e)
+ msg = ds.chat(messages, temperature=0.0, top_p=1.0)
+ return (msg.content or "").strip()
+
+ def _generate_sql_golden_adapt(
+ self,
+ question: str,
+ schema_str: str,
+ dialect: str,
+ golden: ExperienceSample,
+ golden_score: float,
+ validation_feedback: Optional[str],
+ dialog_context: Optional[str],
+ sql_stream_callback: Optional[Callable[[str], None]] = None,
+ ) -> str:
+ """
+ Chroma few-shot 库内 SQL 视为正确答案:在高分相似命中下,仅让模型调整条件/字面量以匹配当前问题。
+ """
+ from config.prompts import GOLDEN_SQL_ADAPT_SYSTEM, GOLDEN_SQL_ADAPT_USER
+ from utils.sql_parser import normalize_sql_for_dialect
+
+ gsql = (golden.sql or "").strip()
+ max_sql = int(os.getenv("FEWSHOT_GOLDEN_SQL_PROMPT_MAX", "16000"))
+ if len(gsql) > max_sql:
+ gsql = gsql[:max_sql] + "\n-- …(标准答案过长,已截断)"
+
+ dialect_label = dialect
+ if dialect == "tsql":
+ dialect_label = "Microsoft SQL Server (T-SQL)"
+
+ dc = (dialog_context or "").strip()
+ prefix = ""
+ if dc:
+ prefix = (
+ "【对话上文】(用于理解指代与续问条件;请结合「当前用户问题」调整 WHERE 等。)\n"
+ f"{dc}\n\n"
+ )
+
+ user_content = prefix + GOLDEN_SQL_ADAPT_USER.format(
+ golden_score=golden_score,
+ golden_question=(golden.question_zh or "").strip(),
+ golden_sql=gsql,
+ schema=schema_str,
+ question=question,
+ dialect=dialect_label,
+ )
+ if dialect == "tsql":
+ user_content += (
+ "\n\n【硬性要求】目标库为 SQL Server(T-SQL):禁止使用 MySQL 反引号 `;"
+ "标识符如需引用请使用方括号,例如 [TableName]、[ColumnName]。"
+ "字符串连接使用 `+`(与系统提示中的标准版式范例一致)。"
+ "「今日」「当天」等与日期列比较时,使用 `CAST(GETDATE() AS DATE)`,"
+ "**禁止** `CURDATE()`、`NOW()`、`CURRENT_DATE`(MySQL)。"
+ "条件请使用 T-SQL 惯用写法(例如 IS NOT NULL)。"
+ "排版仍须遵守:关键字大写、SELECT 每列一行缩进、WHERE 续行以 AND 开头、PascalCase 英文别名。"
+ "\n**禁止**在单引号字符串字面量或 `N'…'` 中出现任何中日韩文字;"
+ "业务中文须映射为 Schema 注释中的代码或通过维表 JOIN,勿写 `= '过户费'` 这类比对。"
+ )
+
+ if validation_feedback:
+ user_content += (
+ "\n\n【上次校验未通过】请在保留标准答案主干的前提下修正 SQL;"
+ "表名、列名必须与「当前 Schema 片段」中完全一致。\n"
+ f"{validation_feedback}"
+ )
+
+ messages = [
+ {"role": "system", "content": GOLDEN_SQL_ADAPT_SYSTEM},
+ {"role": "user", "content": user_content},
+ ]
+
+ logger.info(
+ "[GEN] 黄金 few-shot 条件适配: qid=%s score=%.4f",
+ golden.qid,
+ golden_score,
+ )
+ sql = self._sql_chat_completion_text(
+ messages, sql_stream_callback=sql_stream_callback
+ )
+
+ if "```sql" in sql:
+ sql = sql[sql.find("```sql") + 6 : sql.find("```", sql.find("```sql") + 6)].strip()
+ elif "```" in sql:
+ sql = sql[sql.find("```") + 3 : sql.find("```", sql.find("```") + 3)].strip()
+
+ sql = normalize_sql_for_dialect(sql, dialect)
+
+ 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 _generate_sql(
self,
question: str,
@@ -342,6 +457,7 @@ class Text2SQLOrchestrator:
dialect: str = "tsql",
validation_feedback: Optional[str] = None,
dialog_context: Optional[str] = None,
+ sql_stream_callback: Optional[Callable[[str], None]] = None,
) -> str:
"""
SQL生成(SQL Generator Agent)
@@ -353,6 +469,7 @@ class Text2SQLOrchestrator:
validation_feedback: 非空时附加到用户提示(重试时传入上次校验错误)
dialog_context: 前几轮对话摘要;与 ``question`` 一并供指代消解与续问。
+ sql_stream_callback: 若提供且 LLM 支持 chat_stream,则在 SQL 主生成/黄金适配时流式回传原始文本分片。
Returns:
SQL语句
@@ -360,11 +477,57 @@ class Text2SQLOrchestrator:
from config.prompts import SQL_GENERATOR_SYSTEM, SQL_GENERATOR_USER
from utils.sql_parser import normalize_sql_for_dialect
+ self._last_fewshot_golden = None
dc = (dialog_context or "").strip()
fewshot_question = question
if dc:
fewshot_question = f"{dc}\n\n【当前问】{question}"
+ golden_reuse = os.getenv("FEWSHOT_GOLDEN_REUSE", "true").lower() in (
+ "1",
+ "true",
+ "yes",
+ )
+ only_chroma = os.getenv("FEWSHOT_GOLDEN_ONLY_CHROMA", "true").lower() in (
+ "1",
+ "true",
+ "yes",
+ )
+ golden_min = float(os.getenv("FEWSHOT_GOLDEN_MIN_SCORE", "0.88"))
+ chroma_ok = bool(
+ self.fewshot_selector and self.fewshot_selector.is_chroma_backend
+ )
+ if only_chroma and not chroma_ok:
+ golden_reuse = False
+
+ if (
+ golden_reuse
+ and self.fewshot_enabled
+ and self.fewshot_selector
+ ):
+ best = self.fewshot_selector.select_best_with_score(
+ fewshot_question,
+ min_rating=self.fewshot_min_rating,
+ )
+ if (
+ best
+ and best[1] >= golden_min
+ and (best[0].sql or "").strip()
+ ):
+ ex, sc = best
+ sql_out = self._generate_sql_golden_adapt(
+ question=question,
+ schema_str=schema_str,
+ dialect=dialect,
+ golden=ex,
+ golden_score=sc,
+ validation_feedback=validation_feedback,
+ dialog_context=dialog_context,
+ sql_stream_callback=sql_stream_callback,
+ )
+ self._last_fewshot_golden = (ex.qid, sc)
+ return sql_out
+
# Few-shot 增强
if self.fewshot_enabled and self.fewshot_selector:
try:
@@ -431,8 +594,9 @@ class Text2SQLOrchestrator:
]
# 与选表一致:生成阶段默认贪心解码,减少同一中文问题多次 SQL 不一致
- response = self.deepseek.chat(messages, temperature=0.0, top_p=1.0)
- sql = response.content.strip()
+ sql = self._sql_chat_completion_text(
+ messages, sql_stream_callback=sql_stream_callback
+ )
# 清理可能的markdown代码块
if "```sql" in sql:
@@ -609,6 +773,7 @@ class Text2SQLOrchestrator:
top_k_candidates: int = 20,
include_schema_in_result: bool = False,
dialog_context: Optional[str] = None,
+ sql_stream_callback: Optional[Callable[[str], None]] = None,
) -> GenerationResult:
"""
主生成流程
@@ -619,10 +784,12 @@ class Text2SQLOrchestrator:
top_k_candidates: 粗筛候选表数量
include_schema_in_result: 结果中是否包含使用的Schema字符串
dialog_context: 前几轮对话可读摘要;选表、向量粗筛、SQL 生成与无数据说明会参考
+ sql_stream_callback: 可选;SQL 主生成 LLM 输出分片回调(用于 API SSE)
Returns:
GenerationResult对象
"""
+ self._last_fewshot_golden = None
original_question = (question or "").strip()
translation_meta: Dict = {}
work_question = original_question
@@ -726,6 +893,7 @@ class Text2SQLOrchestrator:
dialect,
validation_feedback=feedback,
dialog_context=dc_raw or None,
+ sql_stream_callback=sql_stream_callback,
)
last_sql = sql
except Exception as e:
@@ -770,6 +938,11 @@ class Text2SQLOrchestrator:
meta["sql_delivery_message"] = None
if dc_raw:
meta["dialog_context_chars"] = len(dc_raw)
+ if self._last_fewshot_golden:
+ gq, gsc = self._last_fewshot_golden
+ meta["fewshot_golden_reuse"] = True
+ meta["fewshot_golden_qid"] = gq
+ meta["fewshot_golden_score"] = gsc
result = GenerationResult(
sql=sql,
@@ -792,6 +965,11 @@ class Text2SQLOrchestrator:
fail_meta["db_execution_status"] = last_db_execution_status
if dc_raw:
fail_meta["dialog_context_chars"] = len(dc_raw)
+ if self._last_fewshot_golden:
+ gq, gsc = self._last_fewshot_golden
+ fail_meta["fewshot_golden_reuse"] = True
+ fail_meta["fewshot_golden_qid"] = gq
+ fail_meta["fewshot_golden_score"] = gsc
return GenerationResult(
sql=last_sql or "",
valid=False,
diff --git a/backend/config/__pycache__/prompts.cpython-312.pyc b/backend/config/__pycache__/prompts.cpython-312.pyc
index 8a4418c..4eae589 100644
Binary files a/backend/config/__pycache__/prompts.cpython-312.pyc and b/backend/config/__pycache__/prompts.cpython-312.pyc differ
diff --git a/backend/config/prompts.py b/backend/config/prompts.py
index 77e6982..9f28354 100644
--- a/backend/config/prompts.py
+++ b/backend/config/prompts.py
@@ -225,6 +225,38 @@ SQL_GENERATOR_USER = """Schema信息:
请生成**有用 SQL**(见系统提示定义):必须与「业务级黄金范例」**同构**——大写关键字、多行缩进版式、PascalCase 别名、该展示对手方/账户等名称时须 LEFT JOIN 维表;禁止输出挤成一行的「极简 SQL」。"""
+# ========== Few-shot 黄金 SQL 条件适配(Chroma 库内为已校验正确答案)==========
+GOLDEN_SQL_ADAPT_SYSTEM = """你是精通 Microsoft SQL Server (T-SQL) 的数据库专家。
+
+**任务背景**:下方「标准答案 SQL」来自向量库中**已校验通过**的业务范例,结构与写法正确。当前用户问题与范例问题**语义高度相似**,仅时间、账户、状态、代码、筛选口径等「特殊条件」可能不同。
+
+**你必须遵守**:
+1. **以标准答案为主干**:优先保留其 `FROM`/`JOIN`/`ON`、主 `SELECT` 列清单与聚合/分组逻辑;**不要随意更换主表、不要拆掉必要 JOIN**,除非当前 Schema 片段中已不存在该表(此时在 Schema 内做最小替换并说明等价关系仅在脑中完成)。
+2. **只改「条件类」内容**:重点调整 `WHERE`/`HAVING`/`ORDER BY`/`TOP` 中的字面量、日期区间、状态码、账户/合约/代码等过滤;将用户问题中的时间范围、业务对象、筛选口径反映到这些条件中。
+3. **Schema 绝对优先**:表名、列名必须来自下方「当前 Schema 片段」;禁止臆造字段。若标准答案中某列在片段中不存在,按片段改写为合法列。
+4. **T-SQL 与版式**:与常规生成一致——关键字大写、多行缩进、`WHERE` 续行以 `AND` 开头、需要时 PascalCase 英文别名;禁止 MySQL 反引号与 `CURDATE()` 等。
+5. **禁止在字符串字面量中写中日韩文字**去匹配代码列;须用 Schema 注释中的代码或 JOIN 维表(与系统提示 SQL_GENERATOR 一致)。
+
+**输出**:仅输出一条完整可执行 SQL,不要解释。"""
+
+GOLDEN_SQL_ADAPT_USER = """【标准答案 SQL】(向量库中的已校验正确答案;相似度 {golden_score:.4f},越高越应保留结构)
+对应范例问题:{golden_question}
+
+```sql
+{golden_sql}
+```
+
+【当前 Schema 片段】(标识符必须与此一致)
+{schema}
+
+【当前用户问题】
+{question}
+
+【数据库方言】{dialect}
+
+请输出一条完整 T-SQL:**在标准答案基础上仅调整条件/字面量/排序等与当前问题相关的部分**,保留正确的 JOIN 与整体查询意图;版式与 SQL_GENERATOR 黄金范例同构。"""
+
+
# ========== Validator Agent Prompt ==========
VALIDATOR_SYSTEM = """你是一个严谨的SQL审核员,负责验证SQL语句的正确性和安全性。
diff --git a/backend/llm/__pycache__/deepseek_client.cpython-312.pyc b/backend/llm/__pycache__/deepseek_client.cpython-312.pyc
index ed78f27..0936ef0 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/llm/__pycache__/openai_client.cpython-312.pyc b/backend/llm/__pycache__/openai_client.cpython-312.pyc
index 91de5d6..46706f4 100644
Binary files a/backend/llm/__pycache__/openai_client.cpython-312.pyc and b/backend/llm/__pycache__/openai_client.cpython-312.pyc differ
diff --git a/backend/llm/deepseek_client.py b/backend/llm/deepseek_client.py
index 920b700..205abd5 100644
--- a/backend/llm/deepseek_client.py
+++ b/backend/llm/deepseek_client.py
@@ -6,7 +6,7 @@ DeepSeek API 客户端封装
import os
import json
import logging
-from typing import Dict, List, Optional, Any, Union
+from typing import Any, Dict, Iterator, List, Optional, Union
from dataclasses import dataclass, field
from openai import OpenAI, AsyncOpenAI
from openai.types.chat import ChatCompletion, ChatCompletionMessage
@@ -104,6 +104,43 @@ class DeepSeekClient:
logger.error(f"DeepSeek API调用失败: {e}")
raise
+ def chat_stream(
+ self,
+ messages: List[Dict[str, str]],
+ **kwargs: Any,
+ ) -> Iterator[str]:
+ """
+ 流式聊天:按 completion 增量产出文本片段(与 chat 同参,固定 stream=True)。
+ """
+ params: Dict[str, Any] = {
+ "model": self.config.model_name,
+ "messages": messages,
+ "temperature": kwargs.get("temperature", self.config.temperature),
+ "max_tokens": kwargs.get("max_tokens", self.config.max_tokens),
+ "top_p": kwargs.get("top_p", self.config.top_p),
+ "frequency_penalty": kwargs.get(
+ "frequency_penalty", self.config.frequency_penalty
+ ),
+ "presence_penalty": kwargs.get(
+ "presence_penalty", self.config.presence_penalty
+ ),
+ "stream": True,
+ "timeout": kwargs.get("timeout", self.config.timeout),
+ }
+ if self.config.extra_headers:
+ params["extra_headers"] = self.config.extra_headers
+ try:
+ stream = self.client.chat.completions.create(**params)
+ for chunk in stream:
+ if not chunk.choices:
+ continue
+ delta = chunk.choices[0].delta
+ if delta and getattr(delta, "content", None):
+ yield delta.content
+ except Exception as e:
+ logger.error(f"DeepSeek API流式调用失败: {e}")
+ raise
+
def chat_with_json(
self,
messages: List[Dict[str, str]],
diff --git a/backend/llm/openai_client.py b/backend/llm/openai_client.py
index a2a94dd..62004e3 100644
--- a/backend/llm/openai_client.py
+++ b/backend/llm/openai_client.py
@@ -10,7 +10,7 @@ import json
import logging
import os
from dataclasses import dataclass
-from typing import Any, Dict, List, Optional
+from typing import Any, Dict, Iterator, List, Optional
from openai import AsyncOpenAI, OpenAI # type: ignore[import-not-found]
from openai.types.chat import ( # type: ignore[import-not-found]
@@ -119,6 +119,69 @@ class OpenAIClient:
logger.error("OpenAI API调用失败: %s", e)
raise
+ def chat_stream(
+ self,
+ messages: List[Dict[str, str]],
+ **kwargs: Any,
+ ) -> Iterator[str]:
+ """流式聊天:按 completion 增量产出文本片段(与 chat 同参,固定 stream=True)。"""
+ max_tokens = kwargs.get("max_tokens", self.config.max_tokens)
+ max_completion_tokens = kwargs.get("max_completion_tokens", None)
+
+ params: Dict[str, Any] = {
+ "model": self.config.model_name,
+ "messages": messages,
+ "temperature": kwargs.get("temperature", self.config.temperature),
+ "top_p": kwargs.get("top_p", self.config.top_p),
+ "frequency_penalty": kwargs.get(
+ "frequency_penalty", self.config.frequency_penalty
+ ),
+ "presence_penalty": kwargs.get(
+ "presence_penalty", self.config.presence_penalty
+ ),
+ "stream": True,
+ "timeout": kwargs.get("timeout", self.config.timeout),
+ }
+ if max_completion_tokens is not None:
+ params["max_completion_tokens"] = max_completion_tokens
+ else:
+ params["max_tokens"] = max_tokens
+
+ if self.config.extra_headers:
+ params["extra_headers"] = self.config.extra_headers
+
+ try:
+ stream = self.client.chat.completions.create(**params)
+ for chunk in stream:
+ if not chunk.choices:
+ continue
+ delta = chunk.choices[0].delta
+ if delta and getattr(delta, "content", None):
+ yield delta.content
+ except Exception as e:
+ msg = str(e)
+ if (
+ "Unsupported parameter" in msg
+ and "max_tokens" in msg
+ and "max_completion_tokens" in msg
+ and "max_completion_tokens" not in params
+ ):
+ params.pop("max_tokens", None)
+ params["max_completion_tokens"] = max_tokens
+ try:
+ stream = self.client.chat.completions.create(**params)
+ for chunk in stream:
+ if not chunk.choices:
+ continue
+ delta = chunk.choices[0].delta
+ if delta and getattr(delta, "content", None):
+ yield delta.content
+ return
+ except Exception:
+ pass
+ logger.error("OpenAI API流式调用失败: %s", e)
+ raise
+
def chat_with_json(
self, messages: List[Dict[str, str]], **kwargs
) -> Dict[str, Any]:
diff --git a/backend/utils/__pycache__/fewshot_chroma_store.cpython-312.pyc b/backend/utils/__pycache__/fewshot_chroma_store.cpython-312.pyc
index 0183183..bfe9137 100644
Binary files a/backend/utils/__pycache__/fewshot_chroma_store.cpython-312.pyc and b/backend/utils/__pycache__/fewshot_chroma_store.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 99f9457..28f8ee2 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/fewshot_selector.py b/backend/utils/fewshot_selector.py
index 61afe86..2a77b23 100644
--- a/backend/utils/fewshot_selector.py
+++ b/backend/utils/fewshot_selector.py
@@ -19,7 +19,7 @@ Few-shot示例选择器 - 基于经验数据集动态选择相关示例
import json
import os
from pathlib import Path
-from typing import List, Dict, Optional
+from typing import List, Dict, Optional, Tuple
from dataclasses import dataclass
import numpy as np
import logging
@@ -158,6 +158,92 @@ class FewShotSelector:
self.cache_path = None
logger.info("[OK] Few-shot 使用 Chroma(%s 条)", self._chroma_store.count())
+ @property
+ def is_chroma_backend(self) -> bool:
+ """是否使用 Chroma 持久化库(`data/embeddings/chroma_fewshot` 等)。"""
+ return self._chroma_store is not None
+
+ def _passes_filters(
+ self,
+ sample: "ExperienceSample",
+ min_rating: Optional[int],
+ required_tags: Optional[List[str]],
+ max_difficulty: str,
+ exclude_qids: Optional[List[str]],
+ ) -> bool:
+ if exclude_qids and sample.qid in exclude_qids:
+ return False
+ if min_rating and sample.rating is not None and sample.rating < min_rating:
+ return False
+ if max_difficulty == "easy" and sample.difficulty != "easy":
+ return False
+ if max_difficulty == "medium" and sample.difficulty == "hard":
+ return False
+ if required_tags and not all(tag in sample.tags for tag in required_tags):
+ return False
+ return True
+
+ def select_best_with_score(
+ self,
+ question: str,
+ min_rating: Optional[int] = None,
+ required_tags: Optional[List[str]] = None,
+ max_difficulty: str = "hard",
+ exclude_qids: Optional[List[str]] = None,
+ ) -> Optional[Tuple[ExperienceSample, float]]:
+ """
+ 返回通过筛选的**相似度最高**一条样本及分数 ``[0,1]``(与向量余弦一致:1 - distance)。
+ 无命中时返回 ``None``。
+ """
+ if self._chroma_store is not None:
+ from utils.fewshot_chroma_store import sample_from_chroma_metadata
+
+ over_fetch = 96
+ rows = self._chroma_store.search_raw(question, top_k=over_fetch)
+ for score, meta, doc in rows:
+ s = sample_from_chroma_metadata(meta, doc)
+ if not self._passes_filters(
+ s, min_rating, required_tags, max_difficulty, exclude_qids
+ ):
+ continue
+ logger.info(
+ "[Few-shot] best_with_score: qid=%s score=%.4f preview=%r",
+ s.qid,
+ float(score),
+ (s.question_zh or "")[:100],
+ )
+ return (s, float(score))
+ return None
+
+ if self._embedder is None or self.embeddings is None or len(self.samples) == 0:
+ return None
+
+ q_emb = self._embedder.encode(
+ [question],
+ batch_size=1,
+ normalize=True,
+ show_progress=False,
+ )[0]
+ scores = np.dot(self.embeddings, q_emb)
+ best: Optional[Tuple[float, ExperienceSample]] = None
+ for idx, (score, sample) in enumerate(zip(scores, self.samples)):
+ if not self._passes_filters(
+ sample, min_rating, required_tags, max_difficulty, exclude_qids
+ ):
+ continue
+ s = float(score)
+ if best is None or s > best[0]:
+ best = (s, sample)
+ if best is None:
+ return None
+ logger.info(
+ "[Few-shot] best_with_score(numpy): qid=%s score=%.4f preview=%r",
+ best[1].qid,
+ best[0],
+ (best[1].question_zh or "")[:100],
+ )
+ return (best[1], best[0])
+
def _load_samples(self):
"""从 JSONL 加载样本到内存(非 Chroma 模式必需;Chroma 空库时用于首次灌库)。"""
if self.samples_path is None:
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 d33cd2c..824aa08 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/text2sql_api.log b/logs/text2sql_api.log
index d59031f..b4d1ded 100644
--- a/logs/text2sql_api.log
+++ b/logs/text2sql_api.log
@@ -1,67 +1 @@
-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;"
+2026-04-16 13:35:24 INFO [utils.repo_logging] repo_logging.py:79 configure_text2sql_api_logging() | 日志文件: C:\Users\24019\Desktop\backman-camel\logs\text2sql_api.log