Airflow从1.10.12迁移至2.2.4后出现'DagModel'对象无'create_dagrun'属性错误的排查求助
Airflow 2.2.4迁移报错:AttributeError: 'DagModel' object has no attribute 'create_dagrun'
问题背景
我把一个在Airflow 1.10.12中运行正常的DAG迁移到Airflow 2.2.4后,触发了如下错误:
AttributeError: 'DagModel' object has no attribute 'create_dagrun'
对应的日志详情和追溯信息如下:
[2022-07-19, 09:55:13 UTC] {taskinstance.py:1272} INFO - 将任务标记为失败。dag_id=backfill.scheduler, task_id=schedule_runners, execution_date=20220719T095000, start_date=20220719T095512, end_date=20220719T095513 [2022-07-19, 09:55:13 UTC] {standard_task_runner.py:89} ERROR - 执行任务schedule_runners的作业76失败 追溯信息(最近的调用在最前): File "/home/airflow/.local/lib/python3.9/site-packages/airflow/task/task_runner/standard_task_runner.py", line 85, in _start_by_fork args.func(args, dag=self.dag) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/cli/cli_parser.py", line 48, in command return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/utils/cli.py", line 92, in wrapper return f(*args, **kwargs) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/cli/commands/task_command.py", line 298, in task_run _run_task_by_selected_method(args, dag, ti) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/cli/commands/task_command.py", line 107, in _run_task_by_selected_method _run_raw_task(args, ti) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/cli/commands/task_command.py", line 180, in _run_raw_task ti._run_raw_task( File "/home/airflow/.local/lib/python3.9/site-packages/airflow/utils/session.py", line 70, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/models/taskinstance.py", line 1334, in _run_raw_task self._execute_task_with_callbacks(context) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/models/taskinstance.py", line 1460, in _execute_task_with_callbacks result = self._execute_task(context, self.task) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/models/taskinstance.py", line 1516, in _execute_task result = execute_callable(context=context) File "/home/airflow/.local/lib/python3.9/site-packages/airflow/operators/python.py", line 174, in execute return_value = self.execute_callable() File "/home/airflow/.local/lib/python3.9/site-packages/airflow/operators/python.py", line 188, in execute_callable
问题原因
这个报错是Airflow 2.x版本对DAG运行管理的核心逻辑做了重构导致的:
- 在Airflow 1.x中,
DagModel类提供了create_dagrun方法来创建DAG运行实例; - 但到了Airflow 2.x,官方将创建DAG运行的职责从
DagModel转移到了DagRun类本身,同时移除了DagModel中的create_dagrun方法,所以直接调用旧方法就会触发属性不存在的错误。
解决方案
你需要找到DAG代码中调用dag_model.create_dagrun()的代码段,替换为Airflow 2.x支持的写法:
替换示例
原来Airflow 1.x的代码可能是这样的:
from airflow.models.dag import DagModel dag_model = DagModel.get_dagmodel(dag_id="your_dag_id") dag_model.create_dagrun( execution_date=your_execution_date, run_type="manual", conf={}, session=session )
改成Airflow 2.x的写法,直接使用DagRun.create()方法:
from airflow.models.dagrun import DagRun DagRun.create( dag_id="your_dag_id", execution_date=your_execution_date, run_type="manual", # 可选值包括manual、scheduled、backfill等,根据你的需求选择 conf={}, # 传递给DAG的配置参数 session=session # 可选,如果你需要指定数据库会话 )
额外建议
如果你的代码中依赖了更多DAG运行相关的旧API,建议逐步替换所有1.x的私有/已弃用API,避免后续出现更多兼容性问题。
内容的提问来源于stack exchange,提问作者hugo95
相关产品推荐
相关产品推荐

