【发布时间】: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