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
相关产品推荐
相关产品推荐

