【问题标题】:Don't understand Update Mode and watermark in Structured Streaming不理解结构化流中的更新模式和水印
【发布时间】:2018-08-19 20:47:27
【问题描述】:

我有以下代码,它输出

number: 1, count: 1
number: 2, count: 1
number: 3, count: 2
number: 6, count: 2
number: 7, count: 1

我认为number: 6, count: 2 不应该输出,因为事件低于水位线。但我不明白它为什么输出

import java.sql.Timestamp

import org.apache.spark.sql.execution.streaming.MemoryStream
import org.apache.spark.sql.{ForeachWriter, Row, SparkSession}

object UpdateModeWithWatermarkTest {
  def main(args: Array[String]): Unit = {
    val spark: SparkSession = SparkSession.builder()
      .appName("UpdateModeWithWatermarkTest")
      .config("spark.sql.shuffle.partitions", 1)
      .master("local[2]").getOrCreate()

    import spark.implicits._


    val inputStream = new MemoryStream[(Timestamp, Int)](1, spark.sqlContext)
    val now = 5000L

    val aggregatedStream = inputStream.toDS().toDF("created", "number")
      .withWatermark("created", "1 second")
      .groupBy("number")
      .count()

    val query = aggregatedStream.writeStream.outputMode("update")
      .foreach(new ForeachWriter[Row] {
        override def open(partitionId: Long, epochId: Long): Boolean = true

        override def process(value: Row): Unit = {
          println(s"number: ${value.getInt(0)}, count: ${value.getLong(1)}")
        }

        override def close(errorOrNull: Throwable): Unit = {}
      }).start()

    new Thread(new Runnable() {
      override def run(): Unit = {
        inputStream.addData(
          (new Timestamp(now + 5000), 1),
          (new Timestamp(now + 5000), 2),
          (new Timestamp(now + 5000), 3),
          (new Timestamp(now + 5000), 3)
        )
        while (!query.isActive) {
          Thread.sleep(50)
        }
        Thread.sleep(10000)

        // At this point, the water mark is (now  + 5000) - 1 second = 9 seconds
        // when adding following two events: (new Timestamp(4000L), 6),  (new Timestamp(now), 6)
        // These two events are below water mark, so that they should be discarded, then should not output number: 6, count: 2
        inputStream.addData((new Timestamp(4000L), 6))
        inputStream.addData(
          (new Timestamp(now), 6),
          (new Timestamp(11000), 7)
        )
      }
    }).start()

    query.awaitTermination(45000)


  }

}

【问题讨论】:

  • 你在这个 Tom 上的表现如何?

标签: apache-spark spark-structured-streaming


【解决方案1】:

其实没那么难。

Watermark 允许考虑将迟到的数据包含在一段时间内使用窗口的已计算结果中。它的前提是它跟踪到一个时间点,在该时间点之前假定不再有迟到的事件应该到达,但如果它们到达了,它们就会被丢弃。

https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#window-operations-on-event-time 上的优秀示例,并附有漂亮的图表。

【讨论】:

  • 感谢@thebluephantom。我通过阅读火花代码来分析结果。请看一下 StateStoreSaveExec#doExecute,对于这种情况(更新),它首先会丢弃低于水位线的事件(baseIterator 变量是通过丢弃低于水位线的事件创建的)
  • 来自附加模式的手册,但我想知道更新模式:必须在与聚合中使用的时间戳列相同的列上调用 withWatermark。例如, df.withWatermark("time", "1 min").groupBy("time2").count() 在 Append 输出模式下无效,因为 watermark 是在与聚合列不同的列上定义的。您的分组依据在号码上。
  • 还有这个:其他聚合完成,更新由于没有定义水印(仅在其他类别中定义),旧的聚合状态不会被丢弃。不支持追加模式,因为聚合可以更新,因此违反了此​​模式的语义。
  • 文档有点模糊。
  • 尝试 aggr 并针对 created 更新,看看会发生什么。
【解决方案2】:

我想关于output mode of structured streaming的官方解释已经回答了你的问题。

更新模式 - (自 Spark 2.1.1 起可用)只有自上次触发后更新的结果表中的行才会输出到接收器。更多信息将在未来版本中添加。

在您的问题中,这意味着 1 秒内到达的数据将更新归档的“计数”

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-05-21
    • 2021-05-24
    • 1970-01-01
    • 2021-05-03
    • 2019-04-06
    • 1970-01-01
    • 2019-06-04
    • 2018-01-16
    相关资源
    最近更新 更多