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

并发调度大量Airflow DAG时出现重复条目错误的技术问询

问题解答:Airflow同一秒触发同DAG多运行的唯一键冲突

这是Airflow旧版本的设计限制(非Bug,已在新版本修复)

早期Airflow(1.x系列)的dag_run表使用(dag_id, execution_date)作为唯一约束键,而外部触发DAG时,若不手动指定execution_date,Airflow会默认生成秒级精度的时间戳作为该值。当你在短时间内批量触发多个同DAG的运行时,就会出现多个运行共享同一个execution_date(秒级重复),触发数据库的唯一键冲突,也就是你遇到的Duplicate entry错误。

这个设计在早期被视为“特性”——Airflow团队当时认为同一DAG在同一秒内不需要多实例运行,但社区反馈这一限制过于僵硬,因此在Airflow 2.0及以上版本中,官方已经将dag_run表的唯一约束修改为(dag_id, run_id),彻底解决了这个问题(因为你代码中已经生成了唯一的run_id,新版本会以此来区分不同的DAG运行)。


快速修复方案

根据你的场景,有几种可行的解决路径:

1. 优先升级到Airflow 2.0+

这是最彻底、最推荐的方案。新版本不仅修复了这个唯一键问题,还在批量调度、性能优化上有大幅提升,非常适合你每5分钟触发200个DAG运行的场景。

2. 旧版本下手动指定execution_date(临时应急)

如果暂时无法升级,可以在触发DAG时,给每个运行分配微秒级精度的execution_date,确保每个运行的该值唯一:

  • 修改测试代码中的触发逻辑:
    from datetime import datetime, timedelta
    from unittest import TestCase
    from backend.tasks.airflow import trigger_dag
    
    class TestTriggerDag(TestCase):
        def test_trigger_dag(self):
            game_ids = [99, 100, 101, 102, 103]
            for idx, game_id in enumerate(game_ids):
                # 给每个运行分配不同的微秒级时间
                exec_date = datetime.now() + timedelta(microseconds=idx * 1000)
                trigger_dag("update_game_dag", game_id=game_id, execution_date=exec_date)
            self.assertTrue(True)
    
  • 或者修改trigger_dag函数内部,自动生成唯一的execution_date:
    from datetime import datetime
    from typing import List
    import random
    import time
    from airflow.api.client.local_client import Client
    from airflow.models.dagrun import DagRun
    
    afc = Client(None, None)
    
    def get_dag_run_state(dag_id: str, run_id: str):
        return DagRun.find(dag_id=dag_id, run_id=run_id)[0].state
    
    def trigger_dag(dag_id: str, wait_for_complete: bool = False, execution_date=None, **kwargs):
        run_hash = '%030x' % random.randrange(16**30)
        kwarg_list = [f"{str(k)}:{str(v)}" for k, v in kwargs.items()]
        run_id = f"{run_hash}-{'_'.join(kwarg_list)}"
        # 用当前微秒级时间作为execution_date,确保唯一
        exec_date = execution_date or datetime.now()
        afc.trigger_dag(dag_id, run_id=run_id, conf=kwargs, execution_date=exec_date)
        
        while wait_for_complete and get_dag_run_state(dag_id, run_id) == "running":
            time.sleep(1)
        return get_dag_run_state(dag_id, run_id)
    

3. 手动修改数据库唯一约束(不推荐)

如果你使用的是Airflow 1.x,也可以直接修改dag_run表的唯一约束,将原有的(dag_id, execution_date)改为(dag_id, run_id)。但这种操作属于侵入性修改,可能会和Airflow 1.x的部分内部逻辑冲突,仅作为临时应急方案,不建议长期使用。


是否需要提交PR或Jira工单?

不需要。这个问题已经在Airflow 2.0+版本中被官方修复。如果你使用的是最新版本仍然遇到该问题,才需要去ASF Jira提交工单并考虑PR,但目前新版本已经完全解决了这个限制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 23:27:33