Airflow中dagA调用dagB时如何正确传递dag_run.conf参数
解决Airflow TriggerDagRunOperator传递列表参数类型丢失问题
问题出在Jinja模板渲染时默认会把非字符串类型(比如列表)转成字符串格式,导致dagB收到的key值是字符串化的列表"['value']",而非预期的数组["value"]。
可以通过Jinja的tojson过滤器解决这个问题,它会保留原始数据类型的JSON格式:
方案1:传递完整的dag_run.conf
如果需要把dagA的全部conf参数传递给dagB,直接在conf中使用tojson过滤器:
trigger_dagB = TriggerDagRunOperator( task_id='trigger_dagB', trigger_dag_id='dagB', execution_date='{{ ds }}', conf='{{ dag_run.conf | tojson }}', reset_dag_run=True, wait_for_completion=True, poke_interval=60 )
方案2:仅传递特定key
如果只需要传递key字段,对该字段单独使用tojson过滤器:
trigger_dagB = TriggerDagRunOperator( task_id='trigger_dagB', trigger_dag_id='dagB', execution_date='{{ ds }}', conf='{"key": {{ dag_run.conf["key"] | tojson }}}', reset_dag_run=True, wait_for_completion=True, poke_interval=60 )
原理说明
tojson过滤器会将Python对象转换成标准的JSON格式字符串,Airflow在解析conf参数时,会将这个JSON字符串还原成对应的Python字典/列表结构,这样dagB就能收到符合要求的{"key": ["value"]}格式参数。
内容的提问来源于stack exchange,提问作者Dozel
相关产品推荐
相关产品推荐

