如何在Airflow流水线中避免Bigquery表已有数据时重复加载CSV文件
实现方案
以下是经过生产验证的可行实现逻辑,核心围绕文件哈希校验+BigQuery元数据管理实现去重:
前置准备
首先在BigQuery中创建一张文件加载元数据表,用于记录所有已成功加载的文件信息,表结构参考:
CREATE TABLE `your_project.your_dataset.file_load_metadata` ( file_path STRING COMMENT '源文件存储路径', file_hash STRING COMMENT '文件SHA256哈希值', load_time TIMESTAMP COMMENT '加载完成时间' ) CLUSTER BY file_hash; -- 可选:给file_hash加唯一约束,避免重复写入 ALTER TABLE `your_project.your_dataset.file_load_metadata` ADD CONSTRAINT unique_file_hash UNIQUE (file_hash) NOT ENFORCED;
DAG实现步骤
整个DAG分为4个核心任务节点:
- 任务1:计算待加载CSV文件的哈希值,用PythonOperator实现即可,大文件可以分块读取避免内存溢出:
def calculate_file_hash(**context): import hashlib # 替换为你的文件读取逻辑,本地文件/GCS文件均可 file_path = context["dag_run"].conf.get("file_path") hash_sha256 = hashlib.sha256() # 如果是GCS文件,用gcsfs库读流即可 with open(file_path, "rb") as f: for chunk in iter(lambda: f.read(4096), b""): hash_sha256.update(chunk) file_hash = hash_sha256.hexdigest() context["ti"].xcom_push(key="file_hash", value=file_hash) context["ti"].xcom_push(key="file_path", value=file_path)
- 任务2:校验文件是否已加载,用BigQueryCheckOperator查询元数据表:
from airflow.providers.google.cloud.operators.bigquery import BigQueryCheckOperator check_duplicate_task = BigQueryCheckOperator( task_id="check_file_duplicate", sql=""" SELECT COUNT(*) = 0 FROM `your_project.your_dataset.file_load_metadata` WHERE file_hash = @file_hash """, query_params=[ { "name": "file_hash", "parameterType": {"type": "STRING"}, "parameterValue": {"value": "{{ ti.xcom_pull(key='file_hash') }}"} } ], use_legacy_sql=False, location="你的BigQuery区域" )
这里直接返回COUNT(*) = 0的布尔结果,不存在就返回True,存在就返回False。
- 任务3:分支判断,用BranchPythonOperator根据校验结果决定走哪个分支:
from airflow.operators.python import BranchPythonOperator def load_branch_decision(**context): is_new_file = context["ti"].xcom_pull(task_ids="check_file_duplicate") if is_new_file: return "load_csv_to_bq_task" return "skip_load_task" branch_task = BranchPythonOperator( task_id="load_branch_decision", python_callable=load_branch_decision, provide_context=True )
- 任务4:两个分支执行对应逻辑:
- 跳过分支:直接用DummyOperator标记即可,无需额外操作
- 加载分支:先调用
GCSToBigQueryOperator(CSV存GCS场景)或者对应加载算子完成数据写入,加载成功后再执行一个BigQueryInsertJobOperator,将当前文件的哈希、路径等信息写入元数据表,建议把加载和写入元数据的操作放在同一个BigQuery事务中执行,避免加载成功但元数据写入失败导致的重复加载问题。
备选方案
如果你不想单独维护元数据表,也可以直接在目标BigQuery表中新增file_hash字段,每次加载前查询目标表是否存在对应哈希值,不存在再执行加载,该方案适合需要保留数据来源溯源的场景。
内容的提问来源于stack exchange,提问作者Amarjeet
相关产品推荐
相关产品推荐

