【问题标题】:Don't understand Update mode and watermarking不懂更新模式和水印
【发布时间】:2021-05-24 19:33:05
【问题描述】:

据我了解,水印是最后一次看到的事件时间 - 延迟阈值。因此,如果最后看到的事件时间是 12:11,而延迟阈值是 10 分钟,则水印是 12:01。由于 12:01 晚于 12:00 的窗口开始时间,因此它的状态被丢弃。

但我写了查询:

stream
  .withWatermark("created", "2 seconds")
  .groupBy(
    window($"created", "2 seconds", "2 seconds"),
    $"animal"
  )
  .count()
  .writeStream
  .format("console")
  .outputMode(OutputMode.Update())

还有输出:

[2021-02-22 16:06:40.0,2021-02-22 16:06:42.0]:dog
[2021-02-22 16:06:40.0,2021-02-22 16:06:42.0]:owl
[2021-02-22 16:06:40.0,2021-02-22 16:06:42.0]:cat
[2021-02-22 16:06:34.0,2021-02-22 16:06:36.0]:pig

上次事件时间:2021-02-22 16:06:41.696 在窗口 40-42 秒内

猪时间:2021-02-22 16:06:35.696

你可以看到,窗口 34-36 中存在猪,但阈值为 2 秒。

为什么我可以看到 pig int 输出?

有趣的事情:如果我将 pig 与其他事件同时推送但使用旧时间戳,则此事件将添加到结果集中。但是如果事件在 2 秒(阈值)后以相同的时间戳推送,则不会显示在结果集中。

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-streaming spark-structured-streaming


    【解决方案1】:

    我将所有数据一次性推送到流中。那时还没有水印,这就是为什么我可以在结果集中看到旧事件。如果我将一些数据推送到流中,设置ProcessingTime,例如100ms,100ms后将推送旧数据,结果是预期的。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-02-04
      • 1970-01-01
      • 2020-06-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-28
      • 1970-01-01
      相关资源
      最近更新 更多