如何部署可自动停止的指定时长Dataflow流式作业?
Dataflow流式作业自动停止的实现方案
Dataflow本身没有类似max_workers的直接流水线配置参数,来指定作业运行时长或到特定时间后自动停止,但可以通过以下几种方式实现需求:
借助Airflow调度机制实现
既然是通过Airflow触发的作业,可以在DAG中设计两个关联任务:- 使用
DataflowStartFlexTemplateOperator启动Dataflow流式作业,将返回的作业ID通过Airflow XCom传递给下一个任务; - 用
PythonOperator(或自定义操作符)延迟指定时长(比如24小时)后,调用Dataflow的取消接口终止作业。可以通过google-cloud-dataflow客户端的cancel_job方法,或者执行gcloud dataflow jobs cancel <JOB_ID>命令完成取消操作。
- 使用
在Dataflow作业内部嵌入停止逻辑
针对流式Pipeline,可以自定义DoFn来实现定时停止:
在数据处理流程中加入一个检查逻辑,判断当前时间是否超过作业启动时间+指定时长,或者是否到达预设的停止时间,满足条件时抛出StopBundleException终止作业。示例代码如下:import time import apache_beam as beam class StopAtTimeDoFn(beam.DoFn): def __init__(self, stop_after_hours=24): self.stop_after_hours = stop_after_hours self.start_time = time.time() def process(self, element): elapsed_seconds = time.time() - self.start_time if elapsed_seconds >= self.stop_after_hours * 3600: raise beam.utils.exceptions.StopBundleException("已达到指定运行时长,终止作业") yield element将该DoFn插入到Pipeline的合适节点(比如数据源读取之后)即可。
通过外部调度工具触发取消
可以结合Cloud Scheduler和Cloud Functions,在指定时间触发函数调用Dataflow取消API,但因为已经使用Airflow作为调度工具,这种方式相比Airflow内部调度更冗余,优先推荐前两种方案。
内容的提问来源于stack exchange,提问作者oikonomiyaki
相关产品推荐
相关产品推荐

