【问题标题】:Can we pass x_com variable to the next task in DAG as a parameter?我们可以将 x_com 变量作为参数传递给 DAG 中的下一个任务吗?
【发布时间】:2021-09-16 11:15:39
【问题描述】:

请考虑这种情况:

我在触发云功能的云作曲家上有一个云 Dag。该函数命中一个 api,然后将表存储在 GCS 中。现在我的 Airflow DAG(使用 Cloud Composer)触发下一阶段,即 Dataproc 作业,它从 GCS 获取表并推送到 BQ,但是当我触发我的 Dataproc 工作流模板时,我传递了一个参数,它是来自 dag 的表的名称本身以及我想从 x_com 中选择的那个参数。

这是一段代码,sn-p 抛出一个未定义的错误

dataproc_job = dataproc_operator.DataprocWorkflowTemplateInstantiateOperator(
# The task id of your job
task_id="dataproc_job",
# The template id of your workflow
template_id="newwf1",
project_id='#######',
region="us-central1",
parameters={"TABLE_NAME":ti.xcom_pull(task_ids=simple_http}
)

如何解决此错误并将 x_com 值作为参数传递给我在 DAG 中的下一步?

【问题讨论】:

    标签: google-cloud-platform google-cloud-composer airflow


    【解决方案1】:

    假设parameters 参数可以被模板化(又名在DataprocWorkflowTemplateInstantiateOperator 中列为templated_field),您可以使用Jinja 表达式来访问XCom 值。

    DataprocWorkflowTemplateInstantiateOperator(
        # The task id of your job
        task_id="dataproc_job",
        # The template id of your workflow
        template_id="newwf1",
        project_id='#######',
        region="us-central1", 
        parameters={"TABLE_NAME": "{{ ti.xcom_pull(task_ids='simple_http' }}"} 
    )
    
    

    详细了解 Airflow here 中的 Jinja 模板以及将 XComsthis doc 中提到的模板一起使用。

    顺便说一句,该运算符看起来像一个非常古老的 Airflow 1 运算符。如果可以的话,我强烈建议升级到Airflow 2。有无数的功能和性能改进,可以让您的 Airflow 和管道执行体验更好。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-04-10
      • 2022-06-15
      • 1970-01-01
      相关资源
      最近更新 更多