【发布时间】:2018-02-13 02:59:44
【问题描述】:
用例:我有具有 messageId 的消息,多条消息可以具有相同的消息 ID,这些消息存在于由 messageId 分区的流式管道(如 kafka)中,因此我确保所有具有相同 messageId 的消息都将进入相同的分区。
所以我需要编写一个作业,它应该将消息缓冲一段时间(比如说 1 分钟),然后将具有相同 messageId 的所有消息合并为单个大消息。
我认为可以使用 spark Datasets 和 spark sql(或其他东西?)来完成。但是我找不到任何关于如何为给定消息 ID 存储消息一段时间然后对这些消息进行聚合的示例/文档。
【问题讨论】:
-
你在想什么样的聚合?你想要一个聚合值(比如一个总和),还是想要一个消息的消息?
-
消息的消息,假设您有 10 条消息具有相同的消息 ID,我的结果应该是 1 条大消息,其中包含所有 10 条消息。希望清除。
标签: apache-kafka streaming spark-streaming buffering apache-samza