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

Airflow DAG开发需求:空CSV移至GCS归档桶,非空文件导入BigQuery

Airflow处理GCS CSV文件:空文件归档/非空文件导入BigQuery方案

现有代码核心问题

  • detect_csv任务误用GoogleCloudStorageToGoogleCloudStorageOperator:该Operator直接执行移动操作,完全跳过了检查文件大小的核心判断逻辑
  • 任务依赖逻辑不符合需求:当前串行流程会导致空文件被错误导入BigQuery,且detect_csv提前移动文件后,后续load_bigquery会找不到源文件
  • 未实现分支判断:没有区分空文件和非空文件的不同处理路径

推荐Operator选型

  • PythonOperator:编写自定义逻辑调用GCS SDK,获取文件大小并完成判断
  • BranchPythonOperator:根据文件大小的判断结果,动态选择后续执行的任务分支
  • GoogleCloudStorageToBigQueryOperator:官方标准Operator,用于将非空CSV文件导入BigQuery
  • GoogleCloudStorageToGoogleCloudStorageOperator:官方标准Operator,用于移动文件到归档桶
  • DummyOperator:作为分支流程的起始/结束节点,统一任务流结构

完整DAG示例代码

from airflow import DAG
from airflow.contrib.operators.gcs_to_bq import GoogleCloudStorageToBigQueryOperator
from airflow.contrib.operators.gcs_to_gcs import GoogleCloudStorageToGoogleCloudStorageOperator
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.operators.dummy import DummyOperator
from airflow.utils.dates import datetime
from google.cloud import storage

# 初始化GCS客户端
storage_client = storage.Client()

def check_csv_file_size(**context):
    """检查GCS中CSV文件的大小,返回对应的分支任务ID"""
    source_bucket = context['params']['source_bucket']
    source_prefix = context['params']['source_prefix']
    
    bucket = storage_client.get_bucket(source_bucket)
    # 获取匹配前缀的CSV文件(单文件场景,多文件可扩展遍历逻辑)
    blobs = bucket.list_blobs(prefix=source_prefix)
    csv_blob = next((b for b in blobs if b.name.endswith('.csv')), None)
    
    if not csv_blob:
        # 无CSV文件,返回对应分支
        return 'no_file_found'
    elif csv_blob.size == 0:
        # 文件为空,返回归档分支
        return 'move_empty_to_archive'
    else:
        # 文件非空,返回导入BigQuery分支
        return 'load_to_bigquery'

# 定义DAG
dag = DAG(
    'gcs_csv_processing',
    description='处理GCS中CSV文件:空文件归档,非空文件导入BigQuery',
    schedule_interval='@daily',
    start_date=datetime(2023, 7, 1),
    catchup=False,
    params={
        'source_bucket': 'your-source-bucket',
        'archive_bucket': 'your-archive-bucket',
        'source_prefix': 'prefix/path/to/csv',
        'archive_prefix_empty': 'prefix/archive/empty/',
        'archive_prefix_processed': 'prefix/archive/processed/',
        'bq_table': 'your-project.your_dataset.your_table'
    }
)

# 起始节点
start = DummyOperator(task_id='start', dag=dag)

# 检查文件大小并分支
check_file_size = BranchPythonOperator(
    task_id='check_file_size',
    python_callable=check_csv_file_size,
    provide_context=True,
    dag=dag
)

# 无文件分支
no_file_found = DummyOperator(task_id='no_file_found', dag=dag)

# 空文件归档分支
move_empty_to_archive = GoogleCloudStorageToGoogleCloudStorageOperator(
    task_id='move_empty_to_archive',
    source_bucket="{{ params.source_bucket }}",
    source_object="{{ params.source_prefix }}*.csv",
    destination_bucket="{{ params.archive_bucket }}",
    destination_object="{{ params.archive_prefix_empty }}",
    move_object=True,
    dag=dag
)

# 非空文件导入BigQuery
load_to_bigquery = GoogleCloudStorageToBigQueryOperator(
    task_id='load_to_bigquery',
    bucket="{{ params.source_bucket }}",
    source_objects=["{{ params.source_prefix }}*.csv"],
    destination_project_dataset_table="{{ params.bq_table }}",
    schema_fields=[
        {'name': 'column1', 'type': 'STRING'},
        {'name': 'column2', 'type': 'INTEGER'},
        {'name': 'column3', 'type': 'FLOAT'},
    ],
    write_disposition='WRITE_APPEND',  # 按需调整,推荐增量导入用APPEND
    source_format='CSV',
    skip_leading_rows=1,  # CSV有表头时启用
    dag=dag
)

# 导入完成后归档文件
move_processed_to_archive = GoogleCloudStorageToGoogleCloudStorageOperator(
    task_id='move_processed_to_archive',
    source_bucket="{{ params.source_bucket }}",
    source_object="{{ params.source_prefix }}*.csv",
    destination_bucket="{{ params.archive_bucket }}",
    destination_object="{{ params.archive_prefix_processed }}",
    move_object=True,
    dag=dag
)

# 结束节点
end = DummyOperator(task_id='end', trigger_rule='none_failed_or_skipped', dag=dag)

# 定义任务依赖
start >> check_file_size
check_file_size >> [no_file_found, move_empty_to_archive, load_to_bigquery]
load_to_bigquery >> move_processed_to_archive
[no_file_found, move_empty_to_archive, move_processed_to_archive] >> end

关键说明

  • 多文件处理:如果GCS路径下有多个CSV文件,可扩展check_csv_file_size函数,遍历所有CSV文件并分类处理
  • 配置管理:使用Airflow Jinja模板语法{{ params.xxx }}统一管理配置,便于后续维护
  • BigQuery写入策略:根据业务需求调整write_disposition,WRITE_TRUNCATE适合全量覆盖场景
  • 触发规则:end节点使用trigger_rule='none_failed_or_skipped',确保任意分支执行完成后都能触发结束节点

内容的提问来源于stack exchange,提问作者user17489358

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 16:35:06