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

Cloud Composer调度Dataflow任务重复执行问题排查求助

问题描述

在GCP Cloud Composer中调度DAG运行Dataflow作业,maf_to_bq_X任务通过Airflow TaskGroup分组,基于BeamRunPythonPipelineOperator实现,执行流程为:

  1. 从variables_conf获取文件路径;
  2. 读取并处理文件;
  3. 将输出写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 14:25:33