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

为何并行SourceFunction添加rebalance后Watermark可正常传递?

关于Flink中Rebalance算子修复Watermark传递问题的疑问

我编写了一段Flink代码用于验证并行SourceFunction的Watermark行为,定义的并行SourceFunction仅第一个分区会产生数据。由于Watermark采用上游所有Watermark的最小值,另外两个无数据的分区导致Watermark无法向前传递,控制台无输出符合预期。但在assignTimestampsAndWatermarks前添加rebalance算子后,Watermark成功传递,控制台出现输出。现请教为何rebalance能让Watermark正常传递?

import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.AssignerWithPunctuatedWatermarks
import org.apache.flink.streaming.api.functions.source.{RichParallelSourceFunction, SourceFunction}
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.scala.function.ProcessWindowFunction
import org.apache.flink.streaming.api.watermark.Watermark
import org.apache.flink.streaming.api.windowing.time.Time
import org.apache.flink.streaming.api.windowing.windows.TimeWindow
import org.apache.flink.util.Collector

import java.text.SimpleDateFormat
import java.util.Date
import java.util.concurrent.TimeUnit

case class MyOrder(name: String, create_date: String)

class MyParallelOrderSourceFunction extends RichParallelSourceFunction[MyOrder] {
  val data = ParallelSourceFunctionWatermarkTest.createData()

  override def run(ctx: SourceFunction.SourceContext[MyOrder]): Unit = {
    val idx = getRuntimeContext.getIndexOfThisSubtask
    //Only the first subTask/partition will have data
    if (idx == 0) {
      data.map(e => MyOrder(e._1, e._2)).foreach(ctx.collect)
    }
    while (true) {
      Thread.sleep(1000)
    }
  }

  override def cancel(): Unit = {
  }
}

object ParallelSourceFunctionWatermarkTest {
  def to_milli(str: String) =
    new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").parse(str).getTime

  def to_char(milli: Long) = {
    val date = if (milli <= 0) new Date(0) else new Date(milli)
    new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(date)
  }

  def createData(): Seq[(String, String)] = {
    //create 100 records, create_date increase by 1 second for each record
    var time = "2023-02-26 17:26:11"
    (1 to 100).map {
      _ =>
        time = to_char(to_milli(time) + 1000)
        ("A", time)
    }
  }


  def main(args: Array[String]) {
    val env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration())

    //MyParallelOrderSourceFunction's parallelism is 3
    //call rebalance before assignTimestampsAndWatermarks, water mark forwards; if doesn't call rebalance, then no water mark forwards
    env.addSource(new MyParallelOrderSourceFunction).setParallelism(3).rebalance/*rebalance method call makes watermark forward works*/.assignTimestampsAndWatermarks(new AssignerWithPunctuatedWatermarks[MyOrder] {

      private var max_seen: Long = Long.MinValue
      private val allowLateness = 4000

      override def checkAndGetNextWatermark(lastElement: MyOrder, extractedTimestamp: Long): Watermark = {
        if (extractedTimestamp > max_seen) {
          max_seen = extractedTimestamp
        }
        new Watermark(max_seen - allowLateness)
      }

      override def extractTimestamp(element: MyOrder, recordTimestamp: Long): Long = {
        to_milli(element.create_date)
      }
    }).keyBy(_.name.hashCode % 3).timeWindow(Time.of(6, TimeUnit.SECONDS)).process(new ProcessWindowFunction[MyOrder, String, Int, TimeWindow] {
      override def process(key: Int, cxt: Context, elements: Iterable[MyOrder], out: Collector[String]): Unit = {
        val ws = to_char(cxt.window.getStart)
        val we = to_char(cxt.window.getEnd)
        val sb = new StringBuilder()

        sb.append(s"($ws, $we), ${elements.map(e => (e.name, e.create_date)).mkString(",")}")
        out.collect(sb.toString())
      }
    }).setParallelism(3).print()

    env.execute()
  }


}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 06:47:45