【问题标题】: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