update by wensicheng on 0629

This commit is contained in:
tao.chen
2026-07-01 17:46:08 +08:00
parent 0fc1c702d9
commit 2f406bd64d
12 changed files with 1012 additions and 1574 deletions
+154 -90
View File
@@ -1,124 +1,188 @@
---
name: pyspark-sql-pipeline
description: 需要把业务需求一路编排到"已被评审过、可执行的 SQL 字符串"时使用。触发词包括"完整跑一遍取数流程"、"生成 SQL 并 review"、"end-to-end 取数"、"full pipeline"、"从头到尾"、"数据需求到 SQL"。把 requirements-analysismetadata-validatorlogic-plannersql-context-builderpyspark-sql-guardrailssql-review 串成一条确定性流水线,阶段之间用显式的用户确认关口隔开。**不执行 SQL** — 只把流水线驱动到一条被评审过的 SQL 字符串
description: 用户希望从业务需求端到端生成一条经过校验和评审的 Spark SQL / PySpark SQL 作业时使用。它编排 requirements-analysismetadata-validatorlogic-plannersql-context-builder、SQL 生成、pyspark-sql-guardrailssql-review,并以正式交付格式输出 DRD、元数据校验、逻辑计划、SQL 上下文、最终 SQL、review 和交付清单
---
# PySpark SQL Pipeline(PySpark SQL 流水线)
# PySpark SQL Pipeline(端到端 SQL 流水线
## 概述(Overview)
## 目标
一份业务数据需求一路编排到"已评审"SQL 字符串的端到端协调器。每个阶段有固定的输入/输出契约;任意两阶段之间,当上游状态不是绿色时,都设有**用户确认关口**
“用户一句业务需求”推进到“可交付、可审查、可追溯的 SQL 和配套审查结果”。它不是单个 SQL 生成器,而是一条带关口的工程流程
核心原则:**绝不跳阶段,也绝不让任何阶段悄无声息地越过它自己的 `NEED_USER_CONFIRMATION`。** 流水线是一连串交接,不是某个 agent 的内心独白
本流水线默认不执行 SQL、不提交 Spark 作业。只有用户明确要求执行,并且 SQL 已通过 guardrails 和 review 后,才允许进入 MCP 工具层
## 何时使用(When to Use)
## 阶段顺序
**使用场景:**
- 用户给出一个数据需求,希望端到端生成 SQL(例如"统计近 30 天各城市贷款总额及客户数")
- 用户要求"跑完整流水线"或"从需求到 SQL 一条龙"
- 一个非平凡的 SQL 改动会同时触碰多个 skill 的契约
**不要使用场景:**
- 用户只需要其中一个阶段(例如"评审这条 SQL" → 直接用 `sql-review`)
- 用户明确说"跳过校验"或"快点给我 SQL"— 主动说明风险并确认一次;若用户确认,可以合并阶段,但**仍必须**应用 `pyspark-sql-guardrails``sql-review`
## 流水线(The Pipeline)
```dot
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"];
}
```text
1. requirements-analysis
2. metadata-validator
3. logic-planner
4. sql-context-builder
5. SQL generation
6. pyspark-sql-guardrails
7. sql-review
8. optional MCP execution handoff
```
每条箭头都是**契约边界**。前序阶段没产出预期状态时,绝不允许下游阶段启动。
## 每阶段放行条件
## 阶段契约(Stage Contracts)
```yaml
requirements-analysis: requirements_output.status == READY_FOR_METADATA
metadata-validator: validation_result.status == VALIDATED
logic-planner: logic_plan.status == PLANNED
sql-context-builder: sql_context.status == CONTEXT_READY
pyspark-sql-guardrails: sql_guard.status == PASS
sql-review: sql_review.verdict == PASS 或用户明确接受非 HIGH 风险
mcp-handoff: user explicitly requested execution
```
| # | 阶段 | 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` |
任何阶段输出 `NEED_USER_CONFIRMATION``FAIL`,流水线必须暂停。
## 关口规则(Gate Rules,务必细读)
## SQL 生成规则
返回停止状态的阶段**不能**被下一阶段跟随。要做的:
SQL 生成没有单独 skill。生成 SQL 时只能翻译以下内容:
1. 把该阶段的 `pending_questions` 暴露给用户
2. 停下流水线
3. 等用户回答后,从**同一个阶段**重跑(不是从头跑),把用户的新回答作为额外输入
4. 一旦该阶段产出绿色状态,推进到下一阶段
- `requirements_output`
- `validation_result`
- `logic_plan.steps`
- `sql_context.aliases`
- `sql_context.fields`
- `sql_context.joins`
- `sql_context.metrics`
- `sql_context.filters`
- `sql_context.time_windows`
- `sql_context.outputs`
- `sql_context.write_intent`
阶段是幂等的。同一份输入重跑必须产出同样的输出;新输入则必须产出能覆盖旧输出的新输出。所有前置的 `field_mapping``validation_result``logic_plan``sql_context` 必须作为流水线的持久化产物一路带下来
不能临时新增表、字段、alias、join key 或指标公式。每个 `alias.column` 必须能在 `sql_context.fields` 中找到
## SQL 生成(阶段 5)
## SQL 时间口径硬规则
本步没有独立的 skill。流水线直接执行,约束如下:
对“近 N 天按 T-1 完整日”,Spark/Hive SQL 必须使用左闭右开窗口:
- `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` 中不修改地跑过
```sql
event_time >= date_sub(current_date(), N)
AND event_time < current_date()
```
## 用户确认消息(User-Confirmation Messages)
禁止用下面写法表示 timestamp 完整日窗口:
每个关口应产出一条结构化消息,格式如下:
```sql
event_time BETWEEN date_add(current_date(), -N) AND date_add(current_date(), -1)
```
如果用户指定自然日、滚动到当前时刻、特定日期或特定时区,必须在 DRD 和 SQL 注释里写明。
## 转化率 SQL 硬规则
漏斗/转化率 SQL 必须保证事件顺序。对“注册后完成付费”,分子应判断存在满足条件的后续付费:
```sql
COUNT(DISTINCT CASE
WHEN ur.user_id IS NOT NULL
AND qp.user_id IS NOT NULL
THEN fl.user_id
END) AS converted_user_cnt
```
其中 `qualified_payment` 这类 CTE 必须在 join 到注册事件后筛选:
```sql
po.pay_status = 'SUCCESS'
AND po.pay_time >= ur.register_time
```
不要先对 `pay_order` 按用户取全局最早支付再和 `register_time` 比较;这会漏掉“历史早付费但注册后也完成付费”的用户。
## 推荐 CTE 结构
对新用户转化率这类需求,优先使用清晰 CTE:
```text
first_login 分母人群:窗口内首次登录用户
registered 注册事件:按用户保留注册时间
qualified_payment 合格付费:注册后成功付费
final_agg 按维度聚合输出
```
## 正式交付格式
流水线完成时,最终回答必须包含以下章节,不能只给 SQL:
```markdown
## 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).
## 需求说明
## 元数据校验
## 逻辑计划
## SQL 上下文
## 最终 SQL
## SQL Review
## 交付清单
```
不要把某个阶段的输出当作最终结果,如果它的状态是停止状态。不要把 `pending_questions` 埋在冗长叙述里。
每章要求:
## 失败处理(Failure Handling)
- `需求说明`:业务目标、统计口径、涉及表、表关系、字段清单。
- `元数据校验`:几张表、多少字段、join key、时间字段是否通过。
- `逻辑计划`:用表格列出 CTE/步骤、说明和关键条件。
- `SQL 上下文`:alias、字段来源、join、指标、时间窗口。
- `最终 SQL`:只在 guardrails PASS 且 review 非 FAIL 后展示。
- `SQL Review`:verdict、最高风险等级、风险条数、关键建议。
- `交付清单`SQL 文件、PySpark 脚本、是否已执行、如何运行。
- **防护栏失败(`pyspark-sql-guardrails`)** — 不要用同一条 SQL 重试。回到阶段 5 重新生成,这一次遵守防护栏
- **评审结论为 FAIL 且有 HIGH 风险** — 回到阶段 5(或更早,如果是 context 类的发现,就回到阶段 4)。**不要**"打个补丁"让发现项闭嘴;要解决根因
- **评审结论为 FAIL 但只有 MEDIUM/LOW** — 把发现项暴露给用户,由用户决定上线还是改
- **阶段间不一致**(例如 SQL 引用了 `sql_context` 里没有的列) — 视为阶段 5 的 bug;SQL 在生成时没遵守 context。**重新生成,不要修补**。
如果生成 PySpark 脚本但当前容器没有 Java/Spark,要明确:
## 禁止行为(Forbidden Behaviors)
```text
已生成脚本,但当前 opencode 容器不是 Spark 执行环境,未本地执行。可在 Spark 集群中使用 spark-submit 运行。
```
| 自我说服 | 现实 |
|---|---|
| "用户要 SQL 快点,关口跳过" | 关口存在的原因就是跳过它们会让数字变错。先问一次;如果用户确认,可以只合并相邻的非校验阶段。**永远不能**跳过防护栏或评审。 |
| "只一个阶段失败,不必问" | 每一个停止状态都是要问用户的问题。永远要问。 |
| "我把校验内联到 SQL 生成里" | 阶段分开是为了能独立重跑,内联会把它们耦合死。 |
| "评审只发现 MEDIUM,直接上线" | 把发现项暴露给用户,决定权在用户,不在 agent。 |
| "字段映射太显然,跳过校验器" | 字段映射就是校验器的工作。绕过它就是幻觉字段名的头号原因。 |
| "状态一样,直接继续" | 状态名是固定字符串。重读实际状态文本,不要凭记忆做模式匹配。 |
## 关口消息
## 完成判定(Completion Criteria)
当阶段暂停时,必须用这个格式:
流水线完成需要**全部**满足:
```markdown
## Pipeline paused at: <stage>
**Status:** NEED_USER_CONFIRMATION
**Open questions:**
1. <question>
**To resume:** answer the questions above, then say "继续".
```
- 7 个阶段按序都跑过
- `sql-review.verdict``PASS`(或用户已显式接受剩余的 MEDIUM/LOW 发现)
- 同一份回复里同时向用户交付:最终 SQL 字符串、`sql_context``logic_plan``validation_result``requirements_output`
## 输出格式
如果用户后续修改了需求,从阶段 1 重新开始 — **不要**打补丁修下游产物。
```yaml
pipeline_result:
status: COMPLETE | PAUSED | FAILED
current_stage: requirements-analysis | metadata-validator | logic-planner | sql-context-builder | sql-generation | pyspark-sql-guardrails | sql-review | mcp-handoff
requirements_output: {}
validation_result: {}
logic_plan: {}
sql_context: {}
sql: "只在 guardrails PASS 后展示"
pyspark_code: "只在 SQL 通过校验和评审后生成"
sql_guard: {}
sql_review: {}
mcp_plan:
should_execute: false
allowed_tools: []
pending_questions: []
blocking_findings: []
```
## 失败处理
- 需求变了:回到 requirements-analysis。
- 缺表缺字段:回到 metadata-validator。
- 指标公式或 join 粒度错:回到 logic-planner。
- SQL 引用了 context 没有的字段:回到 sql-context-builder 或更上游。
- guardrails 失败:重新生成 SQL,不改弱校验器。
- review 有 HIGH:阻断执行,修根因。
- MCP 准备失败:只修执行参数或脚本落盘问题,不反向篡改已评审 SQL。
## 禁止行为
- 不允许跳过 metadata-validator。
- 不允许 guardrails PASS 后跳过 review。
- 不允许 review 失败后只做字符串补丁而不检查上游产物。
- 不允许在用户未确认时继续下一阶段。
- 不允许把 `prepare_submit_job` 当成真正执行成功。
- 不允许在未获用户明确确认时调用 `confirm_submit_job`