Airflow 2.0.2中同一DAG并行运行异常问题求助
问题原因与解决方案
核心问题
你触发DAG B时,所有实例都用了同一个execution_date({{ ts }}即DAG A的执行时间戳)。Airflow的规则是:同一DAG的相同execution_date只能存在一个活跃运行实例,所以即使你设置了max_active_runs=4,也只会有一个DAG B实例实际运行。
解决步骤
给每个触发的DAG B生成唯一的execution_date,确保每个并行实例的执行时间不重复。
修改DAG A的TriggerDagRunOperator代码
把execution_date参数改成带唯一标识的取值,比如结合循环变量i生成不同值:
for i in list: run_task = TriggerDagRunOperator( task_id='trigger_' + i, trigger_dag_id='dag_b', # 方式1:给每个实例设置偏移的时间,符合Airflow时间格式要求 execution_date=f'{{{{ execution_date.add(minutes={int(i)*5}) }}}}', # 方式2:用无分隔符时间戳加循环变量,生成完全唯一标识 # execution_date=f'{{{{ ts_nodash }}}}_{i}', conf={"run_job": i}, reset_dag_run=True, wait_for_completion=True, poke_interval=60 )
两种方式任选其一即可:
- 方式1通过时间偏移保证execution_date唯一,符合Airflow对执行时间的格式规范;
- 方式2用字符串拼接生成唯一标识,Airflow允许非标准格式的execution_date,只要全局唯一即可。
额外检查
- 确认DAG B的
max_active_runs=4和concurrency=4配置生效,你当前的代码已经设置正确; - 检查Airflow全局配置
core.parallelism、core.dag_concurrency的值,确保全局并行数足够支撑多个DAG实例同时运行。
内容的提问来源于stack exchange,提问作者Dozel
相关产品推荐
相关产品推荐

