【问题标题】:Ordering of events in structured streaming window aggregates in append mode附加模式下结构化流窗口聚合中的事件排序
【发布时间】:2020-05-07 17:56:31
【问题描述】:

我在使用 Spark 进行结构化流式传输时遇到问题。

当前设置:我有一个来自 kafka 的数据流。每条消息都有一个事件时间。我正在使用这些事件时间来进行窗口聚合,并使用水印规则来丢弃状态。 输出方式为追加方式。

目标:我需要在它们过期时按顺序获取窗口聚合,以便我可以按事件时间窗口的顺序处理这些事件。由于我的滑动窗口,我希望窗口状态会按顺序过期。

问题:有时打印的消息顺序不是基于 windows 的顺序。例如

|[2020-06-11 08:02:00, 2020-06-11 08:03:00]|

|[2020-06-11 08:01:00, 2020-06-11 08:02:00]|

为什么窗口没有按顺序放置?我想订购这个。 请帮忙

【问题讨论】:

  • 你的 Kafka 主题有多少个分区?如果它有多个,它是如何分区的?
  • 你的意思是我的源卡夫卡主题吗?源卡夫卡和汇卡夫卡都只有一个分区。
  • 流数据集仅在聚合后和完整输出模式下才支持排序操作。
  • @thebluephantom 为什么我不能在微批次中对值进行排序?我不想跨批次排序,我想在微批次中排序

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


【解决方案1】:

它不在架构中(还)。 SO上的许多帖子都消除了这种观念。

正如 cricket_007 在许多帖子中所说,您最好将排序留给一般的下游系统。这种方式也更加灵活,RDBMS 的整个概念是固定数据排序顺序不太有效 - 除了集群。

如果您查看此用例https://mapr.com/blog/real-time-analysis-popular-uber-locations-spark-structured-streaming-machine-learning-kafka-and-mapr-db/,您会发现排序不起作用。也就是说,我看到很多请求,但许多请求无需排序即可实现目标。

如果排序顺序由“生产者”设置并且该顺序是可接受的,则单分区主题是较低数量的结果。此外,您可以考虑使用 KSQL 并写入单个分区 KAFKA 主题以随后读取或使用 Java、Scala 的 KAFKA 流。

我认为问题在于几年前我看到的一篇帖子:

结构化流的基本原则是查询应该返回 流式或批处理模式下的相同答案。我们支持排序 完整模式,因为我们拥有所有数据并且可以正确排序 并返回完整的答案。在更新或追加模式下,排序将 只有当我们可以保证记录时才返回正确的答案 sort lower 将稍后到达(我们不能)。因此它是 不允许。

【讨论】:

  • 我需要检测数据中的顺序模式,因此只有在微批处理中缓冲、排序和处理数据时才有可能。
  • 我基本上需要在微批次中进行顺序模式检测
  • 生活并不总是我们想要的样子
  • 是的,你是对的。顺便说一句,旧的 Dstream api 允许使用 .transform 方法对这些微批次进行排序?
  • 是的,但它是遗留的,所以不相关。很遗憾。大多数 Sp Str Str 可以在没有排序的情况下完成,但可能是一种稍微不同的心态 reqd。成功。点赞就好了
猜你喜欢
  • 1970-01-01
  • 2021-09-07
  • 2018-07-22
  • 2020-09-09
  • 2019-10-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多