6.1 KiB
6.1 KiB
name, description
| name | description |
|---|---|
| pyspark-sql-pipeline | 当用户希望从业务需求端到端生成一条经过校验和评审的 Spark SQL / PySpark SQL 作业时使用。它编排 requirements-analysis、metadata-validator、logic-planner、sql-context-builder、SQL 生成、pyspark-sql-guardrails、sql-review,并以正式交付格式输出 DRD、元数据校验、逻辑计划、SQL 上下文、最终 SQL、review 和交付清单。 |
PySpark SQL Pipeline(端到端 SQL 流水线)
目标
把“用户一句业务需求”推进到“可交付、可审查、可追溯的 SQL 和配套审查结果”。它不是单个 SQL 生成器,而是一条带关口的工程流程。
本流水线默认不执行 SQL、不提交 Spark 作业。只有用户明确要求执行,并且 SQL 已通过 guardrails 和 review 后,才允许进入 MCP 工具层。
阶段顺序
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
每阶段放行条件
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
任何阶段输出 NEED_USER_CONFIRMATION 或 FAIL,流水线必须暂停。
SQL 生成规则
SQL 生成没有单独 skill。生成 SQL 时只能翻译以下内容:
requirements_outputvalidation_resultlogic_plan.stepssql_context.aliasessql_context.fieldssql_context.joinssql_context.metricssql_context.filterssql_context.time_windowssql_context.outputssql_context.write_intent
不能临时新增表、字段、alias、join key 或指标公式。每个 alias.column 必须能在 sql_context.fields 中找到。
SQL 时间口径硬规则
对“近 N 天按 T-1 完整日”,Spark/Hive SQL 必须使用左闭右开窗口:
event_time >= date_sub(current_date(), N)
AND event_time < current_date()
禁止用下面写法表示 timestamp 完整日窗口:
event_time BETWEEN date_add(current_date(), -N) AND date_add(current_date(), -1)
如果用户指定自然日、滚动到当前时刻、特定日期或特定时区,必须在 DRD 和 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 到注册事件后筛选:
po.pay_status = 'SUCCESS'
AND po.pay_time >= ur.register_time
不要先对 pay_order 按用户取全局最早支付再和 register_time 比较;这会漏掉“历史早付费但注册后也完成付费”的用户。
推荐 CTE 结构
对新用户转化率这类需求,优先使用清晰 CTE:
first_login 分母人群:窗口内首次登录用户
registered 注册事件:按用户保留注册时间
qualified_payment 合格付费:注册后成功付费
final_agg 按维度聚合输出
正式交付格式
流水线完成时,最终回答必须包含以下章节,不能只给 SQL:
## 需求说明
## 元数据校验
## 逻辑计划
## SQL 上下文
## 最终 SQL
## SQL Review
## 交付清单
每章要求:
需求说明:业务目标、统计口径、涉及表、表关系、字段清单。元数据校验:几张表、多少字段、join key、时间字段是否通过。逻辑计划:用表格列出 CTE/步骤、说明和关键条件。SQL 上下文:alias、字段来源、join、指标、时间窗口。最终 SQL:只在 guardrails PASS 且 review 非 FAIL 后展示。SQL Review:verdict、最高风险等级、风险条数、关键建议。交付清单:SQL 文件、PySpark 脚本、是否已执行、如何运行。
如果生成 PySpark 脚本但当前容器没有 Java/Spark,要明确:
已生成脚本,但当前 opencode 容器不是 Spark 执行环境,未本地执行。可在 Spark 集群中使用 spark-submit 运行。
关口消息
当阶段暂停时,必须用这个格式:
## Pipeline paused at: <stage>
**Status:** NEED_USER_CONFIRMATION
**Open questions:**
1. <question>
**To resume:** answer the questions above, then say "继续".
输出格式
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。