Airflow默认变量使用方法及dag_run未定义报错解决方案咨询
Let's break down why you're hitting this error and how to fix it properly.
The Root Cause
Your code uses Jinja template syntax ({{ dag_run.task_id }}) directly inside a Python function, but Python doesn't recognize this syntax. Airflow's Jinja variables like dag_run are only resolved during template rendering—they aren't available as native Python variables in your function unless you explicitly pass them via the execution context.
Solution: Access dag_run via the Execution Context
If this function is meant to be used as a python_callable for a BranchPythonOperator (which it looks like, given the return value logic), you need to modify your function to accept the Airflow execution context, then pull dag_run from it.
Here's the corrected code:
def decide_which_task(**context): # Extract dag_run from the provided context current_dag_run = context['dag_run'] if current_dag_run.task_id == "Move_file": return "move_file" else: return "push_to_db"
Then, when defining your BranchPythonOperator, make sure to enable context passing:
- For Airflow 1.x: Add
provide_context=Trueto the operator arguments - For Airflow 2.x: Context is passed by default, but explicitly setting
provide_context=Truestill works for compatibility
Example operator definition:
from airflow.operators.python import BranchPythonOperator branch_task = BranchPythonOperator( task_id='decide_which_task', python_callable=decide_which_task, provide_context=True, # Required for Airflow 1.x, optional but safe for 2.x dag=dag )
Why This Works
The **context parameter tells Airflow to pass the full execution context (including dag_run, task_instance, ds, etc.) to your function as keyword arguments. You can then access any of these variables directly from the context dictionary.
内容的提问来源于stack exchange,提问作者anjum

