数据截至 (上游 commit 69155611c59f)
03 · 主管线 run_query
本章讲主循环
run_query(api/core/text2sql.py:308):把 01/02 章的部件串成一条端到端的 Text2SQL 流水线,并解释一个漂亮的设计——同一条生成器同时喂流式前端和无服务器 SDK。
3.1 一个生成器,两种消费者
它要解决的小问题: 网页要流式看到每一步进度(「正在生成 SQL…」「正在执行…」),而 Python SDK 只想要一个最终结果对象。若为两者各写一套管线,逻辑会漂移。
思路: 让 run_query 做成 async 生成器:
- 每一步
yield一个面向前端的 wire dict({"type": "reasoning_step", "message": ...}这类); - 最后
yield一个_Final(QueryResult)哨兵(api/core/text2sql.py:203),里面是结构化结果。
两种消费者各取所需:
流式路由 _serialize_pipeline (api/routes/graphs.py:37)
每个 dict → json.dumps + MESSAGE_DELIMITER 推给前端;遇 _Final 停
SDK collect_result (api/core/text2sql.py:209)
丢掉所有 dict,只 return _Final.value(QueryResult)
MESSAGE_DELIMITER = "|||FALKORDB_MESSAGE_BOUNDARY|||"(api/core/pipeline.py:33)是前后端约定的切帧符,前端按它把粘在一起的 JSON 拆开。
3.2 主线各步(按执行顺序)
run_query(user_id, graph_id, chat_data)
│
1. 校验&截断对话历史(只留最近 5 轮 SHORT_MEMORY_LENGTH)
│ 并 发起 MemoryTool 创建任务(若 use_memory)
│
2. 取库描述 db_description、库 URL、用户规则;定 db_type/loader
│
3. 并发:find(找表) ‖ RelevancyAgent(相关性门禁)
│ ├─ off-topic → 取消 find,发 followup 事件,收尾(is_valid=False)
│ └─ on-topic → 等 find 拿到 tables;若开记忆则搜记忆上下文
│
4. AnalysisAgent.get_analysis → 一条 SQL + confidence + 缺失/歧义
│ └─ is_sql_translatable=False → FollowUpAgent 反问用户,收尾
│
5. 自动给带特殊字符的表名加引号(auto_quote_sql_identifiers)
│
6. detect_destructive_operation:
│ ├─ demo 库上破坏性 → 直接拒
│ └─ 破坏性 → 发确认事件,收尾(requires_confirmation=True)
│
7. 执行 SQL
│ └─ 抛错 → HealerAgent 自愈(见 04 章),仍失败则抛
│
8. schema 变更?→ 刷新图(_emit_schema_refresh)
│
9. ResponseFormatterAgent → 人话回答
│
10. 后台异步存记忆(save_memory_background)
│
└─ yield _Final(QueryResult)