如何在Airflow DAG中为任务传递字典键参数?
问题修复:Airflow DAG中无法将字典key传递给BashOperator的params参数
你遇到的问题根源在于:BashOperator中引用了{{ params.x }}模板变量,但并未给operator传入params参数,导致Airflow无法识别x这个参数,从而抛出未定义的错误。
修复后的完整代码
from airflow import DAG from airflow.operators.bash_operator import BashOperator from datetime import datetime default_args = { 'owner': 'Airflow', 'start_date': datetime(2023, 1, 1), 'retries': 0, } mydict = { 1: "name1", 2: "name2", 3: "name3" } with DAG('test_dag', default_args=default_args, schedule_interval=None, catchup=False) as dag: for x, y in mydict.items(): run_test_1 = BashOperator( task_id=f'task1_{y}', bash_command="python /path-to-the-script/simple_script.py '{{ params.x }}'", params={'x': x} # 显式传入循环中的key值作为params参数 ) run_test_2 = BashOperator( task_id=f'task2_{y}', bash_command="echo '{{ params.x }}'", params={'x': x} # 同样传入params参数 ) run_test_1 >> run_test_2
关键修复点说明
- 添加
params参数:在每个BashOperator中声明params={'x': x},将循环遍历得到的字典key值x传递给模板变量{{ params.x }},让Airflow能够正确解析该变量。 - 移除无用变量:删除了代码中未实际使用的
INDEX变量,简化逻辑。 - 修正拼写错误:将脚本名
simple_scirpt.py改为simple_script.py,避免后续执行时出现文件找不到的错误。
内容的提问来源于stack exchange,提问作者dips
相关产品推荐
相关产品推荐

