【问题标题】:Aggregation (summation) of number of events from different Kafka topics来自不同 Kafka 主题的事件数量的聚合(总和)
【发布时间】:2020-01-26 23:22:31
【问题描述】:

我的应用程序有三个主题,它们接收一些属于用户的事件:

Event Type A -> Topic A
Event Type B -> Topic B
Event Type C -> Topic C

这将是一个消息流的例子:

Message(user 1 - event A - 2020-01-03) 
Message(user 2 - event A - 2020-01-03) 
Message(user 1 - event C - 2020-01-20)
Message(user 1 - event B - 2020-01-22)

我希望能够生成报告,其中包含每个用户每月的事件总数,聚合来自三个主题的所有事件,例如:

User 1 - 2020-01 -> 3 total events
User 2 - 2020-01 -> 1 total events

拥有三个 KStream(每个主题一个),我如何每月执行此加法以汇总来自三个不同主题的所有事件?你能展示一下代码吗?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams spring-cloud-stream spring-cloud-stream-binder-kafka


    【解决方案1】:

    因为您只对计数感兴趣,所以最简单的方法是将用户 ID 保留为键,并为每个 KStream 设置一些虚拟值,合并所有三个流并在之后进行窗口计数(注意开箱即用不支持基于日历的窗口;您可以使用 31 天的窗口作为近似值或构建自己的自定义窗口):

    // just map to dummy empty string (note, that `null` would not work
    KStream<UserId, String> streamA = builder.stream("topic-A").mapValues(v -> "");
    KStream<UserId, String> streamB = builder.stream("topic-B").mapValues(v -> "");
    KStream<UserId, String> streamC = builder.stream("topic-C").mapValues(v -> "");
    
    streamA.merge(streamB).merge(streamC).groupByKey().windowBy(...).count();
    

    您可能还对suppress() 运算符感兴趣。

    【讨论】:

    • 优秀的答案马蒂亚斯!您能否提供一些关于如何从每个月的第一天到最后一天实现基于日历的窗口的提示?我在其他场景中遇到过这种需求,但仍然不知道该怎么做。
    • 这可能会有所帮助:github.com/confluentinc/kafka-streams-examples/blob/5.4.0-post/… -- 它实现了不同时区的每日窗口(默认情况下,所有窗口都基于 UTC 时区)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-23
    • 2018-12-10
    • 1970-01-01
    • 2020-06-02
    • 2016-01-14
    • 1970-01-01
    • 2011-05-21
    相关资源
    最近更新 更多