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中的流数据集:
- 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 = {} })
- ds2:每1秒写入HBase(对应Update模式)
ds2 .timeWindowAll(Time.seconds(1)) .apply { (window, batchData, out: Collector[YourDataType]) => MyUtils.save2Hbase(batchData) } .addSink(new SinkFunction[YourDataType] {})
- 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
相关产品推荐
相关产品推荐

