Airflow任务如何识别是否通过airflow test CLI命令执行?
解决方案
Airflow 内置了无需手动传参的检测方式,通过任务实例(TaskInstance)的 test_mode 属性即可判断任务是否通过 airflow tasks test 命令启动:
1. 在PythonOperator中使用
在Python任务里,可以通过上下文参数ti(即TaskInstance对象)直接获取test_mode属性,据此调整SQL逻辑:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def test_mode_demo(**context): ti = context["ti"] # 判断是否处于test模式 if ti.test_mode: target_table = "test_output_table" else: target_table = "prod_output_table" # 构造并执行SQL sql = f"INSERT INTO {target_table} SELECT * FROM source_table" print(f"执行SQL: {sql}") # 这里替换为实际的数据仓库执行逻辑 with DAG( dag_id="test_mode_dag", start_date=datetime(2024, 1, 1), schedule_interval=None ) as dag: task = PythonOperator( task_id="check_test_mode", python_callable=test_mode_demo, provide_context=True )
2. 在SQL类Operator中使用(如BigQueryOperator/PostgresOperator)
可以通过Jinja模板直接在SQL语句中引用{{ ti.test_mode }},动态切换目标表:
from airflow import DAG from airflow.providers.postgres.operators.postgres import PostgresOperator from datetime import datetime with DAG( dag_id="sql_test_mode_dag", start_date=datetime(2024, 1, 1), schedule_interval=None ) as dag: sql_task = PostgresOperator( task_id="dynamic_table_sql", postgres_conn_id="your_db_conn", sql=""" INSERT INTO {{ 'test_output_table' if ti.test_mode else 'prod_output_table' }} SELECT * FROM source_table; """ )
原理说明
当使用airflow tasks test命令启动任务时,Airflow会自动将当前TaskInstance的test_mode属性设置为True;正常调度或手动触发的任务,该属性为False。这个属性是Airflow原生支持的,无需额外配置或传参,完全避免了用户遗忘加参数的问题。
内容的提问来源于stack exchange,提问作者Antoine Lefebvre
相关产品推荐
相关产品推荐

