6.9 KiB
6.9 KiB
name, description
| name | description |
|---|---|
| pyspark-sql-pipeline | 当需要把业务需求一路编排到"已被评审过、可执行的 SQL 字符串"时使用。触发词包括"完整跑一遍取数流程"、"生成 SQL 并 review"、"end-to-end 取数"、"full pipeline"、"从头到尾"、"数据需求到 SQL"。把 requirements-analysis → metadata-validator → logic-planner → sql-context-builder → pyspark-sql-guardrails → sql-review 串成一条确定性流水线,阶段之间用显式的用户确认关口隔开。**不执行 SQL** — 只把流水线驱动到一条被评审过的 SQL 字符串。 |
PySpark SQL Pipeline(PySpark SQL 流水线)
概述(Overview)
把一份业务数据需求一路编排到"已评审"SQL 字符串的端到端协调器。每个阶段有固定的输入/输出契约;任意两阶段之间,当上游状态不是绿色时,都设有用户确认关口。
核心原则:绝不跳阶段,也绝不让任何阶段悄无声息地越过它自己的 NEED_USER_CONFIRMATION。 流水线是一连串交接,不是某个 agent 的内心独白。
何时使用(When to Use)
使用场景:
- 用户给出一个数据需求,希望端到端生成 SQL(例如"统计近 30 天各城市贷款总额及客户数")
- 用户要求"跑完整流水线"或"从需求到 SQL 一条龙"
- 一个非平凡的 SQL 改动会同时触碰多个 skill 的契约
不要使用场景:
- 用户只需要其中一个阶段(例如"评审这条 SQL" → 直接用
sql-review) - 用户明确说"跳过校验"或"快点给我 SQL"— 主动说明风险并确认一次;若用户确认,可以合并阶段,但仍必须应用
pyspark-sql-guardrails和sql-review
流水线(The Pipeline)
digraph pipeline {
"1. requirements-analysis" [shape=box];
"2. metadata-validator" [shape=box];
"3. logic-planner" [shape=box];
"4. sql-context-builder" [shape=box];
"5. SQL generation" [shape=box];
"6. pyspark-sql-guardrails" [shape=box];
"7. sql-review" [shape=box];
"1. requirements-analysis" -> "2. metadata-validator" [label="status=READY_FOR_VALIDATION"];
"2. metadata-validator" -> "3. logic-planner" [label="status=VALIDATED"];
"3. logic-planner" -> "4. sql-context-builder" [label="status=PLANNED"];
"4. sql-context-builder" -> "5. SQL generation" [label="status=READY"];
"5. SQL generation" -> "6. pyspark-sql-guardrails" [label="SQL 字符串"];
"6. pyspark-sql-guardrails" -> "7. sql-review" [label="guardrail pass"];
}
每条箭头都是契约边界。前序阶段没产出预期状态时,绝不允许下游阶段启动。
阶段契约(Stage Contracts)
| # | 阶段 | Skill | 输入 | 绿色输出状态 | 停止状态 |
|---|---|---|---|---|---|
| 1 | 需求 | requirements-analysis |
用户消息 | READY_FOR_VALIDATION |
NEED_USER_CONFIRMATION |
| 2 | 元数据 | metadata-validator |
requirements_output + 工作区 | VALIDATED |
NEED_USER_CONFIRMATION |
| 3 | 逻辑 | logic-planner |
validation_result | PLANNED |
NEED_USER_CONFIRMATION |
| 4 | 上下文 | sql-context-builder |
logic_plan + field_mapping | READY |
NEED_USER_CONFIRMATION |
| 5 | 生成 | (无独立 skill — 见下文"SQL 生成") | sql_context | SQL 字符串 | n/a |
| 6 | 防护栏 | pyspark-sql-guardrails |
SQL 字符串 | pass |
fail(抛错) |
| 7 | 评审 | sql-review |
SQL 字符串 + sql_context | verdict: PASS |
verdict: FAIL |
关口规则(Gate Rules,务必细读)
返回停止状态的阶段不能被下一阶段跟随。要做的:
- 把该阶段的
pending_questions暴露给用户 - 停下流水线
- 等用户回答后,从同一个阶段重跑(不是从头跑),把用户的新回答作为额外输入
- 一旦该阶段产出绿色状态,推进到下一阶段
阶段是幂等的。同一份输入重跑必须产出同样的输出;新输入则必须产出能覆盖旧输出的新输出。所有前置的 field_mapping、validation_result、logic_plan、sql_context 必须作为流水线的持久化产物一路带下来。
SQL 生成(阶段 5)
本步没有独立的 skill。流水线直接执行,约束如下:
sql_context是字段名、alias、join key 的唯一来源- 按
logic_plan.steps的顺序组装 SQL:source→filter→join→transform→dedupe→window→aggregate→output - 输出中每个
<alias>.<column>都必须出现在sql_context.fields中 - 生成的 SQL 必须能在
pyspark-sql-guardrails.safe_spark_sql中不修改地跑过
用户确认消息(User-Confirmation Messages)
每个关口应产出一条结构化消息,格式如下:
## Pipeline paused at: <stage name>
**Status:** NEED_USER_CONFIRMATION
**Open questions:**
1. <question 1>
2. <question 2>
**Pending options:** A / B / C / other
**To resume:** answer the questions above, then say "继续" (continue).
不要把某个阶段的输出当作最终结果,如果它的状态是停止状态。不要把 pending_questions 埋在冗长叙述里。
失败处理(Failure Handling)
- 防护栏失败(
pyspark-sql-guardrails) — 不要用同一条 SQL 重试。回到阶段 5 重新生成,这一次遵守防护栏 - 评审结论为 FAIL 且有 HIGH 风险 — 回到阶段 5(或更早,如果是 context 类的发现,就回到阶段 4)。不要"打个补丁"让发现项闭嘴;要解决根因
- 评审结论为 FAIL 但只有 MEDIUM/LOW — 把发现项暴露给用户,由用户决定上线还是改
- 阶段间不一致(例如 SQL 引用了
sql_context里没有的列) — 视为阶段 5 的 bug;SQL 在生成时没遵守 context。重新生成,不要修补。
禁止行为(Forbidden Behaviors)
| 自我说服 | 现实 |
|---|---|
| "用户要 SQL 快点,关口跳过" | 关口存在的原因就是跳过它们会让数字变错。先问一次;如果用户确认,可以只合并相邻的非校验阶段。永远不能跳过防护栏或评审。 |
| "只一个阶段失败,不必问" | 每一个停止状态都是要问用户的问题。永远要问。 |
| "我把校验内联到 SQL 生成里" | 阶段分开是为了能独立重跑,内联会把它们耦合死。 |
| "评审只发现 MEDIUM,直接上线" | 把发现项暴露给用户,决定权在用户,不在 agent。 |
| "字段映射太显然,跳过校验器" | 字段映射就是校验器的工作。绕过它就是幻觉字段名的头号原因。 |
| "状态一样,直接继续" | 状态名是固定字符串。重读实际状态文本,不要凭记忆做模式匹配。 |
完成判定(Completion Criteria)
流水线完成需要全部满足:
- 7 个阶段按序都跑过
sql-review.verdict为PASS(或用户已显式接受剩余的 MEDIUM/LOW 发现)- 同一份回复里同时向用户交付:最终 SQL 字符串、
sql_context、logic_plan、validation_result、requirements_output
如果用户后续修改了需求,从阶段 1 重新开始 — 不要打补丁修下游产物。