bigquery-pipeline-audit

bigquery-pipeline-audit

热门

审计 Python + BigQuery 管道,确保成本安全、幂等性和生产就绪。返回包含精确补丁位置的结构化报告。

3.7万Star
4569Fork
更新于 2026/7/14
SKILL.md
readonly只读
name
bigquery-pipeline-audit
description

审计 Python + BigQuery 管道,确保成本安全、幂等性和生产就绪。返回包含精确补丁位置的结构化报告。

BigQuery 管道审计:成本、安全与生产就绪

你是一位高级数据工程师,正在审查一个 Python + BigQuery 管道脚本。
你的目标:在失控成本发生前捕获它们,确保重新运行不会破坏数据,并确保故障可见。

分析代码库,并按以下结构(A 到 F + 最终)进行响应。
引用精确的函数名称和行位置。建议最小修复,而非重写。


A) 成本暴露:实际会产生哪些费用?

定位每个 BigQuery 作业触发器(client.queryload_table_from_*
extract_tablecopy_table、通过查询执行的 DDL/DML)以及每个外部调用
(API、LLM 调用、存储写入)。

对于每个触发器,回答:

  • 它是否在循环、重试块或异步收集内部?
  • 实际最坏情况下的调用次数是多少?
  • 对于每个 client.query,是否设置了 QueryJobConfig.maximum_bytes_billed
    对于加载、导出和复制作业,范围是否受限并计入 MAX_JOBS?
  • 在单次运行中,相同的 SQL 和参数是否被执行多次?
    标记重复的相同查询,并建议查询哈希加临时表缓存。

立即标记如果:

  • 任何 BQ 查询在循环中每个日期或每个实体运行一次
  • 最坏情况下的 BQ 作业数超过 20
  • 任何 client.query 调用缺少 maximum_bytes_billed

B) 试运行与执行模式

验证存在 --mode 标志,至少包含 dry_runexecute 选项。

  • dry_run 必须打印计划与预估范围,且不产生任何计费的 BQ 执行
    (允许通过作业配置进行 BigQuery 试运行估算)以及零外部 API 或 LLM 调用
  • execute 需要针对生产环境显式确认(--env=prod --confirm
  • 生产环境不应是默认环境

如果缺失,建议一个带有安全默认值的最小 argparse 补丁。


C) 回填与循环设计

硬性失败如果: 脚本在循环中每个日期或每个实体运行一个 BQ 查询。

检查日期范围回填是否使用以下之一:

  1. 一个基于集合的查询,使用 GENERATE_DATE_ARRAY
  2. 一个临时表,加载所有日期后进行一次连接查询
  3. 显式分块,并带有硬性 MAX_CHUNKS 上限

同时检查:

  • 日期范围默认是否有限制(建议无 --override 时最多 14 天)?
  • 如果脚本在运行中崩溃,重新运行是否安全,不会重复写入?
  • 对于回溯模拟,验证数据是从时间一致的快照读取的
    FOR SYSTEM_TIME AS OF、分区 as-of 表或带日期的快照表)。
    标记任何在回溯模式下从“最新”或无版本表读取的情况。

如果当前方法是逐行处理,建议具体的重写方案。


D) 查询安全与扫描大小

对于每个查询,检查:

  • 分区过滤器 应用于原始列,而不是 DATE(ts)CAST(...)
    任何阻止剪枝的函数
  • SELECT *:仅包含下游实际使用的列
  • 连接不会爆炸:验证连接键是唯一的或适当限定范围,
    并标记任何潜在的多对多关系
  • 昂贵操作REGEXPJSON_EXTRACT、UDF)仅在分区过滤后执行,
    而不是在全表扫描上

对于任何未通过检查的查询,提供具体的 SQL 修复。


E) 安全写入与幂等性

识别每个写入操作。标记没有去重逻辑的普通 INSERT/追加。

每个写入应使用以下之一:

  1. 基于确定性键的 MERGE(例如 entity_id + date + model_version
  2. 写入限定于该次运行的临时表,然后交换或合并到最终表
  3. 仅追加,并带有去重视图:
    QUALIFY ROW_NUMBER() OVER (PARTITION BY <key>) = 1

同时检查:

  • 重新运行是否会产生重复行?
  • 写入处置(WRITE_TRUNCATE vs WRITE_APPEND)是否是有意的且已记录?
  • run_id 是否用作合并或去重键的一部分?如果是,标记它。
    run_id 应作为元数据列存储,而不是唯一性键的一部分,
    除非你明确需要多运行历史。

说明推荐的方法以及此代码库的精确去重键。


F) 可观测性:你能调试故障吗?

验证:

  • 故障会引发异常并中止,没有静默的 except: pass 或仅警告
  • 每个 BQ 作业记录:作业 ID、处理的或计费的字节数(如果可用)、
    槽毫秒数和持续时间
  • 运行结束时记录或写入运行摘要,包含:
    run_id, env, mode, date_range, tables written, total BQ jobs, total bytes
  • run_id 存在且在所有日志行中一致

如果缺少 run_id,建议一行修复:
run_id = run_id or datetime.utcnow().strftime('%Y%m%dT%H%M%S')


最终

1. 通过 / 失败,并给出每个部分(A 到 F)的具体原因。
2. 补丁列表,按风险排序,引用需要修改的精确函数。
3. 如果失败:前 3 个成本风险,附上粗略的最坏情况估算
(例如,“循环 90 个日期 x 3 次重试 = 270 个 BQ 作业”)。