如何构建含分支逻辑的GCS-Dataflow-BigQuery数据处理Airflow DAG
构建指定流程的Airflow DAG实现方案
核心流程梳理
对应需求的DAG执行逻辑如下:
- 上传本地
population_YYYYMMDD.csv至GCS存储桶的A文件夹,通过GCSObjectExistenceSensor校验文件存在 - 启动Dataflow模板作业完成列名、数据类型转换
- 分支处理:
- 转换成功:将数据加载至BigQuery的
A数据集下population_YYYYMMDD表,同时将原CSV文件移动至成功文件夹 - 转换失败:将原CSV文件移动至失败文件夹
- 转换成功:将数据加载至BigQuery的
完整DAG代码实现
from airflow import DAG from airflow.providers.google.cloud.sensors.gcs import GCSObjectExistenceSensor from airflow.providers.google.cloud.operators.dataflow import DataflowTemplatedJobStartOperator from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator from airflow.providers.google.cloud.operators.gcs import GCSMoveObjectOperator from airflow.utils.dates import days_ago from airflow.operators.python import BranchPythonOperator from airflow.utils.state import State # 配置参数 PROJECT_ID = "your-gcp-project-id" GCS_BUCKET = "your-gcs-bucket" SOURCE_FOLDER = "A" SUCCESS_FOLDER = "success" FAILURE_FOLDER = "failure" DATAFLOW_TEMPLATE_GCS_PATH = "gs://your-dataflow-template-bucket/templates/csv-to-bq-transform" BQ_DATASET = "A" DATE_SUFFIX = "{{ execution_date.strftime('%Y%m%d') }}" SOURCE_FILE = f"population_{DATE_SUFFIX}.csv" SOURCE_GCS_PATH = f"{SOURCE_FOLDER}/{SOURCE_FILE}" default_args = { 'owner': 'airflow', 'start_date': days_ago(1), 'retries': 1 } with DAG( 'population_data_pipeline', default_args=default_args, schedule_interval='@daily', catchup=False, tags=['gcs', 'dataflow', 'bigquery'] ) as dag: # 1. 校验GCS中上传的文件是否存在 check_file_exists = GCSObjectExistenceSensor( task_id='check_population_file_exists', bucket=GCS_BUCKET, object=SOURCE_GCS_PATH, mode='poke' ) # 2. 启动Dataflow模板作业执行转换 run_dataflow_transform = DataflowTemplatedJobStartOperator( task_id='run_dataflow_transform', project_id=PROJECT_ID, template=DATAFLOW_TEMPLATE_GCS_PATH, parameters={ 'inputFile': f"gs://{GCS_BUCKET}/{SOURCE_GCS_PATH}", 'outputTable': f"{PROJECT_ID}.{BQ_DATASET}.population_{DATE_SUFFIX}" }, location='us-central1' ) # 3. 判断Dataflow作业执行结果,分支处理 def decide_next_task(**context): task_instance = context['task_instance'] job_status = task_instance.xcom_pull(task_ids='run_dataflow_transform')['job']['status'] return 'load_to_bigquery' if job_status == 'JOB_STATE_DONE' else 'move_to_failure_folder' branch_task = BranchPythonOperator( task_id='decide_processing', python_callable=decide_next_task, provide_context=True ) # 4. 成功分支:加载数据到BigQuery load_to_bigquery = BigQueryInsertJobOperator( task_id='load_to_bigquery', project_id=PROJECT_ID, configuration={ "load": { "sourceUris": [f"gs://{GCS_BUCKET}/{SOURCE_GCS_PATH}"], "destinationTable": { "projectId": PROJECT_ID, "datasetId": BQ_DATASET, "tableId": f"population_{DATE_SUFFIX}" }, "sourceFormat": "CSV", "skipLeadingRows": 1, "autodetect": False, "schemaUpdateOptions": ["ALLOW_FIELD_ADDITION"] } } ) # 5. 成功分支:移动文件到成功文件夹 move_to_success_folder = GCSMoveObjectOperator( task_id='move_to_success_folder', source_bucket=GCS_BUCKET, source_object=SOURCE_GCS_PATH, destination_bucket=GCS_BUCKET, # 可改为其他桶 destination_object=f"{SUCCESS_FOLDER}/{SOURCE_FILE}" ) # 6. 失败分支:移动文件到失败文件夹 move_to_failure_folder = GCSMoveObjectOperator( task_id='move_to_failure_folder', source_bucket=GCS_BUCKET, source_object=SOURCE_GCS_PATH, destination_bucket=GCS_BUCKET, # 可改为其他桶 destination_object=f"{FAILURE_FOLDER}/{SOURCE_FILE}" ) # 设置任务依赖 check_file_exists >> run_dataflow_transform >> branch_task branch_task >> load_to_bigquery >> move_to_success_folder branch_task >> move_to_failure_folder
关键细节说明
- 日期变量:通过
{{ execution_date.strftime('%Y%m%d') }}动态生成日期后缀,适配每日调度的文件命名规则 - Dataflow模板:需提前将转换逻辑打包成Dataflow模板存储在GCS,
parameters需根据你的模板需求调整输入输出参数 - 分支判断:通过
BranchPythonOperator读取Dataflow作业的状态XCom数据,决定后续执行分支 - 文件移动:
GCSMoveObjectOperator支持同桶或跨桶移动,只需修改destination_bucket参数即可
内容的提问来源于stack exchange,提问作者Nhu Dao
相关产品推荐
相关产品推荐

