Airflow同步Google Sheet至BigQuery的财务周报数据固化方案咨询
解决方案:Airflow ETL Google Sheet至BigQuery时保留历史财务报告数据
针对你遇到的「每周财务报告提交后数据需固化,同时允许Google Sheet频繁修改」的问题,以下是几个可落地的方案:
方案1:按周创建独立快照表
这是最直接的隔离方式,彻底分开历史报告与日常更新数据:
- 放弃
WRITE_TRUNCATE覆盖主表的逻辑,改为每周生成对应周数的快照表,命名规则用finance_report_yyyyww(例如finance_report_202419代表2024年第19周) - 利用Airflow的
execution_date自动生成周标识,代码示例:from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator create_snapshot_task = BigQueryInsertJobOperator( task_id="create_weekly_finance_snapshot", configuration={ "query": { "query": """ CREATE OR REPLACE TABLE `your-project.your-dataset.finance_report_{{ execution_date.strftime('%Y%W') }}` AS SELECT * FROM `your-project.your-dataset.staging_google_sheet_data`; """, "useLegacySql": False, } }, gcp_conn_id="google_cloud_default", ) - 主表保留最新数据供日常查看,财务团队直接查询对应周的快照表即可获取固化报告;后续Sheet修改只会更新主表,不影响已生成的快照。
方案2:主表加版本字段+历史视图
适合需要统一数据入口、不想维护大量独立表的场景:
- 在BigQuery目标表新增两个字段:
report_week STRING(存储yyyyww格式周标识)、load_timestamp TIMESTAMP(记录数据加载时间) - 修改ETL逻辑:放弃
WRITE_TRUNCATE,改为先删除当前周旧数据,再追加新数据,确保截止节点前当周数据可修改,代码示例:update_current_week_task = BigQueryInsertJobOperator( task_id="update_current_week_finance_data", configuration={ "query": { "query": """ -- 删除当前周旧数据 DELETE FROM `your-project.your-dataset.finance_report` WHERE report_week = '{{ execution_date.strftime('%Y%W') }}'; -- 追加最新数据 INSERT INTO `your-project.your-dataset.finance_report` SELECT *, '{{ execution_date.strftime('%Y%W') }}' AS report_week, CURRENT_TIMESTAMP() AS load_timestamp FROM `your-project.your-dataset.staging_google_sheet_data`; """, "useLegacySql": False, } }, gcp_conn_id="google_cloud_default", ) - 创建历史报告视图,自动保留每周截止节点前的最后一次加载数据:
CREATE OR REPLACE VIEW `your-project.your-dataset.finance_report_history` AS SELECT * EXCEPT(rn) FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY report_week ORDER BY load_timestamp DESC) AS rn FROM `your-project.your-dataset.finance_report` -- 可选:添加截止节点过滤,仅保留截止后的最终数据 WHERE DATE(load_timestamp) <= DATE_TRUNC(DATE_ADD(DATE(load_timestamp), INTERVAL 1 WEEK), WEEK) ) WHERE rn = 1; - 在Airflow中添加截止节点判断:定义每周截止日(比如周一),当
execution_date超过截止日时,禁止更新对应周数据,仅允许维护历史视图。
方案3:利用BigQuery时间旅行(临时回溯场景)
如果仅需临时回溯历史数据,且不需要长期固化,可使用BigQuery时间旅行功能:
- 原逻辑用
WRITE_TRUNCATE覆盖表,可通过时间旅行查询截止节点时刻的表状态,示例SQL:SELECT * FROM `your-project.your-dataset.finance_report` FOR SYSTEM_TIME AS OF TIMESTAMP('2024-05-06 00:00:00'); -- 替换为每周截止节点的时间戳 - 注意:时间旅行默认保留7天,最多可延长至1年,不适合长期存储历史报告数据。
关键细节补充
- 权限控制:确保财务团队能访问所有快照表或历史视图,避免数据权限问题
- 截止节点自动化:在Airflow DAG中定义
FINANCE_CUTOFF_WEEKDAY = 0(0代表周一),通过execution_date.weekday()判断是否已过截止日,自动切换ETL逻辑 - 数据验证:每次生成快照或更新数据后,添加数据校验任务(比如行数匹配、关键指标核对),确保数据准确性
内容的提问来源于stack exchange,提问作者Novan Dwi Atmaja
相关产品推荐
相关产品推荐

