Cloud Composer/Airflow如何仅在DataFusion任务真正完成后启动新任务
问题原因
CloudDataFusionStartPipelineOperator默认采用异步触发逻辑,提交DataFusion流水线请求后不会持续等待作业执行完成,只要作业进入RUNNING状态就会直接判定Airflow任务执行成功,因此会出现DataProc集群预配完成、作业刚进入运行态就被标记成功的现象,无法满足等待作业完全完成再触发下游任务的需求。
修复方案
调整Operator初始化参数,开启同步轮询,显式指定成功判定状态即可:
- 导入流水线状态枚举类
- 关闭异步模式,让Operator持续轮询流水线运行状态
- 明确指定仅
COMPLETED状态为任务成功的判定条件,避免版本默认值差异导致行为不符合预期
修正后的完整代码如下:
from airflow.providers.google.cloud.operators.datafusion import ( CloudDataFusionStartPipelineOperator, PipelineStates ) my_task = CloudDataFusionStartPipelineOperator( location='europe-west1', pipeline_name="my_datafusion_pipeline_name", instance_name="my_datafusion_instance_name", task_id="my_task_name", asynchronous=False, success_states=[PipelineStates.COMPLETED], poke_interval=30, # 可选参数,自定义状态轮询间隔,单位为秒,默认值30 )
如果你使用的Google Cloud Provider版本较低,导入
PipelineStates时报错,可以将状态类的导入替换为from airflow.providers.google.cloud.utils.datafusion import PipelineStatus,对应成功状态参数写PipelineStatus.COMPLETED即可。
内容的提问来源于stack exchange,提问作者Sander van den Oord
相关产品推荐
相关产品推荐

