You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何构建含分支逻辑的GCS-Dataflow-BigQuery数据处理Airflow DAG

构建指定流程的Airflow DAG实现方案

核心流程梳理

对应需求的DAG执行逻辑如下:

  • 上传本地population_YYYYMMDD.csv至GCS存储桶的A文件夹,通过GCSObjectExistenceSensor校验文件存在
  • 启动Dataflow模板作业完成列名、数据类型转换
  • 分支处理:
    • 转换成功:将数据加载至BigQuery的A数据集下population_YYYYMMDD表,同时将原CSV文件移动至成功文件夹
    • 转换失败:将原CSV文件移动至失败文件夹

完整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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 04:40:30