Cloud Composer调度Dataflow任务重复执行问题排查求助
在GCP Cloud Composer中调度DAG运行Dataflow作业,maf_to_bq_X任务通过Airflow TaskGroup分组,基于BeamRunPythonPipelineOperator实现,执行流程为:
- 从
variables_conf获取文件路径; - 读取并处理文件;
- 将输出写入BigQuery。
已设置max_active_runs=1和max_active_tasks=3,但每次运行管道时总会出现一个任务重复执行的问题:set_variables_conf任务完成后,3个任务并行启动,其中一个在Dataflow仍处于执行状态时被标记为up_for_retry,无错误完成后Airflow会再次调度该任务,最终导致BigQuery表中出现重复数据。尝试设置任务优先级权重后问题依旧。
附DAG代码如下:
with models.DAG( "ieo_dima_extraction", default_args={ "start_date": pendulum.today('Europe/Rome'), 'retries': 3, "dataflow_default_options": { "project": project_id, "region": "europe-west1", } }, on_success_callback=cleanup_xcom, is_paused_upon_creation=True, render_template_as_native_obj=True, max_active_runs=1, max_active_tasks=3, schedule_interval=None ) as dag: set_variables = PythonOperator( task_id='set_variables_conf', python_callable=set_variables_conf, op_kwargs={'bucket_name': 'dima_landing_area'}, do_xcom_push=True ) list_of_maf = Variable.get('dima_maf_paths', default_var={}, deserialize_json=True) mapping_file_path = Variable.get('dima_mapping_path', default_var='') dima_bucket = 'dima_landing_area' with TaskGroup(group_id='process_maf_files', prefix_group_id=False) as process_maf_files: if list_of_maf['files']: for index, maf_file in enumerate(list_of_maf['files']): maf_to_bq = maf_to_bq_job(index, processing_path, maf_file, bq_data_table, dataset_name, temp_location, staging_location, project_id, dataflow_region, mapping_file_path) truncate_map_tmp_start = truncate_job("start", dataset_name, bq_mapping_tmp_table, bq_job_region) with TaskGroup(group_id='process_mapping', prefix_group_id=False) as process_mapping: map_to_bq = mapping_to_bq(update_map_path, bq_mapping_tmp_table, dataset_name, sql_view_file, temp_location, staging_location, project_id, dataflow_region, mapping_file_path) merge_map = merge_job(query_merge, dataset_name, bq_mapping_table, bq_job_region) truncate_map_tmp_end = truncate_job("end", dataset_name, bq_mapping_tmp_table, bq_job_region) update_view = create_view_job(sql_view_script, dataset_name, bq_job_region) chain(map_to_bq, merge_map, truncate_map_tmp_end, update_view) chain(set_variables, truncate_map_tmp_start, [process_maf_files, process_mapping])
核心问题1:DAG上下文外的变量读取导致任务生成不稳定
list_of_maf和mapping_file_path的Variable.get()调用写在DAG定义的with块外部,Airflow调度器会周期性解析DAG文件,每次解析都会执行这段代码,可能因变量值变化、解析时机差异导致TaskGroup内的任务重复生成或实例异常,进而触发重复调度。
解决办法:将变量读取逻辑移到DAG上下文内部,确保仅在DAG解析时正确读取一次:
with models.DAG(...) as dag: set_variables = PythonOperator(...) # 移到DAG块内 list_of_maf = Variable.get('dima_maf_paths', default_var={}, deserialize_json=True) mapping_file_path = Variable.get('dima_mapping_path', default_var='') dima_bucket = 'dima_landing_area' with TaskGroup(group_id='process_maf_files', prefix_group_id=False) as process_maf_files: ...
核心问题2:默认重试机制与Dataflow异步作业不兼容
default_args中设置了retries=3,而BeamRunPythonPipelineOperator是异步触发Dataflow作业,Airflow会定期轮询作业状态。如果轮询间隔内Dataflow作业状态未更新(比如长运行作业处于等待资源阶段),Airflow可能误判为任务失败,触发重试机制,导致同一Dataflow作业被重复提交,最终写入重复数据到BigQuery。
解决办法:
- 针对
BeamRunPythonPipelineOperator单独禁用重试:在创建maf_to_bq任务时覆盖retries=0; - 配置合理的
execution_timeout和poke_interval,避免Airflow误判作业状态。例如:
maf_to_bq = BeamRunPythonPipelineOperator( ..., retries=0, execution_timeout=timedelta(hours=4), poke_interval=300 # 5分钟轮询一次状态 )
核心问题3:TaskGroup内任务ID生成逻辑不稳定
当前用index作为任务ID的一部分,若list_of_maf['files']的顺序发生变化,会导致任务ID变更,Airflow会将其视为新任务,旧任务可能残留并被重复调度。
解决办法:使用文件路径/文件名作为任务ID的唯一标识,确保任务ID稳定:
import os for maf_file in list_of_maf['files']: # 从文件名生成唯一task_id,避免依赖index task_id = f"maf_to_bq_{os.path.basename(maf_file).replace('.', '_')}" maf_to_bq = maf_to_bq_job(task_id, processing_path, maf_file, ...)
额外检查点
确认cleanup_xcom回调函数未破坏set_variables任务推送的XCom数据,若XCom数据被意外清理,可能导致后续任务重新触发变量读取逻辑,引发重复执行。
内容的提问来源于stack exchange,提问作者alex-mont

