【问题标题】:How to sessionize / group the events in Akka Streams?如何对 Akka Streams 中的事件进行会话/分组?
【发布时间】:2018-03-12 10:28:43
【问题描述】:

要求是我想编写一个 Akka 流应用程序,它侦听来自 Kafka 的连续事件,然后根据嵌入在每个事件中的一些 id 值在一个时间范围内对事件数据进行会话。

例如,假设我的时间框架窗口是两分钟,在前两分钟我得到以下四个事件:

输入:

{"message-domain":"1234","id":1,"aaa":"bbb"}
{"message-domain":"1234","id":2,"aaa":"bbb"}
{"message-domain":"5678","id":4,"aaa":"bbb"}
{"message-domain":"1234","id":3,"aaa":"bbb"}

然后在输出中,在对这些事件进行分组/会话化之后,根据它们的消息域值,我将只有两个事件。

输出:

{"message-domain":"1234",messsages:[{"id":1,"aaa":"bbb"},{"id":2,"aaa":"bbb"},{"id":4,"aaa":"bbb"}]}
{"message-domain":"5678",messsages:[{"id":3,"aaa":"bbb"}]}

我希望这能实时发生。关于如何实现这一点的任何建议?

【问题讨论】:

    标签: akka-stream reactive-streams akka-kafka


    【解决方案1】:

    要在一个时间窗口内对事件进行分组,您可以使用Flow.groupedWithin

    val maxCount : Int = Int.MaxValue
    
    val timeWindow = FiniteDuration(2L, TimeUnit.MINUTES)
    
    val timeWindowFlow : Flow[String, Seq[String]] =
      Flow[String] groupedWithin (maxCount, timeWindow)
    

    【讨论】:

    • 如何在这里根据特定的 id 对事件进行分组?因此具有相同消息域值的事件将在分组后作为单个事件结束。
    • @dks551 这部分问题涉及更多。它需要大量的字符串操作才能获得您正在寻找的格式。它是否可以返回一个带有键 1234 和值 Seq[String] 的字典,其中序列中的值是字符串?或者,将每个字符串转换为一个案例类?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-10-22
    • 1970-01-01
    • 2019-01-12
    • 2015-07-31
    • 1970-01-01
    相关资源
    最近更新 更多