You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink是否具备Spark awaitAnyTermination类似功能及多流代码转写咨询

Spark Structured Streaming多流任务转Flink实现及API疑问

原Spark Structured Streaming代码

val df1 = ...
val ds2 = ...
val ds3 = ...

ds1
  .writeStream
  .queryName("stream1")
  .outputMode(OutputMode.Append())
  .foreachBatch((ds, _) => {
    MyUtils.save2Redis(ds)
  })
  .trigger(Trigger.ProcessingTime("2 seconds"))
  .start()

ds2
  .writeStream
  .queryName("stream2")
  .outputMode(OutputMode.Update())
  .foreachBatch((ds, _) => {
    MyUtils.save2Hbase(ds)
  })
  .trigger(Trigger.ProcessingTime("1 seconds"))
  .start()

ds3
  .writeStream
  .queryName("stream3")
  .outputMode(OutputMode.Update())
  .foreachBatch((ds, _) => {
    MyUtils.save2Hbase(ds)
    MyUtils.save2Hbase(ds)
  })
  .trigger(Trigger.ProcessingTime("2 seconds"))
  .start()

spark.streams.awaitAnyTermination()

对应的Flink实现代码

环境初始化

首先初始化Flink流执行环境,设置处理时间语义:

import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.windowing.time.Time
import org.apache.flink.streaming.api.TimeCharacteristic

// 初始化流执行环境
val env = StreamExecutionEnvironment.getExecutionEnvironment
// 设置处理时间语义,对应Spark的ProcessingTime Trigger
env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime)
// 可选:根据集群资源调整并行度
env.setParallelism(1)

各流任务实现

假设ds1、ds2、ds3是已定义好的DataStream[YourDataType],对应原Spark中的流数据集:

  1. ds1:每2秒写入Redis(对应Append模式)
    通过timeWindowAll实现固定时间间隔的批量处理,模拟Spark的foreachBatch逻辑:
ds1
  // 每2秒触发一次批量处理
  .timeWindowAll(Time.seconds(2))
  .apply { (window, batchData, out: Collector[YourDataType]) =>
    // 调用工具类写入Redis,batchData为当前窗口内的所有数据
    MyUtils.save2Redis(batchData)
    // 无需后续算子可省略out.collect逻辑
  }
  // 添加空Sink补全流链路(Flink流必须以Sink结尾)
  .addSink(new SinkFunction[YourDataType] {
    override def invoke(value: YourDataType): Unit = {}
  })
  1. ds2:每1秒写入HBase(对应Update模式)
ds2
  .timeWindowAll(Time.seconds(1))
  .apply { (window, batchData, out: Collector[YourDataType]) =>
    MyUtils.save2Hbase(batchData)
  }
  .addSink(new SinkFunction[YourDataType] {})
  1. ds3:每2秒两次写入HBase
ds3
  .timeWindowAll(Time.seconds(2))
  .apply { (window, batchData, out: Collector[YourDataType]) =>
    MyUtils.save2Hbase(batchData)
    MyUtils.save2Hbase(batchData)
  }
  .addSink(new SinkFunction[YourDataType] {})

执行作业

// 启动Flink作业,所有流任务作为同一作业的一部分运行
env.execute("multi jobs")

关于Flink中类似awaitAnyTermination的功能

Spark的awaitAnyTermination()是等待任意一个流查询终止后退出整个应用,Flink的对应实现分为两种场景:

  • 单作业多流场景:若所有流任务在同一个StreamExecutionEnvironment中定义,它们会作为单个Flink作业的一部分运行。当作业中任意一个算子失败或终止,整个作业会停止,env.execute()会抛出异常并退出,效果与awaitAnyTermination()的"任意任务终止则应用退出"逻辑一致。
  • 多独立作业场景:若要完全模拟Spark中多个独立流查询的模式(每个流为独立作业),可创建多个StreamExecutionEnvironment,分别提交作业后监听状态,当任意一个作业终止时停止所有作业:
import org.apache.flink.api.common.JobStatus

val env1 = StreamExecutionEnvironment.getExecutionEnvironment
val env2 = StreamExecutionEnvironment.getExecutionEnvironment
val env3 = StreamExecutionEnvironment.getExecutionEnvironment

// 分别定义三个流任务...

val job1 = env1.executeAsync("stream1")
val job2 = env2.executeAsync("stream2")
val job3 = env3.executeAsync("stream3")

// 等待任意一个作业终止
val jobs = List(job1, job2, job3)
while (jobs.forall(_.getJobStatus.isRunning)) {
  Thread.sleep(1000)
}
// 终止所有剩余作业
jobs.filter(_.getJobStatus.isRunning).foreach(_.cancel())

内容的提问来源于stack exchange,提问作者ohmycloudy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 06:17:01