Airflow中基于dag_run.run_type分支任务报错:KeyError: 'dag_run'
解决Airflow中KeyError: 'dag_run'的问题
错误原因
KeyError: 'dag_run'是因为你的get_the_date函数被调用时,上下文(context)字典中没有传入dag_run对象。这通常是因为函数没有在Airflow的任务上下文环境中执行,或者调用方式不符合Airflow的模板规则。同时你的原代码还存在语法错误(比如DAG参数缺失逗号、with语句未缩进、BashOperator命令未闭合),这些也会导致执行异常。
解决方案
下面提供两种可行的修复方案,按需选择:
方案一:用PythonOperator获取日期并通过XCom传递
这种方式适合需要对日期做复杂处理的场景:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.bash import BashOperator from datetime import datetime def get_the_date(**context): # 从上下文中获取dag_run对象 run_type = context['dag_run'].run_type if run_type == "manual": # 直接返回ds变量的实际值,无需模板语法 return context['ds'] else: return "2025-02-24" with DAG( dag_id="mydag", start_date=datetime(2024, 9, 11), schedule_interval="@daily", # 替换为你的调度规则,比如"0 0 * * *" catchup=False ) as dag: # 先执行Python任务获取目标日期 get_date_task = PythonOperator( task_id='get_date_task', python_callable=get_the_date, provide_context=True, # 确保上下文被传入函数(Airflow 2.x可省略,但显式指定兼容性更好) do_xcom_push=True ) # 通过XCom获取日期并执行Bash命令 my_task = BashOperator( task_id='my_task', bash_command='echo "{{ ti.xcom_pull(task_ids=\'get_date_task\') }}"' ) # 设置任务依赖 get_date_task >> my_task
方案二:直接在BashOperator中用Jinja模板判断(更简洁)
如果仅需简单的分支判断,直接利用Airflow的Jinja模板支持即可,无需额外Python函数:
from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime with DAG( dag_id="mydag", start_date=datetime(2024, 9, 11), schedule_interval="@daily", catchup=False ) as dag: my_task = BashOperator( task_id='my_task', # 直接在模板中判断run_type并输出对应日期 bash_command='echo "{% if dag_run.run_type == \'manual\' %}{{ ds }}{% else %}2025-02-24{% endif %}"' )
关键注意事项
- 确保DAG定义的语法正确:参数末尾加逗号、with语句内的任务要缩进、字符串要正确闭合。
- Airflow 2.x中,PythonOperator默认会传递上下文,但显式设置
provide_context=True可以兼容旧版本。 - 不要在函数返回值中加
{{ ds }}模板语法,直接返回变量值即可,模板语法应在任务的模板字段中使用。
内容的提问来源于stack exchange,提问作者Aviator
相关产品推荐
相关产品推荐

