Airflow中如何将PythonOperator获取的参数传递给BatchOperator
Airflow 向 BatchOperator 传递 dag_run conf 中参数的实现方案
你不需要额外编写Python任务来中转参数,Airflow 自带的 Jinja 模板引擎可以直接在算子参数里读取触发时传入的配置,是最简洁的实现方式:
- 直接在BatchOperator的
parameters字段中使用模板变量读取dag_run.conf里的job_id即可,任务执行时会自动完成值替换:
submit_batch_job = BatchOperator( task_id='submit_batch_job_etl', job_name=JOB_NAME, job_queue=JOB_QUEUE, job_definition=JOB_DEFINITION, # 直接通过模板取触发参数 parameters={"job_id": "{{ dag_run.conf['job_id'] }}"} )
如果你需要保留前置的get_input任务做参数校验、格式转换等逻辑,可以通过XCom在任务间传值,先修正你原有代码的两处语法问题:
PythonOperator初始化时的dag参数传入错误,你写的dag=DAG是传入了DAG类,应该传入你实例化好的dag对象dag=dagget_inputs_config函数中parameters = ("{}".format(kwargs['dag_run'].conf['job_id'])缺少闭合右括号
修正后,PythonOperator的返回值会自动存入XCom,下游BatchOperator直接通过模板拉取对应XCom值即可,完整代码如下:
dag = DAG( dag_id='example_batch_submit_job', schedule_interval=None, start_date=datetime(2022, 7, 14), tags=['batch_job'], catchup=False ) def get_inputs_config(**kwargs): job_id = kwargs['dag_run'].conf['job_id'] print("Remotely received value of {} for key=job_id".format(job_id)) # 这里可以加你需要的参数校验、转换逻辑 return job_id run_this = PythonOperator( task_id='get_input', provide_context=True, python_callable=get_inputs_config, dag=dag, ) submit_batch_job = BatchOperator( task_id='submit_batch_job_etl', job_name=JOB_NAME, job_queue=JOB_QUEUE, job_definition=JOB_DEFINITION, # 拉取上游get_input任务的返回值 parameters={"job_id": "{{ ti.xcom_pull(task_ids='get_input') }}"} ) # 配置任务依赖 run_this >> submit_batch_job
补充说明:AWS BatchOperator的
parameters字段默认在模板渲染字段列表中,不需要额外修改算子的template_fields属性,写Jinja表达式即可直接生效。
内容的提问来源于stack exchange,提问作者Musa
相关产品推荐
相关产品推荐

