Airflow中如何传递DatabricksRunNowOperator返回值给下游任务
可以通过XCom实现参数传递,具体实现步骤如下
1. 调整上游任务配置与Notebook逻辑
- 给
verification_run算子添加do_xcom_push=True参数,开启XCom推送能力 - 在上游
verification_task对应的Databricks Notebook末尾,添加返回逻辑,把生成的date参数输出:
# verification_task对应的Databricks Notebook代码示例 import json # 你的业务逻辑,生成process_date变量 process_date = "2024-01-01" # 输出结果供Airflow拉取 dbutils.notebook.exit(json.dumps({"date": process_date}))
2. 下游任务通过Jinja模板拉取XCom值
DatabricksRunNowOperator的notebook_params参数原生支持Jinja模板渲染,直接在参数中通过任务实例ti拉取上游XCom值即可,拉取到的值会自动传入下游对应的Databricks Notebook。
下游Notebook可以通过dbutils.widgets.get("date")直接获取到传递的date参数。
修改后完整DAG代码
from airflow import DAG from airflow.providers.databricks.operators.databricks import DatabricksRunNowOperator from airflow.utils.dates import days_ago default_args = { 'owner': 'airflow' } with DAG('databricks_dag', start_date = days_ago(2), schedule_interval = None, default_args = default_args ) as dag: verification_run = DatabricksRunNowOperator( task_id = 'verification_task', databricks_conn_id = 'databricks_default', job_id = '-----', do_xcom_push = True # 开启XCom推送 ) insert_run = DatabricksRunNowOperator( task_id = 'insert_task', databricks_conn_id = 'databricks_default', job_id = '-----', # 拉取上游XCom值作为参数传入 notebook_params = { "date": "{{ ti.xcom_pull(task_ids='verification_task') | fromjson | get('date') }}" } ) workspace_run = DatabricksRunNowOperator( task_id = 'workspace_task', databricks_conn_id = 'databricks_default', job_id = '------', notebook_params = { "date": "{{ ti.xcom_pull(task_ids='verification_task') | fromjson | get('date') }}" } ) verification_run >> [insert_run,workspace_run]
版本兼容说明
如果你使用的Airflow Databricks Provider包版本低于4.0,XCom返回的结果会嵌套在notebook_output.result字段中,取值写法调整为以下内容即可:
{{ ti.xcom_pull(task_ids='verification_task')['notebook_output']['result'] | fromjson | get('date') }}
内容的提问来源于stack exchange,提问作者navin619
相关产品推荐
相关产品推荐

