【问题标题】:Cloud Composer / Airflow start new task only when Cloud DataFusion task is really finishedCloud Composer / Airflow 仅在 Cloud DataFusion 任务真正完成时才开始新任务
【发布时间】:2022-11-11 20:05:41
【问题描述】:

我在触发 Cloud DataFusion 管道的 Airflow (Cloud Composer) 中有以下任务。

问题是:
当(在 DataFusion 中)已配置 DataProc 集群并且实际作业已进入 RUNNING 状态时,Airflow 认为此任务已经成功。

但我只希望它在完成时被认为是成功的。

from airflow.providers.google.cloud.operators.datafusion import \
    CloudDataFusionStartPipelineOperator

my_task = CloudDataFusionStartPipelineOperator(
    location='europe-west1',
    pipeline_name="my_datafusion_pipeline_name",
    instance_name="my_datafusion_instance_name", 
    task_id="my_task_name",
)

【问题讨论】:

    标签: python airflow google-cloud-composer google-cloud-data-fusion


    【解决方案1】:

    我不得不查看源代码,但以下状态是默认的success_states:
    [PipelineStates.COMPLETED] + [PipelineStates.RUNNING]

    因此,您必须使用关键字success_states 将succes_states 限制为仅[PipelineStates.COMPLETED],如下所示:

    from airflow.providers.google.cloud.operators.datafusion import 
        CloudDataFusionStartPipelineOperator
    from airflow.providers.google.cloud.hooks.datafusion import PipelineStates
    
    my_task = CloudDataFusionStartPipelineOperator(
        location='europe-west1',
        pipeline_name="my_datafusion_pipeline_name",
        instance_name="my_datafusion_instance_name", 
        task_id="my_task_name",
        success_states=[PipelineStates.COMPLETED], # overwrite default success_states
        pipeline_timeout=3600, # in seconds, default is currently 300 seconds
    )
    

    也可以看看:
    Airflow documentation on the DataFusionStartPipelineOperator

    Airflow source code used for success states of DataFusionStartPipelineOperator

    【讨论】:

      猜你喜欢
      • 2019-07-21
      • 1970-01-01
      • 2021-02-22
      • 1970-01-01
      • 2019-02-20
      • 1970-01-01
      • 2018-06-24
      • 2019-01-20
      • 2012-10-28
      相关资源
      最近更新 更多