diff --git a/skills/bash-linux/SKILL.md b/skills/bash-linux/SKILL.md index 456b77a..146e172 100644 --- a/skills/bash-linux/SKILL.md +++ b/skills/bash-linux/SKILL.md @@ -1,209 +1,88 @@ --- name: bash-linux -description: "Bash/Linux terminal patterns. Critical commands, piping, error handling, scripting. Use when working on macOS or Linux systems." -risk: unknown -source: community -date_added: "2026-02-27" +description: 当需要在 Linux/macOS Bash 环境或容器内查文件、搜文本、跑脚本、看日志、处理进程、使用管道、环境变量或编写小型 shell 命令时使用。它是辅助命令行规范,不属于 SQL 业务流水线。 --- -# Bash Linux Patterns +# Bash Linux(命令行规范) -> Essential patterns for Bash on Linux/macOS. +## 目标 ---- +让 agent 在 Bash/Linux/容器环境下用稳定、可解释、低风险的方式执行命令。先观察,再操作;先缩小范围,再修改。 -## 1. Operator Syntax +## 基本原则 -### Chaining Commands +- 查文件优先用 `rg --files`,不可用时用 `find`。 +- 搜内容优先用 `rg -n`。 +- 读文件用 `sed -n`、`head`、`tail` 控制输出量。 +- 脚本使用 `set -euo pipefail`。 +- 路径和变量要加引号。 +- 删除、移动、改权限前必须确认目标范围。 +- 长日志只摘关键部分,不整屏倾倒。 +- 容器内的 `127.0.0.1` 指向容器自己,不等于宿主机。 -| Operator | Meaning | Example | -|----------|---------|---------| -| `;` | Run sequentially | `cmd1; cmd2` | -| `&&` | Run if previous succeeded | `npm install && npm run dev` | -| `\|\|` | Run if previous failed | `npm test \|\| echo "Tests failed"` | -| `\|` | Pipe output | `ls \| grep ".js"` | - ---- - -## 2. File Operations - -### Essential Commands - -| Task | Command | -|------|---------| -| List all | `ls -la` | -| Find files | `find . -name "*.js" -type f` | -| File content | `cat file.txt` | -| First N lines | `head -n 20 file.txt` | -| Last N lines | `tail -n 20 file.txt` | -| Follow log | `tail -f log.txt` | -| Search in files | `grep -r "pattern" --include="*.js"` | -| File size | `du -sh *` | -| Disk usage | `df -h` | - ---- - -## 3. Process Management - -| Task | Command | -|------|---------| -| List processes | `ps aux` | -| Find by name | `ps aux \| grep node` | -| Kill by PID | `kill -9 ` | -| Find port user | `lsof -i :3000` | -| Kill port | `kill -9 $(lsof -t -i :3000)` | -| Background | `npm run dev &` | -| Jobs | `jobs -l` | -| Bring to front | `fg %1` | - ---- - -## 4. Text Processing - -### Core Tools - -| Tool | Purpose | Example | -|------|---------|---------| -| `grep` | Search | `grep -rn "TODO" src/` | -| `sed` | Replace | `sed -i 's/old/new/g' file.txt` | -| `awk` | Extract columns | `awk '{print $1}' file.txt` | -| `cut` | Cut fields | `cut -d',' -f1 data.csv` | -| `sort` | Sort lines | `sort -u file.txt` | -| `uniq` | Unique lines | `sort file.txt \| uniq -c` | -| `wc` | Count | `wc -l file.txt` | - ---- - -## 5. Environment Variables - -| Task | Command | -|------|---------| -| View all | `env` or `printenv` | -| View one | `echo $PATH` | -| Set temporary | `export VAR="value"` | -| Set in script | `VAR="value" command` | -| Add to PATH | `export PATH="$PATH:/new/path"` | - ---- - -## 6. Network - -| Task | Command | -|------|---------| -| Download | `curl -O https://example.com/file` | -| API request | `curl -X GET https://api.example.com` | -| POST JSON | `curl -X POST -H "Content-Type: application/json" -d '{"key":"value"}' URL` | -| Check port | `nc -zv localhost 3000` | -| Network info | `ifconfig` or `ip addr` | - ---- - -## 7. Script Template +## 常用命令 ```bash -#!/bin/bash -set -euo pipefail # Exit on error, undefined var, pipe fail +pwd +ls -la +find . -type f -name '*.py' +rg -n "keyword" . +sed -n '1,120p' file.txt +head -n 40 file.txt +tail -n 100 app.log +wc -l file.txt +``` -# Colors (optional) -RED='\033[0;31m' -GREEN='\033[0;32m' -NC='\033[0m' +## Docker/容器排查 -# Script directory -SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +```bash +docker ps +docker logs --tail 100 +docker exec -it sh +getent hosts host.docker.internal +curl -fsS http://host.docker.internal:8000/health +``` -# Functions -log_info() { echo -e "${GREEN}[INFO]${NC} $1"; } -log_error() { echo -e "${RED}[ERROR]${NC} $1" >&2; } +## 进程和端口 + +```bash +ps aux | grep '[p]ython' +lsof -i :8000 +kill +``` + +## 脚本模板 + +```bash +#!/usr/bin/env bash +set -euo pipefail -# Main main() { - log_info "Starting..." - # Your logic here - log_info "Done!" + if [ "$#" -lt 1 ]; then + echo "usage: $0 " >&2 + exit 2 + fi + echo "arg=$1" } main "$@" ``` ---- +## 输出报告 -## 8. Common Patterns - -### Check if command exists - -```bash -if command -v node &> /dev/null; then - echo "Node is installed" -fi +```yaml +command_report: + purpose: "为什么运行" + command: "命令摘要" + result: success | failed | skipped + important_output: "关键输出" + next_step: "下一步" ``` -### Default variable value +## 禁止行为 -```bash -NAME=${1:-"default_value"} -``` - -### Read file line by line - -```bash -while IFS= read -r line; do - echo "$line" -done < file.txt -``` - -### Loop over files - -```bash -for file in *.js; do - echo "Processing $file" -done -``` - ---- - -## 9. Differences from PowerShell - -| Task | PowerShell | Bash | -|------|------------|------| -| List files | `Get-ChildItem` | `ls -la` | -| Find files | `Get-ChildItem -Recurse` | `find . -type f` | -| Environment | `$env:VAR` | `$VAR` | -| String concat | `"$a$b"` | `"$a$b"` (same) | -| Null check | `if ($x)` | `if [ -n "$x" ]` | -| Pipeline | Object-based | Text-based | - ---- - -## 10. Error Handling - -### Set options - -```bash -set -e # Exit on error -set -u # Exit on undefined variable -set -o pipefail # Exit on pipe failure -set -x # Debug: print commands -``` - -### Trap for cleanup - -```bash -cleanup() { - echo "Cleaning up..." - rm -f /tmp/tempfile -} -trap cleanup EXIT -``` - ---- - -> **Remember:** Bash is text-based. Use `&&` for success chains, `set -e` for safety, and quote your variables! - -## When to Use -This skill is applicable to execute the workflow or actions described in the overview. - -## Limitations -- Use this skill only when the task clearly matches the scope described above. -- Do not treat the output as a substitute for environment-specific validation, testing, or expert review. -- Stop and ask for clarification if required inputs, permissions, safety boundaries, or success criteria are missing. +- 不要对未确认路径执行递归删除。 +- 不要把不可信文件名直接拼进危险命令。 +- 不要默认 GNU 命令参数在 macOS 上一定可用。 +- 不要留下用户不需要的后台进程。 +- 不要在没确认网络边界时把容器 localhost 当宿主机。 diff --git a/skills/logic-planner/SKILL.md b/skills/logic-planner/SKILL.md index a1dcdaf..35b5244 100644 --- a/skills/logic-planner/SKILL.md +++ b/skills/logic-planner/SKILL.md @@ -1,161 +1,149 @@ --- name: logic-planner -description: 当需要把一份已校验的数据需求(候选表/字段/join 关联已确认)拆解为 SQL 生成器可机械执行的步骤计划时使用。触发词包括"拆解逻辑"、"执行计划"、"logic plan"、"SQL 编排"、"怎么算"、"指标拆解"、"join 顺序"。消费 metadata-validator 的 Validation Result,产出固定步骤分类法(Source / Filter / Join / Transform / Aggregate / Output)的 Logic Plan。**不写 SQL**。 +description: 当需求和元数据都已经确认,需要把业务逻辑拆成 SQL 生成前的确定性执行步骤时使用。它只规划 source、filter、join、dedupe、transform、window、aggregate、output,不写 SQL,不调用 MCP。 --- -# Logic Planner(逻辑规划器) +# Logic Planner(逻辑规划) -## 概述(Overview) +## 目标 -将一份已校验的数据需求分解为固定分类法下的确定性步骤流水线。Logic Plan 是 `sql-context-builder` 拼装 SQL Context **唯一**的输入 — 下游不允许再做任何业务推理。 +把已确认的 DRD 和元数据拆成可执行计算步骤,让后续 SQL 只是翻译计划,而不是重新猜业务逻辑。 -核心原则:**只拆分,不翻译。** 本 skill 绝不写 SQL 语法,只在业务/逻辑层面描述每一步要发生什么。SQL 翻译是另一回事。 +## 输入 -## 何时使用(When to Use) - -**使用场景:** - -- `metadata-validator` 返回状态 `VALIDATED`,需要设计查询如何执行 -- 用户询问"这个应该怎么算"、"逻辑是啥"、"一步步怎么走"等数据需求问题 -- 写 SQL 之前需要先消除粒度、指标公式、去重、窗口、排名的歧义 - -**不要使用场景:** - -- `metadata-validator` 返回 `NEED_USER_CONFIRMATION` — 先回去解决未决问题 -- 用户已经提供了书面执行计划,只是要 SQL — 直接进入 `sql-context-builder` - -## 输入(Inputs) - -1. 来自 `metadata-validator` 的 **Validation Result**(状态为 `VALIDATED`) -2. 原始业务需求(通常嵌在 `requirements-analysis` 的输出中) -3. `field_mapping`(必备) — 每个业务字段在规划开始前必须先映射到物理 `(table, column, type)` - -## 步骤分类法(Step Taxonomy) - -每个 Logic Plan 都是一组步骤。每个步骤有且仅有一个 `kind`,且只能从下表选取。对重复操作(如多次 join、多次 filter)复用同一种 kind。 - -| Kind | 用途 | 必填字段 | -|---|---|---| -| `source` | 从表或子查询读取 | `alias`, `table`, `columns[]` | -| `filter` | 在指定 alias 上应用谓词 | `alias`, `predicates[]` | -| `join` | 在 key 上合并两个 alias | `left`, `right`, `keys[]`, `type` (inner/left/right/full/semi/anti) | -| `transform` | 派生新列(case-when、算术、类型转换) | `alias`, `expressions[]` | -| `aggregate` | 分组聚合 | `alias`, `group_by[]`, `metrics[]` (每个含 `expr`, `alias`, `agg`) | -| `dedupe` | DISTINCT 或基于 row_number 的去重 | `alias`, `keys[]`, `strategy` | -| `window` | 排名 / 累计 / lag-lead | `alias`, `partition_by[]`, `order_by[]`, `window_fn` | -| `output` | 最终投影和排序 | `columns[]`, `order_by[]`, `limit?` | - -不要发明新 kind。如果某个变换无法归入上述分类,作为 `pending_question` 暴露并停止。 - -## 工作流(Workflow) - -```dot -digraph logic_planner { - "读取 Validation Result + field_mapping" [shape=box]; - "区分事实表 / 维度表,选定主源表" [shape=box]; - "列出过滤条件(时间区间、状态、业务谓词)" [shape=box]; - "规划 join(alias A ON key = alias B.key)" [shape=box]; - "规划 transform(派生列、类型转换)" [shape=box]; - "如需去重 / 窗口,一并规划" [shape=box]; - "规划 aggregate(粒度 + 指标)" [shape=box]; - "规划 output 投影 + 排序" [shape=box]; - "还有歧义?" [shape=diamond]; - "输出 Logic Plan" [shape=box]; - "输出 Logic Plan + pending_questions" [shape=box]; - - "读取 Validation Result + field_mapping" -> "区分事实表 / 维度表,选定主源表"; - "区分事实表 / 维度表,选定主源表" -> "列出过滤条件(时间区间、状态、业务谓词)"; - "列出过滤条件(时间区间、状态、业务谓词)" -> "规划 join(alias A ON key = alias B.key)"; - "规划 join(alias A ON key = alias B.key)" -> "规划 transform(派生列、类型转换)"; - "规划 transform(派生列、类型转换)" -> "如需去重 / 窗口,一并规划"; - "如需去重 / 窗口,一并规划" -> "规划 aggregate(粒度 + group_by)"; - "规划 aggregate(粒度 + group_by)" -> "规划 output 投影 + 排序"; - "规划 output 投影 + 排序" -> "还有歧义?"; - "还有歧义?" -> "输出 Logic Plan" [label="no"]; - "还有歧义?" -> "输出 Logic Plan + pending_questions" [label="yes"]; -} +```yaml +requirements_output: + status: READY_FOR_METADATA +validation_result: + status: VALIDATED + field_mapping: {} + joins: [] ``` -## 必须给出的推理(Required Reasoning,不可跳过) +如果元数据没有 `VALIDATED`,停止并回到 `metadata-validator`。 -在输出 plan 之前,你必须显式陈述并记录: +## 步骤类型 -1. **粒度(Grain)** — 输出一行代表什么?(例如:一个客户、一个客户×月、一笔订单) -2. **指标公式(Metric formula)** — 对每个指标写出精确表达式:分子与分母都要定义 -3. **时间粒度(Time grain)** — 确认时间字段与粒度(天 / 月 / 小时),用闭/开区间说明窗口 -4. **Join 基数(Join cardinality)** — 对每个 join,预测 1:1 / 1:N / N:N。如果是 N:N,必须标记并询问用户是否在 join 前去重 -5. **去重位置(Dedup position)** — 如果事实表存在行重复(例如一个客户多个地址),在 join 之前决定哪一侧、在哪个 key 上先去重 -6. **空值策略(Null strategy)** — 对可空 join key 或度量,决定用 `INNER` 还是 `LEFT + COALESCE` +只允许使用这些 `kind`: -以上任意一项不清楚时,加入 `pending_question` 并停止。绝不臆测。 +- `source`:读取哪张表、哪些字段。 +- `filter`:时间、分区、状态、业务条件过滤。 +- `join`:多表关联。 +- `dedupe`:去重,通常在 join 或 aggregate 前。 +- `transform`:派生字段、分类、类型转换、标志位。 +- `window`:首笔、末笔、排名、累计、lag/lead。 +- `aggregate`:分组和指标计算。 +- `output`:最终列、排序、limit 或写入意图。 -## 输出契约(Output Contract) +## 工作顺序 + +1. 定义输出粒度:一行代表什么。 +2. 选事实表:承载核心事件或分母人群的表。 +3. 列出 source:每张表只取必要字段。 +4. 尽早放 filter:尤其是时间、分区和状态条件。 +5. 固定时间窗口边界:下界包含、上界排除。 +6. 规划 join:顺序、key、类型、基数假设。 +7. 判断是否需要 dedupe/window:首登、首次注册、首笔成功事件等。 +8. 规划 transform:标志位、事件顺序、空值处理。 +9. 规划 aggregate:group by、指标公式、别名。 +10. 规划 output:输出字段、排序、写入模式。 + +## 时间窗口硬规则 + +当需求为“近 N 天按 T-1 完整日”时,必须规划为: + +```yaml +time_window: + lower_bound_sql: date_sub(current_date(), N) + upper_bound_sql: current_date() + boundary: left_closed_right_open + predicate_sql: event_time >= date_sub(current_date(), N) AND event_time < current_date() +``` + +禁止规划为 `BETWEEN ... AND date_add(current_date(), -1)`,尤其是 timestamp 字段。 + +## 转化/漏斗硬规则 + +转化类需求必须维护事件顺序。对“注册后成功付费”,推荐计划: + +1. `first_login`:取窗口内每个用户的首次登录,形成分母。 +2. `registered`:取用户注册事件,必要时去重到一个注册时间。 +3. `qualified_payment`:在关联注册事件后筛选 `pay_status = 'SUCCESS' AND pay_time >= register_time`。 +4. `final_agg`:按维度计算分母、分子和转化率。 + +不要先对全量 `pay_order` 按用户取 `MIN(pay_time)` 再与 `register_time` 比较;这可能漏掉“历史早付费但注册后也付费”的用户。 + +## 必须显式写出的风险 + +- 1:N 或 N:N join 是否会放大事实行。 +- 去重发生在 join 前还是 join 后。 +- 指标是否需要 `count_distinct` 而不是 `count`。 +- 时间过滤是否能命中分区字段。 +- left join 后维度缺失如何处理。 +- 转化事件是否严格满足前后顺序。 +- 分母为 0 时比率如何处理。 + +## 输出格式 ```yaml logic_plan: - grain: "<一句话描述输出一行代表什么>" - fact_table: - dimension_tables: - - + status: PLANNED | NEED_USER_CONFIRMATION + grain: "结果一行代表什么" + fact_table: "事实表" + dimension_tables: [] + time_window: + lower_bound_sql: "date_sub(current_date(), N)" + upper_bound_sql: "current_date()" + predicate_sql: "event_time >= date_sub(current_date(), N) AND event_time < current_date()" + boundary: left_closed_right_open steps: - - kind: source - alias: o - table: loan_order - columns: [customer_id, loan_amount, apply_time, status] - - kind: filter - alias: o - predicates: - - "o.apply_time >= ''" - - "o.apply_time < ''" - - "o.status = 'SUCCESS'" - - kind: join - left: { alias: o } - right: { alias: c, table: customer_info } - keys: ["o.customer_id = c.customer_id"] + - id: S1 + kind: source + name: first_login + table: "表名" + columns: ["字段"] + reason: "为什么需要" + - id: S2 + kind: filter + target: "表名或中间结果" + predicates: ["业务过滤条件"] + partition_pruning: true + - id: S3 + kind: join + left: "左表或CTE" + right: "右表或CTE" + keys: ["左字段 = 右字段"] type: left - cardinality_assumption: 1:N # if known; else state "unknown" - - kind: transform - alias: o - expressions: - - expr: "CASE WHEN o.loan_amount >= 100000 THEN 'large' ELSE 'small' END" - alias: loan_bucket - - kind: aggregate - alias: o - group_by: ["c.city"] + cardinality: "1:1 | 1:N | N:1 | N:N | unknown" + row_growth_risk: low | medium | high + - id: S4 + kind: aggregate + group_by: ["维度字段"] metrics: - - expr: "SUM(o.loan_amount)" - alias: total_loan - agg: sum - - expr: "COUNT(DISTINCT o.customer_id)" - alias: customer_cnt - agg: count_distinct - - kind: output - columns: ["c.city AS city", "total_loan", "customer_cnt"] - order_by: ["total_loan DESC"] - limit: 100 - pending_questions: - - "<尚未解决的歧义,会阻塞 SQL 生成>" - status: + - name: "指标名" + agg: count_distinct | sum | avg | min | max | ratio + expression: "业务表达式" + alias: "指标别名" + - id: S5 + kind: output + columns: ["输出列"] + order_by: [] + write_intent: query_only | insert_into | insert_overwrite + assumptions: [] + risks: [] + pending_questions: [] ``` -## 禁止行为(Forbidden Behaviors) +## 用户可见摘要 -| 自我说服 | 现实 | -|---|---| -| "我直接写 SELECT 列表,用户能读懂 SQL" | Logic Plan 用业务术语,SQL 是下游的事。 | -| "聚合太显然,公式就省了" | 显式给出分子/分母才能避免静默的指标漂移。 | -| "1:N 还是 N:N 对 SELECT 来说无所谓" | 它决定了是否要去重以及是否有膨胀风险,影响很大。 | -| "时间粒度是用户的问题" | 必须由规划器选定一个粒度,选错了数字就全错。 | -| "HAVING 我到 SQL 那一步再加" | 所有过滤(含聚合后)都必须在 Logic Plan 中列出。 | +除了 YAML,必须输出这张表: -## 完成判定(Completion Criteria) +```markdown +## 逻辑计划 +| 步骤 | 中间结果 | 类型 | 说明 | 关键条件 | +|---|---|---|---|---| +``` -要输出状态 `PLANNED`,必须满足: +## 停止条件 -- 粒度用一句话描述,不是一段话 -- 每个 `field_mapping` 条目都被恰好一个 step 引用 -- 每个 join 都有 cardinality 假设 -- 每个指标都有显式 `agg` 和 `alias` -- `pending_questions` 为空 - -否则:状态为 `NEED_USER_CONFIRMATION` 并停止。**不得**进入 `sql-context-builder`。 +如果粒度、指标公式、join 基数、去重策略、空值策略、时间归属或写入模式不明确,输出 `NEED_USER_CONFIRMATION`。 diff --git a/skills/metadata-validator/SKILL.md b/skills/metadata-validator/SKILL.md index b13cd1a..c56ddeb 100644 --- a/skills/metadata-validator/SKILL.md +++ b/skills/metadata-validator/SKILL.md @@ -1,165 +1,98 @@ --- name: metadata-validator -description: 当需要校验某个数据需求中候选的表、字段、join 关联和时间字段是否真实存在于工作区元数据中时使用。触发词包括"validate metadata"、"check if table exists"、"field mapping"、"join key validation"、"确认表/字段"、"校验 schema"、"这个字段在哪张表"、"口径对得上吗"。递归扫描工作区中的元数据文件(JSON / Markdown / CSV / Excel / YAML / DDL / 数据字典),并产出标准化的 Validation Result 和 Field Mapping。 +description: 当需要确认候选表、字段、字段归属、字段类型、时间列、分区列和 join key 是否真实存在于工作区元数据中时使用。消费 requirements-analysis 的 requirements_output,产出 validation_result 和 field_mapping;不写 SQL,不调用 MCP。 --- -# Metadata Validator(元数据校验器) +# Metadata Validator(元数据校验) -## 概述(Overview) +## 目标 -本 skill 在需求的"期望表/字段"和工作区元数据中的"实际 schema"之间架起桥梁。输出会被 `logic-planner` 和 `sql-context-builder` 原样消费,二者**不得**重复校验。 +把需求里“想用的表和字段”与真实元数据对齐,防止模型凭感觉编表名、字段名和关联键。下游只能使用本 skill 确认过的字段映射。 -核心原则:**绝不臆测表名、字段名或 join 关系。** 如果某候选在元数据中无法证实存在,将其作为 `pending_question` 暴露,状态停在 `NEED_USER_CONFIRMATION`。 +## 输入 -## 何时使用(When to Use) - -**使用场景:** - -- `requirements-analysis` 已经产出候选表/字段,需要确认它们在真实元数据中存在 -- 用户询问"这张表/这个字段是否存在"、"表 X 的元数据在哪份文件里"、"业务字段映射到哪个物理列" -- 需要检测 join key、时间列、可用于聚合的数值列 -- 即将写 SQL,必须防止出现幻觉字段名 - -**不要使用场景:** - -- 用户已经提供并确认了精确的表名+列名 — 直接进入 `logic-planner` -- 任务是无数据仓库上下文的纯代码生成 - -## 输入(Inputs) - -1. 来自 `requirements-analysis` 的: - - `candidate_tables`: 用户提到或隐含的表名列表 - - `candidate_fields`: 业务字段名及其建议所在表 - - `business_logic`: 计算逻辑的简短描述 -2. 工作区元数据:当前工作目录下任何符合元数据约定的文件(见下文) - -## 工作流(Workflow) - -```dot -digraph metadata_validator { - "接收候选表/字段" [shape=box]; - "递归扫描工作区中的元数据文件" [shape=box]; - "将每个文件解析为规范化的 (table, field, type) 记录" [shape=box]; - "检查候选表是否存在" [shape=box]; - "检查候选字段是否在其所属表中存在" [shape=box]; - "对缺失字段做模糊匹配" [shape=box]; - "检测 join key 与时间列" [shape=box]; - "所有检查都通过?" [shape=diamond]; - "构建 Field Mapping" [shape=box]; - "输出 Validation Result (VALIDATED)" [shape=box]; - "输出 Validation Result (NEED_USER_CONFIRMATION)" [shape=box]; - - "接收候选表/字段" -> "递归扫描工作区中的元数据文件"; - "递归扫描工作区中的元数据文件" -> "将每个文件解析为规范化的 (table, field, type) 记录"; - "将每个文件解析为规范化的 (table, field, type) 记录" -> "检查候选表是否存在"; - "检查候选表是否存在" -> "检查候选字段是否在其所属表中存在"; - "检查候选字段是否在其所属表中存在" -> "对缺失字段做模糊匹配"; - "对缺失字段做模糊匹配" -> "检测 join key 与时间列"; - "检测 join key 与时间列" -> "所有检查都通过?"; - "所有检查都通过?" -> "构建 Field Mapping" [label="yes"]; - "所有检查都通过?" -> "输出 Validation Result (NEED_USER_CONFIRMATION)" [label="no"]; - "构建 Field Mapping" -> "输出 Validation Result (VALIDATED)"; -} +```yaml +requirements_output: + status: READY_FOR_METADATA + candidate_tables: [] + candidate_fields: [] + join_hints: [] +metadata_files: "当前工作区中的数据字典、DDL、schema、CSV 表头、Markdown 表格、Excel 等" ``` -## 元数据发现(Metadata Discovery) +如果 `requirements_output.status` 不是 `READY_FOR_METADATA`,停止并回到 `requirements-analysis`。 -- 递归扫描**当前工作目录**(及其所有子目录) — 不要假设存在 `metadata/` 目录 -- 文件名常常形似表名(例如 `loan_order.json`、`customer_info.md`、`dim_product.csv`),但**绝不要只信任文件名** — 必须打开文件,确认里面的表标识符 -- 支持的格式(按内容自动识别,不按扩展名): - - JSON / JSONL - - YAML - - Markdown 表格 / 标题 - - CSV(带表头) - - Excel(`.xlsx`、`.xls`) — 用 `openpyxl` 或 `pandas.read_excel` - - DDL / `CREATE TABLE` 语句 - - 自由形式的数据字典(解析键值行或 `field | type | comment` 表格) -- 在校验之前,先把每个文件转成规范化的记录格式: +## 元数据查找顺序 - ```yaml - tables: - - table: <物理表名> - source_file: <相对路径> - fields: - - name: <字段名> - type: - comment: <可选> - ``` +1. 递归查看工作区中的元数据文件,不只相信文件名。 +2. 优先读取结构化文件:DDL、JSON/YAML schema、CSV/Excel 数据字典、Markdown 表格。 +3. 对每张表记录证据:文件路径、表名出现位置、字段名、类型、注释、枚举。 +4. 文件名只能作为线索,不能作为存在性证据。 -## 校验规则(Validation Rules) +## 校验规则 -**全部**以下规则都要执行。任意一条失败 ⇒ 状态为 `NEED_USER_CONFIRMATION`。 +- 字段在别的表里存在,不等于在目标表里存在。 +- 字段名相似不等于字段可替换,必须保留候选并让用户确认。 +- 时间字段有多个时必须问清楚事件时间、分区时间和统计归属。 +- join key 有多个候选时必须问清楚,并标注基数风险。 +- 金额、数量、比率类指标要确认类型、单位和精度。 +- 结果表如果要写入,必须校验目标表存在性和字段兼容性;`INSERT OVERWRITE` 还需要用户显式确认。 -1. **表存在性** — 每个候选表必须出现在解析后的元数据中 -2. **字段存在性** — 每个候选字段必须出现在其候选所属表中 -3. **字段归属** — 字段可能存在但属于另一张表;绝不能悄悄替换 -4. **类型合理性** — 聚合目标(sum / avg / count)必须为数值型;时间过滤字段必须为 date/timestamp/字符串型日期 -5. **必要字段** — 如果业务逻辑中常见字段缺失(统计指标、时间、join key、维度),要标记出来 -6. **模糊匹配** — 对每个缺失字段,按名称相似度打分(`customer_id` ↔ `cust_id`、`cust_no`、`customer_no`)和 Levenshtein 距离,给出 Top-N -7. **join 检测** — 在已校验的表之间,查找共享的 `*_id`、`*_no`、`*_code` 列。如果存在多个可能的 join key,必须询问 -8. **时间列消歧** — 列出每张表中所有时间类字段(`create_time`、`apply_time`、`txn_date`、`update_time`、`dt`);询问哪一个决定时间粒度 -9. **覆盖度** — 验证维度、指标、过滤列、排序键是否都存在 +## 工作顺序 -## 禁止行为(Forbidden Behaviors) +1. 整理所有候选表和候选字段。 +2. 从元数据中确认表是否存在。 +3. 校验字段是否存在于正确表中,并记录类型和来源。 +4. 校验指标字段是否适合聚合,时间字段是否适合时间过滤。 +5. 校验 join key 两边都存在,记录类型一致性和基数假设。 +6. 校验分区字段、状态枚举、金额单位等会影响结果的细节。 +7. 对缺失项给出证据和候选替代,不猜。 +8. 如果缺表、缺字段或 join/time 不明确,停止。 -| 自我说服 | 现实 | -|---|---| -| "文件名是 `loan_order.json`,肯定就是这张表" | 文件名只是提示,必须解析文件内容确认。 | -| "字段存在于某处,所以 join 没问题" | 错误的表归属会让 join 失效。永远要校验归属。 | -| "用户八成指的是 `cust_no`,直接用吧" | 必须作为候选暴露并询问,绝不能自动替换。 | -| "时间字段我猜一个就行" | 多个时间字段时必须询问。粒度决定整条查询。 | -| "这份元数据够了,不必再扫" | 永远要递归扫描。一个文件可能描述多张表。 | - -## 输出契约(Output Contract) - -输出是一份单一的 YAML 文档,下游 skill 原样消费。 +## 输出格式 ```yaml validation_result: - validated_tables: - - - validated_fields: - - . - missing_tables: [] - missing_fields: - - - candidate_fields: - : - - - - - candidate_tables: - : - - - join_candidates: - - left: . - right: . - confidence: - time_field_candidates: - : - - - - + status: VALIDATED | NEED_USER_CONFIRMATION metadata_sources: - - <相对路径指向元数据文件> - pending_questions: - - "<面向用户的问题,最好给出 A/B/C 选项>" + - path: "元数据文件路径" + evidence: "表/字段证据摘要" + validated_tables: + - table: "物理表名" + role: fact | dimension | lookup | target | unknown + source_file: "元数据来源" + validated_fields: + - table: "表名" + column: "字段名" + type: "字段类型" + comment: "字段注释" + source_file: "元数据来源" field_mapping: - : - table: <物理表> - column: <物理字段> - type: <数据类型> - status: + "业务字段名": + table: "物理表名" + column: "物理字段名" + type: "字段类型" + role: dimension | metric | filter | time | partition | join_key | output + source_file: "元数据来源" + confidence: high | medium | low + joins: + - left: "表A.字段" + right: "表B.字段" + type: inner | left | right | full | unknown + cardinality: "1:1 | 1:N | N:1 | N:N | unknown" + confidence: high | medium | low + time_fields: + - table: "表名" + event_time: "事件时间字段" + partition_time: "分区字段,可为空" + missing_tables: [] + missing_fields: [] + alternatives: + "缺失项": ["可能候选"] + pending_questions: + - "需要用户确认的问题" ``` -`field_mapping` 是下游 `logic-planner` 和 `sql-context-builder` 的**唯一权威来源**。如果某个业务字段没有映射,该需求无法继续 — 把它加入 `pending_questions`。 +## 停止条件 -## 完成判定(Completion Criteria) - -仅当**以下全部**成立时,才允许状态 `VALIDATED`: - -- 每个候选表都出现在 `validated_tables` 中 -- 每个候选字段都有 `field_mapping` 条目 -- 所有 join key 都已确定(或显然到可以用 `high` 置信度推断) -- 查询粒度对应的时间列已确定 -- `pending_questions` 为空 - -其余情况 ⇒ `NEED_USER_CONFIRMATION` 并停止。**不得**继续进入 `logic-planner`。 +只在所有必需表、字段、join key、时间字段和目标写入字段都确认后输出 `VALIDATED`。否则输出 `NEED_USER_CONFIRMATION`。 diff --git a/skills/pyspark-sql-guardrails/SKILL.md b/skills/pyspark-sql-guardrails/SKILL.md index bdf5249..b3b79c6 100644 --- a/skills/pyspark-sql-guardrails/SKILL.md +++ b/skills/pyspark-sql-guardrails/SKILL.md @@ -1,391 +1,85 @@ --- name: pyspark-sql-guardrails -description: 在编写或评审 PySpark 脚本、Spark SQL 字符串、`spark.sql` 调用、从配置文件加载的 SQL、ETL 作业、数仓维护脚本,或任何必须把 SQL 限制为 SELECT 和 INSERT 的数据流水线代码时使用。 +description: 当生成、展示、评审或执行 PySpark/Spark SQL 字符串时使用。它负责将 SQL 限制在 SELECT、WITH...SELECT、INSERT...SELECT 安全白名单内,并阻止危险 DDL/DML/维护语句、堆叠语句和未确认的 INSERT OVERWRITE。 --- -# PySpark SQL Guardrails(PySpark SQL 防护栏) +# PySpark SQL Guardrails(SQL 安全闸门) -## 校验关口(先读这一段 — 不可妥协) +## 目标 -**本 skill 产出的每一条 SQL 字符串,在跑过自带的校验器之前都是不可信的。** 校验器只需要一条命令,在 SQL 生成之后、**展示给用户之前**立即执行。不存在绕过此步骤的路径。 +在 SQL 被展示、写入 PySpark 脚本或交给 `spark.sql` 前,先通过白名单校验。通过本 skill 只代表 SQL 语句类型安全,不代表业务逻辑正确。 -整条流程就是两条命令 + 一个条件: +## 允许 -```bash -# Step 1 — Generate the SQL (do this in your head / draft) -# Step 2 — Run the validator on it (do this IMMEDIATELY, in the same turn) -echo "" | python3 skills/pyspark-sql-guardrails/scripts/validate_sql.py - -# Step 3 — Read the output: -# PASS - first keyword: select, body length: N → you may now show the SQL -# FAIL - → STOP. Fix the SQL. Re-run. -``` - -**如果你还没把 `PASS -` 这一行贴到对话里,就不允许贴 SQL。** 校验器就是关口,SQL 站在关口的另一边,绝不能空着手闯关。 - -这条规则由本 skill 自己强制,而非外部系统。遵守它是任何含 SQL 字符串的输出的前提。 - ---- - -## 概述(Overview) - -本环境下的 PySpark SQL 只接受白名单:生成或执行的 SQL 只能是 `SELECT` 或 `INSERT`。在通过校验之前,所有 SQL 字符串、格式化后的 SQL 模板、从配置加载的 SQL、notebook 单元格,以及 `spark.sql(...)` 调用一律视为不安全。 - -核心原则:**绝不能靠"看上去不像危险"来执行 SQL;只有在证明它就是一条以 SELECT 或 INSERT 开头的、合法的语句后,才能执行。** - -## 强制校验(Hard Rule) - -**每条 SQL 字符串在到达 `spark.sql(...)` 之前都必须通过 `assert_select_or_insert()` — 不许有例外。** 这条规则不可妥协,适用于: - -- 你刚生成的 SQL(字面量、f-string、`.format(...)`、`+` 拼接) -- 从 YAML / JSON / INI / 环境变量 / 命令行 / 配置文件加载的 SQL -- 通过参数、函数参数、notebook 变量传入的 SQL -- 从聊天消息里粘贴进来的 SQL -- 由模板引擎(Jinja、string.Template 等)产出的 SQL - -唯一可接受的调用形式: - -```python -# Preferred: wrapper that validates + executes atomically -safe_spark_sql(spark, sql) - -# Or: validate first, then execute explicitly -assert_select_or_insert(sql) # raises on violation -spark.sql(sql) -``` - -**禁止的调用形式(出现即视为流水线失败):** - -```python -spark.sql(sql) # direct, no validation -spark.sql(f"SELECT ... {user_input} ...") # f-string into spark.sql -spark.sql(config["sql"]) # config-driven without validation -spark.sql(open("queries/xxx.sql").read()) # file-loaded without validation -``` - -如果 `assert_select_or_insert()` 抛错,流水线立即停止。**不要**削弱规则,**不要**用正则去剥除禁用关键字,**不要**把 SQL"改写"成看似安全的样子。要么拒绝,要么重新生成。 - -## 必守规则(Required Rule) - -**允许:** - `SELECT ...` - `WITH ... SELECT ...` - `INSERT INTO ... SELECT ...` -- `INSERT OVERWRITE ... SELECT ...` — 仅当用户明确允许本次任务使用 overwrite,否则必须先问 +- `INSERT OVERWRITE ... SELECT ...`,但必须拿到用户明确允许 overwrite 的确认,并使用 `allow_overwrite=true` -**禁止**(即便是维护或元数据刷新): -- `DELETE`、`TRUNCATE`、`DROP`、`ALTER`、`CREATE`、`REPLACE`、`MERGE`、`UPDATE` -- `MSCK REPAIR`、`REFRESH`、`ANALYZE`、`CACHE`、`UNCACHE`、`VACUUM`、`OPTIMIZE` -- 用分号分隔的多条语句 -- 从 YAML/JSON/env/CLI 来的、未经任何处理直接喂给 `spark.sql` 的 SQL +## 禁止 -**执行关口:** 每次 `spark.sql` 调用都必须先调 `assert_select_or_insert(sql)`(或 `safe_spark_sql(spark, sql)`)。没有快速通道。详见上面的 *强制校验(Hard Rule)*。 +- `DROP`、`DELETE`、`UPDATE`、`MERGE`、`ALTER`、`CREATE`、`REPLACE` +- `TRUNCATE`、`MSCK`、`REFRESH`、`ANALYZE`、`CACHE`、`UNCACHE` +- `VACUUM`、`OPTIMIZE`、`CALL`、`GRANT`、`REVOKE` +- 用分号堆叠多条语句 +- `INSERT ... VALUES` 或不基于 `SELECT` 的写入 +- 未校验就直接进入 `spark.sql(sql)` -## 速查表(Quick Reference) +## 使用脚本 -| 场景 | 做法 | -|---|---| -| 想要取数 | `safe_spark_sql(spark, "SELECT ...")` | -| 想要写数 | `safe_spark_sql(spark, "INSERT INTO table SELECT ...")` | -| 想要清理/去重/删除 | 用 SELECT 派生干净数据,再用 `safe_spark_sql` + INSERT 写入已批准的目标/staging 表;**不要**直接删除旧数据 | -| 想要做 schema/表/分区维护 | 停下来询问用户;**不要**输出 DDL 或元数据 SQL | -| SQL 来自配置文件 | 先 `assert_select_or_insert(sql)`,**再** `spark.sql(sql)` | -| 用户要求"快速修一下" | 安全规则照旧 — 先跑 `assert_select_or_insert` | -| SQL 由 `sql-context-builder` 生成 | 执行前必须跑 `assert_select_or_insert` | - -## 安全模板(Safe Pattern) - -**必须:** 每个 `spark.sql` 调用点都走 `safe_spark_sql`(或在 `spark.sql` 之前立即调 `assert_select_or_insert`)。在校验器面前,一条生成的 SQL 和集群之间只隔这一道闸。 - -在 SQL 进入执行的那条边界做校验。**集群环境不要 import `sql_guard`** — 文件在 driver 节点上存在,在 executor 节点上可能不存在,`import` 会报 `ModuleNotFoundError`。直接把下面这段代码**内联**到 PySpark 脚本里(或通过 `--files` 分发后 import)。 - -```python -import re - -_FORBIDDEN_SQL = re.compile( - r"\b(delete|truncate|drop|alter|create|replace|merge|update|msck|refresh|analyze|cache|uncache|vacuum|optimize)\b", - re.IGNORECASE, -) -_LEADING_COMMENTS = re.compile(r"\A\s*(?:--[^\n]*\n|/\*.*?\*/\s*)*", re.DOTALL) -_STRING_LITERAL = re.compile(r"'(?:[^']|'')*'|\"(?:[^\"]|\"\")*\"") - - -def _strip_string_literals(sql: str) -> str: - """Replace quoted string literals with empty strings so that the - forbidden-keyword regex does not false-positive on text content - inside quotes (e.g. ``SELECT 'drop' AS note``). - """ - return _STRING_LITERAL.sub("''", sql) - - -def assert_select_or_insert(sql: str) -> str: - text = sql.strip() - if not text: - raise ValueError("SQL is empty") - - body = text[:-1].strip() if text.endswith(";") else text - if ";" in body: - raise ValueError("Multiple SQL statements are not allowed") - - normalized = _LEADING_COMMENTS.sub("", body).lstrip() - first = normalized.split(None, 1)[0].lower() if normalized else "" - - if first not in {"select", "with", "insert"}: - raise ValueError(f"Only SELECT and INSERT SQL are allowed, got {first!r}") - - if _FORBIDDEN_SQL.search(_strip_string_literals(normalized)): - raise ValueError("Forbidden SQL keyword found") - - if first == "with" and not re.search(r"\bselect\b", normalized, re.IGNORECASE): - raise ValueError("WITH statements must be SELECT queries") - - return body - - -def safe_spark_sql(spark, sql: str): - return spark.sql(assert_select_or_insert(sql)) - - -# Good -safe_spark_sql(spark, """ -INSERT INTO analytics.daily_customer_snapshot -SELECT * FROM staging.daily_customer_snapshot -""") - -# Bad: raises before Spark sees it -safe_spark_sql(spark, "ALTER TABLE analytics.daily_customer_snapshot DROP PARTITION (dt='2026-06-10')") -``` - -如果项目里已经有 `sqlglot` 之类的 SQL 解析器,优先用 AST 校验而非正则。但白名单要保持一致:只能有一条语句、顶层是 SELECT 或 INSERT、不能出现任何禁用的 DDL/DML/维护命令。 - -## 本地验证脚手架(每次生成 SQL 后都要跑) - -校验器如果从来不被运行,就是废物。每当产出一条 SQL 字符串(无论由你、子 agent、模板或工具产出),下一步就是在 Python 里用 `assert_select_or_insert` 跑一遍,然后才把结果给用户看。 - -本 skill 把自己的校验器和测试集内置在 skill 目录里,使契约自包含且与 skill 一起版本化。 - -**场景区分:** - -- **本地开发 / CI / notebook:** 从内置文件 import(`from sql_guard import assert_select_or_insert`)是首选,确保你和 skill 用的是同一版校验器。 -- **集群执行(executor 节点):** **不要 import `sql_guard`。** 文件在 driver 节点上存在,在 executor 节点上往往不存在,`import` 会报 `ModuleNotFoundError`。**直接把下面展示的那段代码内联到 PySpark 脚本里**,或通过 `--files` 将 `sql_guard.py` 分发给 executor 后再 import。 - -``` -skills/pyspark-sql-guardrails/ -├── SKILL.md ← this file -└── scripts/ - ├── sql_guard.py ← canonical validator (assert_select_or_insert + safe_spark_sql) - ├── validate_sql.py ← one-line CLI: pipe SQL in, get PASS/FAIL out - └── test_sql_guard.py ← 18-case standard test set -``` - -### 第 1 步 — 用内置校验器(不复制、不重写) - -skill 自带一个一行 CLI:`validate_sql.py`。用它。不要写不同的调用方式,不要贴内联正则,不要让子 agent 直接调 `assert_select_or_insert`,除非这个 CLI 也坏了。CLI 是 `assert_select_or_insert` 的薄壳,它调用的函数才是真正的权威。 - -**两种把 SQL 喂给校验器的方法:** - -```bash -# Way A (preferred for multi-line): pipe via stdin -cat <<'EOF' | python3 skills/pyspark-sql-guardrails/scripts/validate_sql.py -SELECT - c.city AS city, - SUM(o.loan_amount) AS total_loan -FROM loan_order o -INNER JOIN customer_info c ON o.customer_id = c.customer_id -WHERE o.status = 'SUCCESS' -GROUP BY c.city -EOF - -# Way B (single-line only): pass as first argument -python3 skills/pyspark-sql-guardrails/scripts/validate_sql.py "SELECT 1 FROM dual" -``` - -**退出码:** -- `0` → `PASS - first keyword: , body length: ` -- `1` → `FAIL - `(reason 会指明被违反的规则) -- `2` → 用法错误(没有提供 SQL) - -内置的 `sql_guard.py` 和 `validate_sql.py` 一并展示: - -```python -# sql_guard.py (bundled) -import re - -_FORBIDDEN_SQL = re.compile( - r"\b(delete|truncate|drop|alter|create|replace|merge|update|msck|refresh|analyze|cache|uncache|vacuum|optimize)\b", - re.IGNORECASE, -) -_LEADING_COMMENTS = re.compile(r"\A\s*(?:--[^\n]*\n|/\*.*?\*/\s*)*", re.DOTALL) -_STRING_LITERAL = re.compile(r"'(?:[^']|'')*'|\"(?:[^\"]|\"\")*\"") - - -def _strip_string_literals(sql: str) -> str: - return _STRING_LITERAL.sub("''", sql) - - -def first_keyword(sql: str) -> str: - """Inspection helper for CLI reporting only — does not raise.""" - text = sql.strip() - body = text[:-1].strip() if text.endswith(";") else text - normalized = _LEADING_COMMENTS.sub("", body).lstrip() - return normalized.split(None, 1)[0].lower() if normalized else "" - - -def assert_select_or_insert(sql: str) -> str: - text = sql.strip() - if not text: - raise ValueError("SQL is empty") - body = text[:-1].strip() if text.endswith(";") else text - if ";" in body: - raise ValueError("Multiple SQL statements are not allowed") - normalized = _LEADING_COMMENTS.sub("", body).lstrip() - first = normalized.split(None, 1)[0].lower() if normalized else "" - if first not in {"select", "with", "insert"}: - raise ValueError(f"Only SELECT and INSERT SQL are allowed, got {first!r}") - if _FORBIDDEN_SQL.search(_strip_string_literals(normalized)): - raise ValueError("Forbidden SQL keyword found") - if first == "with" and not re.search(r"\bselect\b", normalized, re.IGNORECASE): - raise ValueError("WITH statements must be SELECT queries") - return body - - -def safe_spark_sql(spark, sql: str): - return spark.sql(assert_select_or_insert(sql)) - - -# validate_sql.py (bundled, abridged) -import sys, os -SKILL_DIR = os.path.dirname(os.path.abspath(__file__)) -sys.path.insert(0, SKILL_DIR) -from sql_guard import assert_select_or_insert, first_keyword - -def main(): - sql = sys.argv[1] if len(sys.argv) > 1 else sys.stdin.read() - try: - body = assert_select_or_insert(sql) - except ValueError as e: - print(f"FAIL - {e}"); return 1 - print(f"PASS - first keyword: {first_keyword(sql)}, body length: {len(body)}") - return 0 -``` - -### 第 2 步 — 每次生成 SQL 后的必用调用模式 - -同一条命令,每次都用,SQL 改成你刚生成的那一条。不许变体,不许走捷径,不许说"待会儿再跑": - -```bash -echo "" | python3 skills/pyspark-sql-guardrails/scripts/validate_sql.py -``` - -**先把校验器的输出贴给用户,再贴 SQL。** 预期的回复样式: +本 skill 自带脚本: ```text -PASS - first keyword: select, body length: 420 +scripts/sql_guard.py 校验库 +scripts/validate_sql.py 命令行入口 +scripts/test_sql_guard.py 标准测试集 ``` -如果出现的是 `FAIL - `,SQL 不允许通过关口。把失败原样报给用户,修 SQL,再跑校验器。循环到 `PASS` 为止。 - -### 第 3 步 — 标准测试集(每次改校验器时跑) - -skill 自带 `test_sql_guard.py`,含 6 个 EXPECT_OK + 12 个 EXPECT_RAISE 用例。每次校验器改动,以及在 CI 中,都要在 skill 目录下跑一遍: +校验一条 SQL: ```bash -# From project root -python3 skills/pyspark-sql-guardrails/scripts/test_sql_guard.py - -# Or from the scripts/ directory -cd skills/pyspark-sql-guardrails/scripts -python3 test_sql_guard.py +python scripts/validate_sql.py "SELECT 1" ``` -预期输出:6 行 `[OK] EXPECT_OK`、12 行 `[OK] EXPECT_RAISE: raised as expected`,以 `ALL TESTS PASSED` 收尾。出现任何偏差都说明校验器有回归 — 修校验器,不是修测试。 +从文件校验: -### 第 4 步 — 失败处理协议 - -当校验器抛错(在 CLI 形式下就是 `FAIL - ` 这一行): - -1. **停下流水线**,不要继续到下一阶段 -2. **读 reason**,它会指明被违反的规则(空 / 多语句 / 错的 verb / 禁用关键字 / WITH 后没 SELECT) -3. **不要**为了放过这条 SQL 而去改校验器。校验器是权威,SQL 是错的 -4. **不要**用正则剥除禁用关键字来"清理"SQL。要拒绝,要么重新生成 SQL -5. **重跑校验器**用修好的 SQL,循环到 `PASS` - -### 禁止行为(校验器执行相关) - -| 自我说服 | 现实 | -|---|---| -| "我就 `spark.sql(sql)`,信它" | 这正是防护栏要拦的。跑校验器。 | -| "校验器的输出在上一轮里" | 校验器只有通过函数调用才"有状态",再跑一遍。 | -| "就一行,肉眼看看就够了" | 肉眼判断正是这个防护栏要防的失败模式。 | -| "我已经对一条类似的 SQL 跑过了" | 每条 SQL 字符串都是新值,要针对实际要交付的字符串重跑。 | -| "我直接内联一段正则,不用 `validate_sql.py`" | skill 只给一个权威 CLI。内联副本会漂移、错过更新。 | -| "SQL 在代码块里,还没'执行'" | 把 SQL 展示给用户就已经是执行了。关口在代码块之前,不在之后。 | - -## 常见自我说服(Common Rationalizations) - -| 借口 | 现实 | -|---|---| -| "只是元数据刷新而已" | `REFRESH`、`MSCK`、`ANALYZE` 都不在 SELECT/INSERT 范围,先问。 | -| "MERGE 是非破坏性的" | 策略只允许 SELECT/INSERT,`MERGE` 一律禁止。 | -| "DELETE+INSERT 是已有模式" | 已有不安全模式不能凌驾于规则之上。 | -| "配置是可信的" | 配置里的 SQL 在 `spark.sql` 前一样要过白名单。 | -| "我可以用正则剥除禁用词" | 不要把不安全的 SQL 改写成看似安全的样子,直接拒绝。 | -| "我们要在 insert 之前做清理" | 用 SELECT 派生干净数据,再用 INSERT 写入;**不要**变更/删除已有数据。 | -| **"SQL 是硬编码的,不用校验。"** | **硬编码字符串一样要走 `assert_select_or_insert`。该函数检查的是字符串本身,不是它的出处。源码里的字面量并不比配置里的值更安全。** | -| **"我肉眼已经看过了,verb 是 SELECT。"** | **肉眼判断正是这个防护栏要防的失败模式。信函数,别信眼睛。`assert_select_or_insert` 没跑过,SQL 在定义上就是不可信的。** | -| **"我到下一轮/下一个 PR 再包。"** | **不行,现在就在调用点加上 wrapper。已提交代码里裸 `spark.sql(...)` 是回归,不是 TODO。** | - -## 红旗信号(Red Flags) - -看到下面任一项就要停下,加校验或先问用户: - -- `spark.sql(config["sql"])`、YAML SQL、env SQL、CLI SQL,或来自外部输入的 f-string SQL -- 任何不是 SELECT/WITH/INSERT 的 SQL verb -- `INSERT OVERWRITE` 没拿到针对 overwrite 语义的明确批准 -- 用分号分批的 SQL -- 分区修复、schema 迁移、过期分区清理、合并压实、vacuum、统计信息收集的请求 -- 措辞为"快速清理一下"、"就刷一下元数据"、"删掉旧分区"、"复用 DELETE+INSERT" 的请求 -- **直接的 `spark.sql(sql)` 调用,前面没有 `assert_select_or_insert(sql)`(也没用 `safe_spark_sql`)— 即便 SQL 字符串是硬编码字面量** -- **某条 PR、diff 或代码评审引入了新的 `spark.sql(...)` 调用点,却没加 wrapper** - -## 常见错误(Common Mistakes) - -- 只检查 `DELETE` 和 `TRUNCATE`;Spark 的风险还包括 `DROP`、`ALTER`、`MERGE`、`UPDATE`、`MSCK`、`REFRESH`、`ANALYZE`、`VACUUM`、`OPTIMIZE` -- 只校验了生成的 SQL,没校验配置里加载的 SQL -- 因为首条语句是 SELECT,就允许了多语句堆叠 -- 把注释当作无害,但后面跟着一句禁用的 SQL -- 在 SQL 里用 `CREATE OR REPLACE TEMP VIEW`;如果需要临时视图,改用 DataFrame 的 `createOrReplaceTempView` -- **直接调 `spark.sql(sql)`,不走 `assert_select_or_insert` / `safe_spark_sql`。白名单在调用点生效,不在 SQL 生成时生效。一条"看起来"安全的 SQL,在函数跑过它之前都不算"已校验"安全。** -- **写内联的 `if sql.startswith("SELECT"): spark.sql(sql)`。这不是校验,只有 `assert_select_or_insert`(或等价且白名单一致的 AST 解析器)才算。** - ---- - -## 在流水线中的位置(Pipeline Position) - -本 skill 是 PySpark SQL 流水线的**写入/执行关口**: - -``` -requirements-analysis → metadata-validator → logic-planner → sql-context-builder → pyspark-sql-guardrails → sql-review +```bash +python scripts/validate_sql.py --file query.sql ``` -**强依赖的上游 skill:** `sql-context-builder` — 每条生成的 SQL 必须带有一份固化了 alias、join key 和字段来源的 `sql_context`。没有伴随 context 的 SQL 在 `sql-review` 中会命中 `HIGH` 风险(见 check #1:reference integrity)。 +JSON 输出: -**强依赖的下游 skill:** `sql-review` — 本防护栏校验"允许哪些语句";`sql-review` 校验"这条被允许的语句是否正确"。两者都通过之后才能 `spark.sql(...)`。**不要**因为防护栏过了就跳过 `sql-review`。 +```bash +python scripts/validate_sql.py --json "SELECT 1" +``` -**执行关口(本 skill 强制):** SQL 字符串通往集群的唯一通道是 `assert_select_or_insert(sql)`(或其封装 `safe_spark_sql(spark, sql)`)。直接的 `spark.sql(sql)` 是流水线违规。这条规则凌驾于流水线顺序 — 即便 `sql-review` 已经通过,执行时仍要过这一道防护栏。防护栏与评审是相互独立的检查,一道过不能代替另一道。 +允许 overwrite 时: -**本防护栏不检查:** +```bash +python scripts/validate_sql.py --allow-overwrite "INSERT OVERWRITE target SELECT * FROM source" +``` -- 列是否真的存在于表(那是 `sql-review` 通过 SQL Context 校验的事) -- join key 是否正确(那是 `sql-review` 的事) -- 聚合是否被重复、行是否会膨胀(那是 `sql-review` 的事) -- 是否带时间/分区过滤(那是 `sql-review` 的事) +## 输出格式 -**本防护栏检查:** +```yaml +sql_guard: + status: PASS | FAIL + first_keyword: select | with | insert | unknown + reason: "失败原因,PASS 时为空" +``` -- 第一个真实关键字是 `SELECT`、`WITH ... SELECT` 或 `INSERT ... SELECT` -- 没有 `;` 堆叠的多语句 -- 体内不出现任何禁用的 verb -- `INSERT OVERWRITE` 必须基于用户的明确批准 +## 必须执行的关口 -校验器抛错时,改 SQL,**不要**削弱校验器。拿不准时,在生成 `INSERT OVERWRITE` 或任何 DDL 之前先问用户。 +1. SQL 生成后、展示给用户前,先运行 guard。 +2. PySpark 代码中出现 `spark.sql(...)` 时,SQL 字符串必须先通过 guard。 +3. guard 失败时,回到 SQL 生成步骤修 SQL,不要削弱校验器。 +4. guard 通过后仍然必须进入 `sql-review`。 + +## 规则 + +- 校验必须针对最终要展示或执行的那条 SQL。 +- 肉眼看过不算通过。 +- 一条类似 SQL 通过,不代表当前 SQL 通过。 +- 字符串字面量和注释里的危险词不应误报。 +- 反引号字段名里的危险词不应误报,但不建议这样命名。 +- `INSERT OVERWRITE` 默认失败,除非用户明确确认 overwrite。 diff --git a/skills/pyspark-sql-guardrails/scripts/sql_guard.py b/skills/pyspark-sql-guardrails/scripts/sql_guard.py index b863305..876028e 100644 --- a/skills/pyspark-sql-guardrails/scripts/sql_guard.py +++ b/skills/pyspark-sql-guardrails/scripts/sql_guard.py @@ -1,80 +1,137 @@ -"""PySpark SQL Guardrails — validator and safe execution wrapper. +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +"""Dependency-free Spark SQL allowlist guard. -This module is the canonical implementation referenced by -.claude/skills/pyspark-sql-guardrails/SKILL.md. Do not redefine the -regex inline at call sites; do not write "lighter" versions. Always -import this module and route `spark.sql(...)` through `safe_spark_sql`. - -The validator is pure-Python and has no Spark / I/O dependencies. It -is safe to import in unit tests, CI, notebooks, and ad-hoc scripts. +The guard answers one question only: is this SQL statement type safe enough +to show or pass to spark.sql? It does not prove business correctness. """ +from __future__ import annotations + import re +from dataclasses import dataclass -_FORBIDDEN_SQL = re.compile( - r"\b(delete|truncate|drop|alter|create|replace|merge|update|msck|refresh|analyze|cache|uncache|vacuum|optimize)\b", - re.IGNORECASE, -) -_LEADING_COMMENTS = re.compile(r"\A\s*(?:--[^\n]*\n|/\*.*?\*/\s*)*", re.DOTALL) -_STRING_LITERAL = re.compile(r"'(?:[^']|'')*'|\"(?:[^\"]|\"\")*\"") +ALLOWED_FIRST = {"select", "with", "insert"} +FORBIDDEN = { + "delete", "truncate", "drop", "alter", "create", "replace", "merge", "update", + "msck", "refresh", "analyze", "cache", "uncache", "vacuum", "optimize", + "call", "grant", "revoke", +} -def _strip_string_literals(sql: str) -> str: - """Replace quoted string literals with empty strings so that the - forbidden-keyword regex does not false-positive on text content - inside quotes (e.g. ``SELECT 'drop' AS note``). +@dataclass(frozen=True) +class GuardResult: + status: str + first_keyword: str + reason: str = "" - Spark uses single quotes for strings and double quotes for delimited - identifiers; both are replaced because the forbidden-keyword set is - verb-shaped (DROP, DELETE, ...) and never a legitimate identifier. - The doubled-quote escape (``''`` / ``""``) is handled inside the regex. - """ - return _STRING_LITERAL.sub("''", sql) + def as_dict(self) -> dict[str, str]: + return { + "status": self.status, + "first_keyword": self.first_keyword, + "reason": self.reason, + } + + +def _mask(sql: str) -> str: + """Mask strings, comments, and backtick identifiers before token checks.""" + out: list[str] = [] + i = 0 + state = "normal" + quote = "" + while i < len(sql): + ch = sql[i] + nxt = sql[i + 1] if i + 1 < len(sql) else "" + if state == "normal": + if ch in {"'", '"', "`"}: + state = "quote" + quote = ch + out.append(" ") + i += 1 + elif ch == "-" and nxt == "-": + state = "line_comment" + out.extend(" ") + i += 2 + elif ch == "/" and nxt == "*": + state = "block_comment" + out.extend(" ") + i += 2 + else: + out.append(ch) + i += 1 + elif state == "quote": + out.append(" ") + if ch == quote: + if nxt == quote and quote in {"'", '"'}: + out.append(" ") + i += 2 + else: + state = "normal" + i += 1 + elif ch == "\\" and nxt: + out.append(" ") + i += 2 + else: + i += 1 + elif state == "line_comment": + out.append("\n" if ch == "\n" else " ") + if ch == "\n": + state = "normal" + i += 1 + else: + out.append("\n" if ch == "\n" else " ") + if ch == "*" and nxt == "/": + out.append(" ") + state = "normal" + i += 2 + else: + i += 1 + return "".join(out) + + +def _body(masked: str) -> str: + text = masked.strip() + return text[:-1].strip() if text.endswith(";") else text def first_keyword(sql: str) -> str: - """Return the lowercased first real SQL keyword after stripping a - leading comment and a trailing semicolon. Pure inspection helper — - does not raise and does not enforce the allowlist. - - Intended for CLI reporting so ``first keyword:`` reflects what the - validator actually saw, not the leading comment marker (``--``). - """ - text = sql.strip() - body = text[:-1].strip() if text.endswith(";") else text - normalized = _LEADING_COMMENTS.sub("", body).lstrip() - return normalized.split(None, 1)[0].lower() if normalized else "" + match = re.search(r"\b([A-Za-z_][\w]*)\b", _body(_mask(sql))) + return match.group(1).lower() if match else "unknown" -def assert_select_or_insert(sql: str) -> str: - """Return the SQL body if it is a single SELECT/INSERT, else raise ValueError. +def validate_sql(sql: str, *, allow_overwrite: bool = False) -> GuardResult: + if not sql or not sql.strip(): + return GuardResult("FAIL", "unknown", "SQL is empty") - Hard-allowlist. No I/O, no logging side effects, no Spark dependency. - Safe to call in unit tests, CI, notebooks, and ad-hoc scripts. - """ - text = sql.strip() - if not text: - raise ValueError("SQL is empty") + masked_body = _body(_mask(sql)) + first = first_keyword(sql) - body = text[:-1].strip() if text.endswith(";") else text - if ";" in body: - raise ValueError("Multiple SQL statements are not allowed") + if ";" in masked_body: + return GuardResult("FAIL", first, "Multiple SQL statements are not allowed") + if first not in ALLOWED_FIRST: + return GuardResult("FAIL", first, f"Only SELECT, WITH, or INSERT is allowed; got {first!r}") - normalized = _LEADING_COMMENTS.sub("", body).lstrip() - first = normalized.split(None, 1)[0].lower() if normalized else "" + lowered = masked_body.lower() + found = sorted(word for word in FORBIDDEN if re.search(rf"\b{word}\b", lowered)) + if found: + return GuardResult("FAIL", first, "Forbidden keyword found: " + ", ".join(found)) + if first == "with" and not re.search(r"\bselect\b", lowered): + return GuardResult("FAIL", first, "WITH must contain SELECT") + if first == "insert": + if not re.search(r"\bselect\b", lowered): + return GuardResult("FAIL", first, "INSERT must be based on SELECT") + if re.search(r"\binsert\s+overwrite\b", lowered) and not allow_overwrite: + return GuardResult("FAIL", first, "INSERT OVERWRITE requires --allow-overwrite") - if first not in {"select", "with", "insert"}: - raise ValueError(f"Only SELECT and INSERT SQL are allowed, got {first!r}") - - if _FORBIDDEN_SQL.search(_strip_string_literals(normalized)): - raise ValueError("Forbidden SQL keyword found") - - if first == "with" and not re.search(r"\bselect\b", normalized, re.IGNORECASE): - raise ValueError("WITH statements must be SELECT queries") - - return body + return GuardResult("PASS", first, "") -def safe_spark_sql(spark, sql: str): - """Validate then execute. The ONLY sanctioned path to spark.sql().""" - return spark.sql(assert_select_or_insert(sql)) +def assert_allowed_sql(sql: str, *, allow_overwrite: bool = False) -> str: + result = validate_sql(sql, allow_overwrite=allow_overwrite) + if result.status != "PASS": + raise ValueError(result.reason) + return sql.strip().rstrip(";") + + +def safe_spark_sql(spark, sql: str, *, allow_overwrite: bool = False): + return spark.sql(assert_allowed_sql(sql, allow_overwrite=allow_overwrite)) diff --git a/skills/pyspark-sql-guardrails/scripts/test_sql_guard.py b/skills/pyspark-sql-guardrails/scripts/test_sql_guard.py index 55ca0d8..80facc0 100644 --- a/skills/pyspark-sql-guardrails/scripts/test_sql_guard.py +++ b/skills/pyspark-sql-guardrails/scripts/test_sql_guard.py @@ -1,67 +1,77 @@ -"""Standard test set for sql_guard.assert_select_or_insert. +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- -Run from the scripts/ subdirectory of the skill: - cd .claude/skills/pyspark-sql-guardrails/scripts - python3 test_sql_guard.py +from sql_guard import assert_allowed_sql, validate_sql -Or from the project root: - python3 .claude/skills/pyspark-sql-guardrails/scripts/test_sql_guard.py +PASS_CASES = [ + ("plain SELECT", "SELECT 1"), + ("WITH SELECT", "WITH x AS (SELECT 1 AS a) SELECT a FROM x"), + ("INSERT INTO SELECT", "INSERT INTO target_table SELECT * FROM source_table"), + ("trailing semicolon", "SELECT 1;"), + ("leading line comment", "-- harmless\nSELECT 1"), + ("leading block comment", "/* harmless */ SELECT 1"), + ("danger word in string", "SELECT 'drop table is text' AS note"), + ("danger word in double quote", 'SELECT "delete" AS note'), + ("danger word in backtick identifier", "SELECT `drop` FROM t"), + ("semicolon in string", "SELECT ';' AS semi"), + ("semicolon in comment", "-- ;\nSELECT 1"), +] -Any change to assert_select_or_insert must be followed by running this -set: all EXPECT_RAISE cases must raise, all EXPECT_OK cases must return -cleanly. Failures are loud (non-zero exit). -""" +PASS_WITH_OVERWRITE = [ + ("INSERT OVERWRITE when explicitly allowed", "INSERT OVERWRITE target SELECT * FROM source"), +] -import sys -import os - -# Allow running this file directly from the scripts/ subdirectory. -sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) - -from sql_guard import assert_select_or_insert +FAIL_CASES = [ + ("DROP", "DROP TABLE x"), + ("UPDATE", "UPDATE t SET a = 1"), + ("DELETE", "DELETE FROM x"), + ("TRUNCATE", "TRUNCATE TABLE x"), + ("ALTER", "ALTER TABLE x ADD COLUMN a int"), + ("CREATE", "CREATE TABLE x(id int)"), + ("REPLACE", "CREATE OR REPLACE VIEW v AS SELECT 1"), + ("MERGE", "MERGE INTO t USING s ON t.id=s.id WHEN MATCHED THEN UPDATE SET a=1"), + ("MSCK", "MSCK REPAIR TABLE x"), + ("REFRESH", "REFRESH TABLE x"), + ("ANALYZE", "ANALYZE TABLE x COMPUTE STATISTICS"), + ("CACHE", "CACHE TABLE x"), + ("UNCACHE", "UNCACHE TABLE x"), + ("VACUUM", "VACUUM x"), + ("OPTIMIZE", "OPTIMIZE x"), + ("CALL", "CALL system.proc()"), + ("GRANT", "GRANT SELECT ON TABLE x TO user"), + ("REVOKE", "REVOKE SELECT ON TABLE x FROM user"), + ("stacked statements", "SELECT 1; DELETE FROM x"), + ("empty", " "), + ("WITH without SELECT", "WITH x AS (DELETE FROM t)"), + ("INSERT VALUES", "INSERT INTO target VALUES (1)"), + ("INSERT OVERWRITE default blocked", "INSERT OVERWRITE target SELECT * FROM source"), +] -def expect_ok(name, sql): - try: - assert_select_or_insert(sql) - print(f" [OK] {name}") - except ValueError as e: - raise AssertionError(f"{name} should have passed but raised: {e}") +def main() -> int: + print("=== EXPECT PASS ===") + for name, sql in PASS_CASES: + result = validate_sql(sql) + assert result.status == "PASS", (name, sql, result) + assert assert_allowed_sql(sql) + print("[OK]", name) + + print("\n=== EXPECT PASS WITH OVERWRITE ===") + for name, sql in PASS_WITH_OVERWRITE: + result = validate_sql(sql, allow_overwrite=True) + assert result.status == "PASS", (name, sql, result) + assert assert_allowed_sql(sql, allow_overwrite=True) + print("[OK]", name) + + print("\n=== EXPECT FAIL ===") + for name, sql in FAIL_CASES: + result = validate_sql(sql) + assert result.status == "FAIL", (name, sql, result) + print("[OK]", name, "->", result.reason) + + print("\nALL TESTS PASSED") + return 0 -def expect_raise(name, sql): - try: - assert_select_or_insert(sql) - raise AssertionError(f"{name} should have raised but passed") - except ValueError: - print(f" [OK] {name}: raised as expected") - - -print("=== EXPECT_OK ===") -expect_ok("plain SELECT", "SELECT 1") -expect_ok("WITH ... SELECT", "WITH t AS (SELECT 1 AS x) SELECT x FROM t") -expect_ok("trailing semicolon", "SELECT 1 FROM dual;") -expect_ok("leading -- comment", "-- comment\nSELECT 1") -expect_ok("leading /* comment", "/* hi */ SELECT 1") -expect_ok("INSERT ... SELECT", "INSERT INTO t SELECT 1") -expect_ok("string 'drop'", "SELECT 'drop' AS note FROM dual") -expect_ok("quoted identifier", 'SELECT "drop" AS note FROM dual') -expect_ok("escaped quote", "SELECT 'don''t drop' AS note FROM dual") - -print() -print("=== EXPECT_RAISE ===") -expect_raise("DROP", "DROP TABLE loan_order") -expect_raise("UPDATE", "UPDATE loan_order SET status = 1") -expect_raise("DELETE", "DELETE FROM loan_order") -expect_raise("TRUNCATE", "TRUNCATE TABLE loan_order") -expect_raise("ALTER", "ALTER TABLE loan_order ADD COLUMN x INT") -expect_raise("MERGE", "MERGE INTO t USING s ON t.id = s.id") -expect_raise("MSCK", "MSCK REPAIR TABLE loan_order") -expect_raise("REFRESH", "REFRESH TABLE loan_order") -expect_raise("VACUUM", "VACUUM TABLE loan_order") -expect_raise("stacked ;", "SELECT 1; DROP TABLE loan_order") -expect_raise("empty", " ") -expect_raise("comment-hidden", "-- ok\nDROP TABLE loan_order") - -print() -print("ALL TESTS PASSED") +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/skills/pyspark-sql-guardrails/scripts/validate_sql.py b/skills/pyspark-sql-guardrails/scripts/validate_sql.py old mode 100755 new mode 100644 index 507dad1..db00a62 --- a/skills/pyspark-sql-guardrails/scripts/validate_sql.py +++ b/skills/pyspark-sql-guardrails/scripts/validate_sql.py @@ -1,60 +1,43 @@ -#!/usr/bin/env python3 -"""One-line SQL guard for the pyspark-sql-guardrails skill. +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- -Usage: - # Pipe SQL via stdin (preferred for multi-line SQL) - echo "SELECT ..." | python3 validate_sql.py - cat query.sql | python3 validate_sql.py - - # Or pass SQL as the first argument (single-line, quoting-friendly) - python3 validate_sql.py "SELECT 1" - -Exit code: - 0 -> PASS (SQL is allowed) - 1 -> FAIL (SQL violates the allowlist; reason printed to stdout) - -This script is the ONLY sanctioned way to validate a generated SQL -string before showing it to the user. It is intentionally short so it -can be invoked in a single Bash call from the assistant. -""" +from __future__ import annotations +import argparse +import json import sys -import os +from pathlib import Path -# Allow running from anywhere; resolve bundled sql_guard.py. -# The skill layout is: -# .claude/skills/pyspark-sql-guardrails/ -# ├── SKILL.md -# └── scripts/ -# ├── sql_guard.py (this script's sibling) -# ├── validate_sql.py -# └── test_sql_guard.py -SCRIPTS_DIR = os.path.dirname(os.path.abspath(__file__)) -sys.path.insert(0, SCRIPTS_DIR) - -from sql_guard import assert_select_or_insert, first_keyword # noqa: E402 +from sql_guard import validate_sql -def read_sql() -> str: - if len(sys.argv) > 1: - return sys.argv[1] +def read_sql(args: argparse.Namespace) -> str: + if args.file: + return Path(args.file).read_text(encoding="utf-8") + if args.sql: + return " ".join(args.sql) if not sys.stdin.isatty(): return sys.stdin.read() - print(__doc__, file=sys.stderr) - sys.exit(2) + raise SystemExit("Provide SQL as arguments, --file, or stdin") def main() -> int: - sql = read_sql() - try: - body = assert_select_or_insert(sql) - except ValueError as e: - print(f"FAIL - {e}") - return 1 - first = first_keyword(sql) - print(f"PASS - first keyword: {first}, body length: {len(body)}") - return 0 + parser = argparse.ArgumentParser(description="Validate Spark SQL against a safe allowlist.") + parser.add_argument("sql", nargs="*", help="SQL string; omitted when using --file or stdin") + parser.add_argument("--file", help="Read SQL from a UTF-8 file") + parser.add_argument("--json", action="store_true", help="Emit machine-readable JSON") + parser.add_argument("--allow-overwrite", action="store_true", help="Allow INSERT OVERWRITE ... SELECT") + args = parser.parse_args() + + result = validate_sql(read_sql(args), allow_overwrite=args.allow_overwrite) + if args.json: + print(json.dumps(result.as_dict(), ensure_ascii=False)) + elif result.status == "PASS": + print(f"PASS - first_keyword={result.first_keyword}") + else: + print(f"FAIL - first_keyword={result.first_keyword} reason={result.reason}") + return 0 if result.status == "PASS" else 1 if __name__ == "__main__": - sys.exit(main()) + raise SystemExit(main()) diff --git a/skills/pyspark-sql-pipeline/SKILL.md b/skills/pyspark-sql-pipeline/SKILL.md index edb88f2..e73dbcd 100644 --- a/skills/pyspark-sql-pipeline/SKILL.md +++ b/skills/pyspark-sql-pipeline/SKILL.md @@ -1,124 +1,188 @@ --- name: pyspark-sql-pipeline -description: 当需要把业务需求一路编排到"已被评审过、可执行的 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 字符串。 +description: 当用户希望从业务需求端到端生成一条经过校验和评审的 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(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` -- 输出中每个 `.` 都必须出现在 `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: -**Status:** NEED_USER_CONFIRMATION -**Open questions:** -1. -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: +**Status:** NEED_USER_CONFIRMATION +**Open questions:** +1. +**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`。 diff --git a/skills/requirements-analysis/SKILL.md b/skills/requirements-analysis/SKILL.md index 052ebb1..f998d96 100644 --- a/skills/requirements-analysis/SKILL.md +++ b/skills/requirements-analysis/SKILL.md @@ -1,273 +1,158 @@ --- name: requirements-analysis -description: 数据需求解读与拆解 skill。当用户提出"统计/计算/分析/取数/拉个数/做个报表/口径/指标/转化率/留存/活跃/DAU/GMV/漏斗/同环比"等数据类需求,或提供数据分析数据开发需求文档,或贴出表 schema 询问怎么写 SQL/Hive/Spark/取数逻辑时使用。本 skill 不直接写 SQL,而是先把需求拆解成可执行的 DRD(数据需求规格说明书):识别核心实体、解析表结构、推导字段、识别歧义口径、构建 join 关系,并主动向用户索要缺失信息直到口径完全明确。无论用户是否提供 schema、是否催促"直接给 SQL",都先走完澄清流程再交付。 +description: 当用户提出统计、取数、报表、指标、转化率、漏斗、留存、活跃、GMV、同环比、SQL/Hive/Spark/PySpark 数据开发需求时使用。先把模糊业务问题澄清成可交付的数据需求规格说明书 DRD,不直接写 SQL,不调用 MCP,不提交作业。 --- -# 数据需求解读 (Requirements Analysis) +# Requirements Analysis(需求澄清) -## 你的角色 +## 目标 -你是数据开发工程师面前的"需求接口人"。用户(通常是产品、运营、数据分析师,甚至业务方)抛过来一句模糊的统计需求,你的任务是:把它翻译成另一位数据开发同学拿到就能直接写 SQL 的规格说明书(DRD)。 +把一句可能含糊的数据需求整理成正式 DRD。重点不是快点写 SQL,而是先把统计对象、时间窗口、指标分子分母、维度、表字段、join 和输出要求固定下来,避免“SQL 看起来对但数字口径错”。 -这意味着你的产出**不是 SQL**,而是一份口径明确、字段对齐、join 关系清晰的需求文档。 +本 skill 的输出是 `requirements_output`,下游只能基于它推进,不得重新解释业务口径。 -## 为什么先澄清再动手 +## 工作顺序 -数据需求最大的坑不在 SQL 写错,而在口径理解错。"近 30 天新用户转化率"这一句话里,"近 30 天"、"新用户"、"转化"三个词每一个都至少有 3 种合理解释。如果不先澄清就动手,最后跑出来的数字可能和业务方期望差几个数量级,返工成本远高于多问几轮。 +1. 复述业务目标:这个统计用于评估什么,最终交付什么。 +2. 定义输出粒度:结果一行代表什么维度组合。 +3. 明确统计对象:用户、订单、交易、设备、账户等。 +4. 明确时间口径:自然日、T-1 完整日、滚动窗口、时区、事件时间字段。 +5. 明确指标口径:分子、分母、去重键、过滤条件、空值策略。 +6. 明确维度口径:渠道、城市、产品、日期等分组字段来自哪张表。 +7. 明确数据来源:候选表、候选字段、字段含义、状态枚举、金额单位。 +8. 明确关联关系:主表、维表、join key、join 类型、1:N 风险。 +9. 明确结果要求:输出列、排序、是否落表、是否允许 overwrite。 +10. 如果任一关键口径不明确,停止并提问;如果已明确,输出 DRD 和 `requirements_output`。 -所以本 skill 的核心动作是:**发现歧义 → 主动提问 → 等待确认 → 再推进**。宁可多问一轮,不要默认假设。 +## 必问歧义 -## 输入分流 +- “近 30 天”是自然日、T-1 完整日,还是滚动到当前时刻。 +- “新用户”是首次注册、首次登录、首次下单、首次激活还是首次付费。 +- “成功”对应哪些状态码,退款、撤销、部分成功是否计入。 +- “转化率”的分子、分母、事件顺序、去重键和转化窗口是什么。 +- 分组维度来自事实表还是维表;维表缺失时是否保留 NULL。 +- 多表 join 后是否会产生重复行,是否要先去重或聚合。 +- 输出是临时查询还是要写入目标表。 -接到需求后,先判断属于哪种情况: +## 时间口径硬规则 -### 情况 A:只有业务需求,没有表结构 - -例如:"统计近 30 天新用户转化率" - -主动索要: -- 涉及哪些表(让用户列出表名即可) -- 这些表的 schema(字段名 + 类型) -- 关键字段的业务含义(特别是状态、时间、金额类) -- 指标口径定义("新用户"、"转化"分别指什么) - -### 情况 B:需求 + 单表/少量表 schema - -理解业务逻辑 → 推导所需字段 → 检查是否缺关键字段 → 向用户确认口径。 - -### 情况 C:需求 + 多表 schema - -额外要做: -- 区分主表 / 维表 / 事实表 -- 推导 join 关系(on 哪些字段、内连接还是左连接) -- 检查 join 字段在两边是否都存在、类型是否一致 - -## 核心动作(每轮对话都要走) - -### 1. 需求拆解 - -按这 7 个维度把用户的话拆开,**任何一个不明确就要追问**: +对“近 N 天按 T-1 完整日”,统一解释为左闭右开窗口: ```text -业务目标:为什么要看这个数? -统计对象:是用户、订单、商品、还是别的? -统计范围:哪些数据进入统计?哪些被过滤? -统计时间:从什么时候到什么时候?按自然日还是滚动? -统计维度:按什么分组?(渠道 / 城市 / 产品线 / 时间粒度) -统计指标:算 count? sum? 比率?分子分母分别是什么? -输出形式:一个数?一张表?带哪些列? +[start, end) = [current_date() - N days, current_date()) ``` -### 2. 识别核心实体 +Spark/Hive SQL 应表达为: -从需求里抽出涉及的业务实体(用户、订单、商品、商户、设备、渠道、活动、贷款、客户、账户、交易……),并推测每个实体对应哪张表。 +```sql +event_time >= date_sub(current_date(), N) +AND event_time < current_date() +``` -如果用户没给表名,**列出你的猜测**让用户确认或补充。 +不要使用下面写法表示 timestamp 完整日窗口: -### 3. Schema 解析(拿到表结构后) +```sql +event_time BETWEEN date_add(current_date(), -N) AND date_add(current_date(), -1) +``` -对每张表,要分析清楚: +因为它很容易只覆盖到结束日期的 00:00:00,漏掉 T-1 白天的数据。 -| 维度 | 含义 | -|------|------| -| 表用途 | 这张表记录什么业务事件/状态 | -| 主键 | 唯一标识一行的字段 | -| 业务主键 | 业务上的唯一键(可能和主键不同) | -| 时间字段 | create_time / update_time / 业务时间,用哪个? | -| 状态字段 | 订单状态、用户状态等枚举字段,可取值是什么 | -| 维度字段 | 可用于分组的字段 | -| 金额字段 | 单位是元还是分?是否含税? | +## 转化口径硬规则 -### 4. 推导所需字段 +转化类需求必须写清: -对照需求拆解,列出每张表要用到的字段。**显式标注哪些字段你假设存在但还没确认**,让用户打勾或补充。 +- 分母人群,例如“窗口内首次登录的新用户”。 +- 分子人群,例如“该新用户注册后存在成功付费”。 +- 事件顺序,例如 `pay_time >= register_time`。 +- 后置事件是否也需要在统计窗口内。 +- 去重键,例如 `user_id`。 -### 5. 识别缺失字段 +如果用户说“注册后付费”,必须把“后”写进 DRD,不允许只写“注册且付费”。 -常见缺失:用户 ID、订单 ID、时间字段、状态字段、金额字段、渠道字段、产品字段。如果发现关键字段缺失,**直接指出**,不要绕开它继续推进。 +## 正式 DRD 输出 -### 6. 识别需求歧义(重要!) - -这一步是本 skill 的核心价值。**对每一个模糊词,主动列出可能的解释让用户选择**。常见歧义点: - -**时间口径** -- "近 30 天" → ① 自然日 [今天-30, 今天);② T-30 不含今天 [昨天-30, 昨天];③ 最近 30 个完整日;④ 滚动 30 天 -- "本月" → 自然月还是结算月? -- "凌晨数据归哪天" → 0 点切分还是 4 点切分? - -**用户口径** -- "新用户" → 注册用户 / 首次下单用户 / 首次登录用户 / 首次激活用户 -- "活跃用户" → 登录过 / 操作过 / 下单过 -- "去重" → 按 user_id 还是按设备 ID 去重? - -**状态/转化口径** -- "转化" → 下单 / 支付成功 / 放款 / 激活 / 完成首单 -- "成功订单" → 状态码具体是哪些?是否包含部分退款? -- "支付金额" → 应付 / 实付 / 实收(扣手续费后) - -**关联口径** -- "用户的订单" → 创建人 / 收货人 / 实际付款人? -- 多表 join → inner / left / 是否要去重防止笛卡尔积? - -输出形式: +当关键口径已明确时,必须先输出可读 DRD,再输出 YAML 摘要。 ```markdown -## 待确认问题 +## 需求说明 -1. "新用户"指的是? - - [ ] A. 在统计窗口内首次注册的用户 - - [ ] B. 在统计窗口内首次下单的用户 - - [ ] C. 其他(请说明) +### 业务目标 +... -2. "近 30 天"是? - - [ ] A. [今天-30, 今天) 含今天 - - [ ] B. [昨天-30, 昨天] T-1 口径 - - [ ] C. 其他 +### 统计口径 +- 统计对象:... +- 时间范围:近 N 天 T-1 完整日,`event_time >= date_sub(current_date(), N)` 且 `event_time < current_date()` +- 维度:... +- 指标:... + +### 涉及表 +| 表名 | 用途 | 粒度 | 使用字段 | +|---|---|---|---| + +### 表关系 +| 左表字段 | 右表字段 | Join 类型 | 风险 | +|---|---|---|---| + +### 字段清单 +| 表名 | 字段 | 用途 | 角色 | +|---|---|---|---| + +### 实现逻辑 +1. ... + +### 待确认问题 +无,或列出阻塞问题。 ``` -### 7. 构建 join 关系 - -多表场景下,明确画出关联: - -``` -user_info.user_id = order_info.user_id (left join, 左表为主) -order_info.order_id = pay_info.order_id (inner join) -``` - -要标明: -- 关联字段及类型是否一致 -- join 类型(inner / left / right / full) -- 是否会产生一对多导致重复计算 - -### 8. 确认机制 - -**严格禁止**: -- 自行脑补未说明的业务逻辑 -- 默认字段含义(哪怕字段名看上去很标准) -- 默认时间口径、用户口径、转化口径 -- 在关键问题没确认前就给出最终 SQL - -每轮回复都要包含三块: -- ✅ **已确认**:到目前为止双方对齐的内容 -- ❓ **待确认**:还需要用户回答的问题(编号列出,方便用户对照回答) -- 📋 **当前推导**:基于已知信息你推出的字段清单 / join 关系(标注哪些是假设) - -## 完成判定 - -只有以下 7 项**全部**明确,才算需求分析完成、可以交付 DRD: - -- [ ] 业务逻辑明确 -- [ ] 指标口径明确(分子分母、聚合方式) -- [ ] 时间范围明确(窗口定义、时区、是否含端点) -- [ ] 维度明确(group by 哪些字段) -- [ ] 表来源明确(每个数据来自哪张表) -- [ ] join 关系明确(关联键、关联类型) -- [ ] 字段明确(每个字段都已被用户确认存在) - -任何一项打问号,都继续走澄清流程,不要交付最终 DRD。 - -## 最终交付:DRD 模板 - -当 7 项全部确认后,按这个模板输出: - -```markdown -# 数据需求规格说明书 (DRD) - -## 一、业务目标 -(1-2 句话说明这个数据是做什么用的、给谁看、支撑什么决策) - -## 二、统计口径 -- 统计对象:xxx -- 时间范围:xxx(精确到时区、端点) -- 过滤条件:xxx -- 聚合维度:xxx -- 指标定义: - - 指标 A = 分子 / 分母,分子定义为 xxx,分母定义为 xxx - -## 三、涉及表 - -### 表 1: -- 用途:xxx -- 粒度:一行代表 xxx -- 使用字段: - - `field_a`: 含义 - - `field_b`: 含义 - -### 表 2: -(同上) - -## 四、表关联关系 -``` -table_a.id = table_b.a_id (left join) -``` -说明:xxx - -## 五、最终字段清单 - -| 表名 | 字段名 | 用途 | 备注 | -|------|--------|------|------| -| table_a | id | 主键 | | -| table_a | created_at | 时间过滤 | UTC+8 | - -## 六、实现逻辑(伪代码层面) -1. 从 table_a 取 [time_range] 内的数据,过滤 status = xxx -2. 按 user_id left join table_b -3. group by xxx,聚合 xxx -4. 输出列:xxx - -## 七、待开发 SQL 所需信息 -全部已确认 ✓ -``` - -## 风格提醒 - -- 每轮回复都用结构化 markdown,不要长段落散文 -- 提问要给选项(A/B/C),降低用户回答成本 -- 字段名、表名用反引号 `code` 包起来 -- 时间相关问题特别仔细,时区、端点、自然日 vs 滚动这三件事最容易翻车 -- 如果用户催"直接给 SQL",温和坚持:"为了避免数字跑出来不对,先把这 N 个口径敲定,几分钟就能确认完,然后直接交付准确的 SQL" -- 始终记得:你的产出物是 DRD,不是 SQL - -## 例外情况 - -如果用户**明确表示**"我已经想清楚口径了,不用再问,按以下定义直接生成 SQL",并且把所有 7 项完成判定都写明了,那么可以跳过澄清直接进入 DRD/SQL 阶段。但即便这样,也要在回复里复述一遍你理解的口径,让用户最后过一眼。 - ---- - -## 下游契约(Pipeline Handoff) - -本 skill 是 PySpark SQL 流水线的**第一站**。澄清完成后,必须把结构化结果以标准化形式交付给 `metadata-validator`,由其继续推进。 - -**Required next step:** `metadata-validator` — 它会递归扫描工作区元数据文件并校验候选表/字段/关联/时间字段是否真实存在。 - -在交付 DRD 之前,必须同时输出以下标准化结构(与第七节 DRD 内容一致): +## 结构化输出 ```yaml requirements_output: - business_goal: "<一句话业务目标>" - business_logic: "<一段话业务逻辑,含统计对象/时间窗口/分组/指标分子分母>" - candidate_tables: - - # 用户提及或根据业务推断的物理表 - candidate_fields: - - name: # 业务字段名(如 客户号 / 贷款金额) - suggested_table: # 建议所在表 - suggested_column: # 建议字段名(如未确认可省略) - role: + status: READY_FOR_METADATA | NEED_USER_CONFIRMATION + business_goal: "一句话业务目标" + grain: "结果一行代表什么" time_window: - field_hint: # 用户提到的时间字段 - window: "<窗口描述,如 近 30 天 / 2026-05-01 ~ 2026-05-31>" - join_hints: # 用户已说明的关联(如有) - - left_table: - left_field: - right_table:
- right_field: - join_type: - status: + description: "时间窗口" + timezone: "时区" + event_time_hint: "候选时间字段" + lower_bound_sql: "date_sub(current_date(), N)" + upper_bound_sql: "current_date()" + boundary: "left_closed_right_open" + metrics: + - name: "指标名" + formula: "业务公式" + numerator: "分子" + denominator: "分母" + dedupe_key: "去重键" + order_constraint: "例如 pay_time >= register_time" + null_policy: "空值处理" + dimensions: + - "分组维度" + filters: + - "过滤条件" + candidate_tables: + - business_role: "事实表/维表/码表/结果表" + table_hint: "候选物理表" + candidate_fields: + - business_name: "业务字段名" + role: dimension | metric | filter | time | join_key | output + table_hint: "候选表" + column_hint: "候选字段" + required: true + join_hints: + - left: "表A.字段" + right: "表B.字段" + type_hint: left | inner | right | full | unknown + cardinality_hint: "1:1 | 1:N | N:1 | N:N | unknown" + output_columns: + - "输出列" + write_intent: + mode: query_only | insert_into | insert_overwrite | unknown + target_table: "可为空" + overwrite_confirmed: false + pending_questions: [] ``` -**契约约束:** +## 停止条件 -- `status: READY_FOR_VALIDATION` 只有在第 7 项完成判定全部打勾时才允许输出。 -- `status: NEED_USER_CONFIRMATION` 时,必须同时在回复里保留第七节的 DRD「待确认问题」清单,`metadata-validator` **不得**在澄清完成前介入。 -- `candidate_tables` / `candidate_fields` 是 `metadata-validator` 的**唯一**输入;不要让 `metadata-validator` 自行反推业务字段。 -- 本 skill 不直接进入 SQL 生成。任何“直接给 SQL”的要求都必须先走完澄清 → DRD → `metadata-validator` → `logic-planner` → `sql-context-builder` 这条链路。 +只要 `pending_questions` 非空,就输出 `NEED_USER_CONFIRMATION`,不要进入元数据校验,更不要写 SQL。 diff --git a/skills/requirements-analysis/evals/evals.json b/skills/requirements-analysis/evals/evals.json index faacb9b..851e849 100644 --- a/skills/requirements-analysis/evals/evals.json +++ b/skills/requirements-analysis/evals/evals.json @@ -2,25 +2,19 @@ "skill_name": "requirements-analysis", "evals": [ { - "id": 1, - "name": "vague-no-schema", - "prompt": "帮我统计下近30天新用户转化率", - "expected_output": "应主动索要相关表/schema/字段说明,并对'近30天'、'新用户'、'转化'这三个词分别提出多选式澄清问题。不应直接生成 SQL 或假设口径。", - "files": [] + "id": "ambiguous-conversion-rate", + "prompt": "帮我看下最近一个月新用户转化率", + "expected_behavior": "必须追问最近一个月的时间边界、新用户定义、转化动作、分子分母、去重键和候选表字段;不能直接写 SQL。" }, { - "id": 2, - "name": "schema-with-ambiguity", - "prompt": "统计上周每个渠道的支付成功订单数和支付金额。表结构如下:\n\norder_info(\n order_id bigint,\n user_id bigint,\n channel_id int,\n create_time timestamp,\n pay_time timestamp,\n status int, -- 1待支付 2已支付 3已发货 4已完成 5已退款 6部分退款\n amount decimal(18,2), -- 应付金额\n paid_amount decimal(18,2) -- 实付金额\n)\n\nchannel_dim(\n channel_id int,\n channel_name varchar\n)", - "expected_output": "应识别出'支付成功'状态码歧义(2/3/4/6 是否都算)、'支付金额'是 amount 还是 paid_amount、'上周'是自然周还是滚动 7 天、退款订单是否计入。应提议 left join channel_dim。不应直接生成 SQL。", - "files": [] + "id": "orders-with-schema", + "prompt": "统计上周各渠道支付成功订单数和金额,并给出 order_info/channel_dim schema", + "expected_behavior": "必须确认上周口径、支付成功状态、金额字段、退款是否计入、join 类型和输出粒度。" }, { - "id": 3, - "name": "multi-table-loan", - "prompt": "我要算下贷款产品的首贷转化漏斗,从注册到首次申请到首次放款的转化率,按月看趋势。表:\n\nuser_register(user_id bigint, register_time timestamp, channel varchar)\nloan_apply(apply_id bigint, user_id bigint, apply_time timestamp, product_id int, apply_amount decimal)\nloan_disburse(disburse_id bigint, apply_id bigint, disburse_time timestamp, disburse_amount decimal, status int)\nproduct_dim(product_id int, product_name varchar, product_type varchar)", - "expected_output": "应识别多表 join 关系(user → apply 按 user_id, apply → disburse 按 apply_id),追问首贷定义(按用户首次还是按产品首次)、放款 status 含义、'按月看'是注册月还是事件发生月、是否限定产品类型、漏斗分母是否包含没注册成功的用户。", - "files": [] + "id": "loan-funnel", + "prompt": "做注册到申请到放款的贷款漏斗,按月看", + "expected_behavior": "必须确认月份归属、首贷定义、放款成功状态、漏斗分母、用户去重和 user/apply/disburse 关联关系。" } ] -} +} \ No newline at end of file diff --git a/skills/sql-context-builder/SKILL.md b/skills/sql-context-builder/SKILL.md index d03e8fd..fbc7925 100644 --- a/skills/sql-context-builder/SKILL.md +++ b/skills/sql-context-builder/SKILL.md @@ -1,178 +1,140 @@ --- name: sql-context-builder -description: 当需要组装 SQL 生成器与评审器共同消费的、唯一可信的 SQL Context 时使用。触发词包括"build SQL context"、"固化 join key"、"table aliases"、"列出来自哪张表"、"准备写 SQL"、"context 已经齐了"。消费已校验的 Logic Plan 与 Field Mapping,产出带有稳定 alias、固化 join key、显式 fact/dimension 标记的标准化 SQL Context。SQL 生成器**不得**再重新推导其中任何一项。 +description: 当逻辑计划已经完成,需要固定 SQL 生成和 SQL review 共同依赖的 alias、字段来源、join key、指标、维度、过滤、时间窗口、事实表/维表上下文时使用。消费 logic_plan 和 validation_result,不写 SQL,不调用 MCP。 --- -# SQL Context Builder(SQL 上下文构建器) +# SQL Context Builder(SQL 上下文构建) -## 概述(Overview) +## 目标 -把 SQL 生成器本来要重复发明的所有决策一次性固化下来:alias、join key、字段来源、类型、角色标签(fact / dimension / filter)、时间列,都在这里只此一次地钉死。输出是 `sql-review` 及任何 SQL 输出步骤的**唯一可信输入**。 +把 SQL 生成时可能被临时猜测的内容提前固定下来。下游只能使用本 skill 输出的 alias、字段、join、指标、过滤、时间窗口和输出列。 -核心原则:**任何不在 SQL Context 中的字段,下游 skill 视同不存在。** 如果下游 skill 还需要新内容,请求必须先回退到 `metadata-validator` 和 `logic-planner`。 +## 输入 -## 何时使用(When to Use) - -**使用场景:** - -- `logic-planner` 返回状态 `PLANNED` -- 用户即将写 Spark SQL,需要一份规范化的 context(alias、join key、字段来源) -- 同一份逻辑要产出多版 SQL 草稿 — context 保证它们保持一致 - -**不要使用场景:** - -- `logic-plan.status` 为 `NEED_USER_CONFIRMATION` — 回到 `logic-planner` -- SQL 极简单(单表无 join) — 仍然建议构建 context,除非用户明确跳过 -- 任务是把 plan 翻译成 SQL — 那是另一回事,不是本 skill 的职责 - -## 输入(Inputs) - -1. 来自 `logic-planner` 的 `logic_plan`(状态 `PLANNED`) -2. 来自 `metadata-validator` 的 `field_mapping`(状态 `VALIDATED`) -3. 原始 `validation_result`(用于表级元数据:类型、来源) - -## 工作流(Workflow) - -```dot -digraph sql_context_builder { - "读取 logic_plan + field_mapping + validation_result" [shape=box]; - "为每张表分配稳定 alias" [shape=box]; - "固化 join key(left_alias.field = right_alias.field)" [shape=box]; - "把每个字段解析为 (alias, column, type, source_file)" [shape=box]; - "为每个字段打标签:dimension | metric | filter | time | join_key" [shape=box]; - "标记事实表 / 维度表" [shape=box]; - "校验 Logic Plan 的字段在 context 中全部存在" [shape=box]; - "所有必需字段都齐全?" [shape=diamond]; - "输出 SQL Context" [shape=box]; - "输出 SQL Context + missing_fields" [shape=box]; - - "读取 logic_plan + field_mapping + validation_result" -> "为每张表分配稳定 alias"; - "为每张表分配稳定 alias" -> "固化 join key(left_alias.field = right_alias.field)"; - "固化 join key(left_alias.field = right_alias.field)" -> "把每个字段解析为 (alias, column, type, source_file)"; - "把每个字段解析为 (alias, column, type, source_file)" -> "为每个字段打标签:dimension | metric | filter | time | join_key"; - "为每个字段打标签:dimension | metric | filter | time | join_key" -> "标记事实表 / 维度表"; - "标记事实表 / 维度表" -> "校验 Logic Plan 的字段在 context 中全部存在"; - "校验 Logic Plan 的字段在 context 中全部存在" -> "所有必需字段都齐全?"; - "所有必需字段都齐全?" -> "输出 SQL Context" [label="yes"]; - "所有必需字段都齐全?" -> "输出 SQL Context + missing_fields" [label="no"]; -} +```yaml +logic_plan: + status: PLANNED +validation_result: + status: VALIDATED + field_mapping: {} + joins: [] ``` -## 规则(Rules) +## 工作顺序 -### Alias 分配 +1. 给每张表分配唯一、稳定、短小的 alias。 +2. 标记事实表、维表、码表和目标表。 +3. 把每个业务字段解析成 `alias.column`。 +4. 固定 join 谓词、join 类型和基数假设。 +5. 固定指标表达式、维度字段、时间字段、分区字段和过滤字段。 +6. 固定时间窗口 SQL 片段,尤其是 T-1 完整日边界。 +7. 固定输出列顺序和别名。 +8. 检查 logic plan 中引用的字段都能在 context 中找到。 +9. 对派生字段记录来源步骤和依赖字段。 -- 一张表一个 alias,绝不重用 -- 事实表(承载主指标的表)优先拿到短 alias(`o`、`f`、`t1`) -- 维度表用有意义的 alias(`c` 代表 customer,`p` 代表 product,`ch` 代表 channel) -- alias **在此处固化**,SQL 生成器与评审器不得再发明新 alias +## 时间窗口上下文 -### Join Key 固化 +当存在时间过滤时,必须输出 `time_windows`。对“近 N 天按 T-1 完整日”: -- 对每条 `logic_plan.steps[kind=join]`,固化精确的谓词: - - `left_alias.col1 = right_alias.col2` -- 包含 join 类型(inner / left / right / full / semi / anti) -- 沿用 Logic Plan 中的 cardinality 假设 -- 如果两张表有多种可能的 join key,在此处固化为所选的那一个,其他候选不再纳入 SQL 范畴 +```yaml +time_windows: + - name: t_minus_1_complete_days + event_field: fl.login_time + lower_bound_sql: date_sub(current_date(), N) + upper_bound_sql: current_date() + predicate_sql: fl.login_time >= date_sub(current_date(), N) AND fl.login_time < current_date() + boundary: left_closed_right_open +``` -### 字段解析(Field Resolution) +下游 SQL 必须复用 `predicate_sql`,不得临时改成 `BETWEEN`。 -对 Logic Plan 用到的每个字段,记录: +## 转化指标上下文 -- `alias`(字段所在表) -- `column`(物理列名) -- `type`(取自 `validation_result`) -- `source_file`(元数据文件来源 — 用于追溯) -- `role`:取自 `dimension` | `metric` | `filter` | `time` | `join_key` | `derived` 之一 +转化率必须固定分母、分子和顺序约束: -一个字段可以有多个 role(例如 `customer_id` 同时是 `join_key` 和 `dimension`),全部列出。 +```yaml +metrics: + - name: conversion_rate + numerator: converted_user_cnt + denominator: new_user_cnt + order_constraint: pay_time >= register_time + null_policy: denominator_zero_returns_0 + exists_after_anchor_event: true +``` -### Fact vs Dimension 标记 - -- 承载主要事件/度量的表为 `fact` -- 起丰富作用的表(customer、product、channel、city)为 `dimension` -- 一条 query 只能有一个 fact。如果 Logic Plan 隐含多张事实表,作为 `pending_question` 抛出 — 那是规划问题,不是 context 问题 - -### 覆盖度校验 - -- `logic_plan.steps` 中的每一步引用的字段都必须存在于本 context -- 每条 `field_mapping` 都必须出现在本 context -- 每个 `group_by`、`order_by` 字段都必须有标签 -- `aggregate` 中的每个指标都必须有 `metric` 标签 -- `filter` 中的每条谓词都必须引用 `filter` 或 `time` 标签的字段 - -## 输出契约(Output Contract) +## 输出格式 ```yaml sql_context: + status: CONTEXT_READY | NEED_USER_CONFIRMATION aliases: - - alias: o - table: loan_order - role: fact - source_file: metadata/loan_order.json - - alias: c - table: customer_info - role: dimension - source_file: metadata/customer_info.md - joins: - - left: { alias: o, column: customer_id } - right: { alias: c, column: customer_id } - type: left - cardinality_assumption: 1:N + - alias: fl + table: user_login + role: fact | dimension | lookup | target + source_file: "元数据来源" fields: - - alias: o - column: customer_id - type: string - role: [join_key] - source_file: metadata/loan_order.json - - alias: o - column: loan_amount - type: decimal(18,2) - role: [metric] - source_file: metadata/loan_order.json - - alias: o - column: apply_time - type: timestamp - role: [time, filter] - source_file: metadata/loan_order.json - - alias: c - column: city - type: string - role: [dimension] - source_file: metadata/customer_info.md - metrics: - - alias: total_loan - expr: "SUM(o.loan_amount)" - agg: sum - - alias: customer_cnt - expr: "COUNT(DISTINCT o.customer_id)" - agg: count_distinct - dimensions: ["c.city"] - time_grain: - column: o.apply_time - granularity: day - filters: - - "o.status = 'SUCCESS'" - missing_fields: [] - status: + - id: fl.user_id + alias: fl + table: user_login + column: user_id + type: bigint + role: [join_key, metric] + nullable: false + source_file: "元数据来源" + joins: + - left: fl.user_id + right: ur.user_id + type: left + cardinality: "N:1" + source_step: S3 + time_windows: [] + metrics: [] + dimensions: [] + filters: [] + outputs: [] + write_intent: + mode: query_only | insert_into | insert_overwrite + target_table: "可为空" + overwrite_confirmed: false + unresolved: [] ``` -## 禁止行为(Forbidden Behaviors) +## 用户可见摘要 -| 自我说服 | 现实 | -|---|---| -| "alias 让 SQL 生成器自己选" | alias 就在此处固化。下游重新起 alias 会破坏 `sql-review`。 | -| "join key 太显然,跳过固化" | 整个 skill 的意义就在于此 — 把 join key 钉死。 | -| "字段类型无所谓,SQL 里都是字符串" | 类型决定了聚合、cast、分区裁剪的策略。 | -| "派生列到 SQL 阶段再发明就行" | 每个派生列都必须在 Logic Plan 里就出现。 | -| "多张事实表?全 join 上就行" | 多事实是 Logic Plan 的问题,要暴露并停止。 | +除了 YAML,必须输出: -## 完成判定(Completion Criteria) +```markdown +## SQL 上下文 +### 表别名 +| Alias | 表名 | 角色 | +|---|---|---| -仅当**以下全部**成立时,`status: READY`: +### 字段来源 +| 字段 | 来源 | 用途 | +|---|---|---| -- Logic Plan 的每一步在本 context 中都有对应条目 -- `missing_fields` 为空 -- 每个指标都有 `agg` 和 `expr` -- 每个 join 都有固化的 key 和 type -- 时间粒度已设定 +### Join 关系 +| 左字段 | 右字段 | 类型 | 基数风险 | +|---|---|---|---| + +### 指标公式 +| 指标 | 分子 | 分母 | 空值策略 | +|---|---|---|---| + +### 时间窗口 +| 字段 | 下界 | 上界 | 边界 | +|---|---|---|---| +``` + +## 规则 + +- alias 一旦生成,下游不能重新命名。 +- context 中不存在的字段,下游视为不存在。 +- join key 只能来自元数据校验和逻辑计划。 +- 指标表达式必须能追溯到 logic plan 的 aggregate 步骤。 +- 输出列必须来自 `fields`、`metrics` 或明确派生字段。 +- 如果 SQL 生成需要新字段,回到 metadata-validator 或 logic-planner。 + +## 停止条件 + +只在 alias、字段、join、指标、时间字段、过滤字段、输出列和写入意图都完整时输出 `CONTEXT_READY`。 diff --git a/skills/sql-review/SKILL.md b/skills/sql-review/SKILL.md index 11f4ce2..5828c00 100644 --- a/skills/sql-review/SKILL.md +++ b/skills/sql-review/SKILL.md @@ -1,127 +1,116 @@ --- name: sql-review -description: 对执行前生成好的 PySpark SQL 字符串做静态评审时使用。触发词包括"review this SQL"、"SQL 评审"、"静态审查"、"risk level"、"check join cardinality"、"数据膨胀"、"重复聚合"、"分区裁剪"。消费 SQL 字符串及其产出的 SQL Context,产出按风险分级的评审结果(HIGH / MEDIUM / LOW)以及具体修复建议。**不修改 SQL** — 只报告发现。 +description: 当一条 Spark SQL 已经生成、通过 guardrails 并准备展示或执行前使用。它基于 sql_context 静态检查字段引用、alias、join key、数据膨胀、重复聚合、时间过滤、转化顺序、GROUP BY、空值、类型和写入风险,只报告问题,不直接改 SQL。 --- -# SQL Review(SQL 评审) +# SQL Review(SQL 静态评审) -## 概述(Overview) +## 目标 -对一条生成好的 PySpark SQL 字符串做执行前的静态评审。评审者**不**改写 SQL — 它只输出发现项、严重程度和建议修复。由 SQL 的作者决定采纳哪些。 +在 SQL 执行前做质量检查,发现可能导致数字错误、执行失败或性能风险的问题。评审只报告风险和建议,不负责直接改写 SQL。 -核心原则:**对照 SQL Context 评审,而不是凭直觉。** SQL 中引用的每一个字段、alias、join key 都必须存在于产出它的 SQL Context。如果不存在,这就是一项 `HIGH` 风险(SQL 绕过了契约)。 +## 输入 -## 何时使用(When to Use) - -**使用场景:** - -- 一条 SQL 字符串已经从 `sql_context` 生成出来,即将执行 -- 用户粘贴一条 SQL 字符串,要求评审、优化或风险评估 -- 代码改动为 PySpark 脚本新增了 `spark.sql(...)` 调用 - -**不要使用场景:** - -- 这条 SQL 不是 Spark SQL(例如纯 Postgres、HiveQL 在 Spark 之外)— 至少要把 Spark 特有的检查替换成对应引擎的等价检查,并标注出来 -- 是一次性、不会被复用的交互查询 — 用户可以主动选择跳过 - -## 输入(Inputs) - -1. 要评审的 **SQL 字符串** -2. 产出这条 SQL 的 **SQL Context**(如果缺失,需要显式声明"该 SQL 是不可信 / 临时拼凑的") - -如果 SQL Context 缺失,抛出一项 `HIGH` 风险:*"SQL 没有对应的 SQL Context;下游 skill 无法保证字段 / join 的正确性。"* 然后继续做结构性检查。 - -## 严重程度模型(Severity Model) - -| 等级 | 含义 | 处置 | -|---|---|---| -| `HIGH` | 会导致数字错误、执行失败或违反策略 | 阻断执行,必须修。 | -| `MEDIUM` | 大概率对,但低效、脆弱或大规模场景有风险 | 上线前修,开发环境可放过。 | -| `LOW` | 风格 / 最佳实践小瑕疵 | 可选修。 | - -按以下顺序输出发现项:`HIGH` → `MEDIUM` → `LOW`。任一 `HIGH` 让整体评审结论为 `FAIL`。 - -## 工作流(Workflow) - -```dot -digraph sql_review { - "解析 SQL(逻辑层面,不只是正则)" [shape=box]; - "对照 SQL Context 交叉校验 alias、字段、join key" [shape=box]; - "检测多语句 / 禁用 verb(防护栏)" [shape=box]; - "分析 join 基数与膨胀风险" [shape=box]; - "检测同一度量的重复聚合" [shape=box]; - "检测全表扫描与缺失分区过滤" [shape=box]; - "校验 GROUP BY 与 SELECT 一致" [shape=box]; - "检查 join key 与度量的空值处理" [shape=box]; - "检查字段类型与运算的匹配" [shape=box]; - "汇总成评审结论" [shape=box]; - - "解析 SQL(逻辑层面,不只是正则)" -> "对照 SQL Context 交叉校验 alias、字段、join key"; - "对照 SQL Context 交叉校验 alias、字段、join key" -> "检测多语句 / 禁用 verb(防护栏)"; - "检测多语句 / 禁用 verb(防护栏)" -> "分析 join 基数与膨胀风险"; - "分析 join 基数与膨胀风险" -> "检测同一度量的重复聚合"; - "检测同一度量的重复聚合" -> "检测全表扫描与缺失分区过滤"; - "检测全表扫描与缺失分区过滤" -> "校验 GROUP BY 与 SELECT 一致"; - "校验 GROUP BY 与 SELECT 一致" -> "检查 join key 与度量的空值处理"; - "检查 join key 与度量的空值处理" -> "检查字段类型与运算的匹配"; - "检查字段类型与运算的匹配" -> "汇总成评审结论"; -} +```yaml +sql: "待评审的 Spark SQL" +sql_context: + status: CONTEXT_READY +sql_guard: + status: PASS +logic_plan: + status: PLANNED ``` -## 检查项(Checks,全部执行) +如果没有 `sql_context`,也可以继续做基础结构检查,但必须报告一个 HIGH 风险:缺少上下文,无法保证字段和 join 正确。 -1. **引用完整性(Reference integrity)** — SQL 中的每个 `.` 都必须存在于 SQL Context。不匹配即 `HIGH` -2. **Join key 匹配** — 每个 `ON` 谓词都必须对得上 SQL Context 中已固化的 join。凭空发明的 join 算 `HIGH` -3. **禁用 verb / 多语句** — 交给 `pyspark-sql-guardrails`。任何 DDL/DML/维护 verb,或任何 `;` 堆叠的语句,均为 `HIGH` -4. **Join 基数风险** — 对每个 join,如果 SQL Context 标记基数为 `1:N` 或 `N:N` 且没有去重,分别抛 `MEDIUM`(1:N)或 `HIGH`(N:N) -5. **重复聚合** — 同一度量被聚合两次(例如子查询里一次、外层再一次)且非有意重新聚合,抛 `HIGH` -6. **分区 / 时间过滤** — 对任何有已知时间字段的事实表,缺少对该字段的 `WHERE` 过滤即为 `HIGH`(全表扫描风险) -7. **GROUP BY 完整性** — SELECT 中每个未聚合的列都必须出现在 `GROUP BY` 中。不匹配为 `HIGH` -8. **空值处理** — 可空 join key 上做 `INNER JOIN` 会静默丢行(`MEDIUM`);可空度量未用 `COALESCE` 聚合,可能产出 `NULL` 进而破坏下游 `WHERE ... > 0`(`MEDIUM`) -9. **类型不匹配** — 字符串列未做 cast 就参与数值聚合(`MEDIUM`);日期字符串和 `TIMESTAMP` 字面量比较(`MEDIUM`) -10. **Spark 最佳实践** — `COUNT(*)` vs `COUNT(1)` 属于风格(`LOW`);生产环境用 `SELECT *` 算 `MEDIUM`;数据倾斜 join 缺 `DISTRIBUTE BY` / `CLUSTER BY` 算 `LOW` 到 `MEDIUM` +## 检查项 -## 输出契约(Output Contract) +1. 字段引用:每个 `alias.column` 是否存在于 `sql_context.fields`。 +2. alias:SQL 是否发明了 context 中没有的 alias。 +3. join key:ON 条件是否匹配 context 中固定的 join。 +4. join 基数:1:N 或 N:N join 是否会造成数据膨胀。 +5. 指标公式:聚合表达式是否符合 logic plan 和 context。 +6. 重复聚合:同一指标是否被重复 sum/count。 +7. 时间过滤:事实表是否缺少必要时间或分区过滤。 +8. 时间边界:T-1 完整日是否使用左闭右开窗口。 +9. 转化顺序:后置事件是否满足前置事件之后发生。 +10. GROUP BY:非聚合输出列是否都在 group by 中。 +11. 空值策略:join key 和度量字段是否需要 coalesce 或保留策略。 +12. 类型匹配:字符串日期、字符串金额是否需要 cast。 +13. 写入风险:INSERT 目标列、overwrite 确认、分区覆盖范围是否明确。 +14. 执行成本:是否存在无分区裁剪的大表全扫、笛卡尔积、窗口无 partition。 + +## 时间边界检查 + +当需求为“近 N 天按 T-1 完整日”且字段是 timestamp/date-like: + +推荐: + +```sql +event_time >= date_sub(current_date(), N) +AND event_time < current_date() +``` + +风险写法: + +```sql +event_time BETWEEN date_add(current_date(), -N) AND date_add(current_date(), -1) +``` + +如果 SQL 使用风险写法,至少标为 `MEDIUM`;如果会明显漏掉结束日期白天数据,标为 `HIGH`。 + +## 转化顺序检查 + +对注册、付费、激活、下单等漏斗指标,必须检查后置事件是否在前置事件之后。若口径是“注册后付费”,SQL 必须有等价约束: + +```sql +pay_time >= register_time +``` + +如果 SQL 先按用户取 `MIN(pay_time)`,再和 `register_time` 比较,而没有在聚合前加入 `pay_time >= register_time`,标为 `HIGH` 或 `MEDIUM`,因为这可能漏算注册后仍然完成付费的用户。 + +## 输出格式 ```yaml sql_review: - verdict: - highest_severity: - summary: "<一句话总结>" + verdict: PASS | FAIL + highest_severity: NONE | LOW | MEDIUM | HIGH + summary: "一句话结论" findings: - - id: F1 - severity: HIGH - check: reference_integrity - location: "<行/列号或 SQL 片段>" - message: "<错在哪>" - fix: "<具体的建议修改>" - - id: F2 - severity: MEDIUM - check: join_cardinality - location: "JOIN customer_info" - message: "1:N join without pre-aggregation; risk of row inflation" - fix: "Aggregate loan_order by customer_id before joining" + - id: R1 + severity: HIGH | MEDIUM | LOW + check: reference_integrity | join_cardinality | metric_formula | time_filter | time_boundary | conversion_order | group_by | type_check | write_risk | performance | other + location: "SQL 片段或行号" + message: "问题描述" + recommendation: "建议修复方向" notes: - - "SQL Context provided: yes (sql_context.aliases = [...])" - - "Reviewed by: sql-review v1" + - "评审假设" ``` -## 禁止行为(Forbidden Behaviors) +## 用户可见摘要 -| 自我说服 | 现实 | -|---|---| -| "SQL 看起来没问题,不用对一遍 context" | 字段可能"看起来对"但源表里根本没有,必须做交叉校验。 | -| "我顺手在评审时把 SQL 修了" | 评审只报告,修是作者的事。职责分离。 | -| "没给 SQL Context,那字段检查就跳过" | 缺 context 本身就是 `HIGH`,绝不能悄悄跳过。 | -| "只是开发查询,分区检查跳过" | 这一项成本极低,跳过就会留下静默的全表扫描。 | -| "1:N join 是常态,没风险" | 在保留粒度的前提下 1:N 是 OK 的,但评审必须验证,不能假设。 | -| "用正则就够了,能查 GROUP BY" | GROUP BY 成员关系是结构性的,要用解析器,不要用正则。 | +除了 YAML,必须输出: -## 完成判定(Completion Criteria) +```markdown +## SQL Review +- 结论:PASS / FAIL +- 最高风险:NONE / LOW / MEDIUM / HIGH +- 已检查:字段引用、join、GROUP BY、时间窗口、转化顺序、空值、类型、性能 +- 主要风险:... +- 是否可交付:... +``` -评审完成的条件: +## 严重程度 -- 全部 10 项检查都已评估 -- 每一条 `HIGH` 风险都有 `fix` -- 当且仅当 `highest_severity` 为 `MEDIUM` 或 `LOW` 时,`verdict` 才为 `PASS` -- 输出的 YAML 可直接贴到 PR 评论,无需修改 +- `HIGH`:会导致数字错误、执行失败、越权写入或违反策略,必须阻断。 +- `MEDIUM`:大概率可跑,但有性能、稳定性或边界风险,需要用户接受。 +- `LOW`:风格或可维护性问题。 + +只要存在 HIGH,`verdict` 就是 `FAIL`。只有 LOW/MEDIUM 时,可以把风险交给用户决定是否接受。 + +## 禁止行为 + +- 不要直接改 SQL;只输出发现和建议。 +- 不要凭直觉审字段,必须对照 `sql_context`。 +- 不要因为 SQL 能通过 guardrails 就跳过业务正确性审查。 +- 不要忽略写入模式,尤其是 `INSERT OVERWRITE`。