【问题标题】:Use Kafka Streams for windowing data and processing each window at once使用 Kafka Streams 对数据进行窗口化并同时处理每个窗口
【发布时间】:2018-07-19 15:49:58
【问题描述】:

我想要实现的目的是按用户对我从 Kafka 主题收到的一些消息进行分组并将它们窗口化,以便汇总我在(5 分钟)窗口中收到的消息。然后我想收集每个窗口中的所有聚合,以便立即处理它们,并将它们添加到我在 5 分钟间隔内收到的所有消息的报告中。

最后一点似乎是困难的部分,因为 Kafka Streams 似乎没有提供(至少我找不到它!)任何可以在“有限”流中收集所有与窗口相关的内容以进行处理的东西在一个地方。

这是我实现的代码

StreamsBuilder builder = new StreamsBuilder();
KStream<UserId, Message> messages = builder.stream("KAFKA_TOPIC");

TimeWindowedKStream<UserId, Message> windowedMessages =
        messages.
                groupByKey().windowedBy(TimeWindows.of(SIZE_MS));

KTable<Windowed<UserId>, List<Message>> messagesAggregatedByWindow =
        windowedMessages.
                aggregate(
                        () -> new LinkedList<>(), new MyAggregator<>(),
                        Materialized.with(new MessageKeySerde(), new MessageListSerde())
                );

messagesAggregatedByWindow.toStream().foreach((key, value) -> log.info("({}), KEY {} MESSAGE {}",  value.size(), key, value.toString()));

KafkaStreams streams = new KafkaStreams(builder.build(), config);
streams.start();

结果是这样的

KEY [UserId(82770583)@1531502760000/1531502770000] Message [Message(userId=UserId(82770583),message="a"),Message(userId=UserId(82770583),message="b"),Message(userId=UserId(82770583),message="d")]
KEY [UserId(77082590)@1531502760000/1531502770000] Message [Message(userId=UserId(77082590),message="g")]
KEY [UserId(85077691)@1531502750000/1531502760000] Message [Message(userId=UserId(85077691),message="h")]
KEY [UserId(79117307)@1531502780000/1531502790000] Message [Message(userId=UserId(79117307),message="e")]
KEY [UserId(73176289)@1531502760000/1531502770000] Message [Message(userId=UserId(73176289),message="r"),Message(userId=UserId(73176289),message="q")]
KEY [UserId(92077080)@1531502760000/1531502770000] Message [Message(userId=UserId(92077080),message="w")]
KEY [UserId(78530050)@1531502760000/1531502770000] Message [Message(userId=UserId(78530050),message="t")]
KEY [UserId(64640536)@1531502760000/1531502770000] Message [Message(userId=UserId(64640536),message="y")]

每个窗口都有许多日志行,它们与其他窗口混合在一起。

我想要的是这样的:

// Hypothetical implementation
windowedMessages.streamWindows((interval, window) -> process(interval, window));

方法过程类似于:

// Hypothetical implementation

void process(Interval interval, WindowStream<UserId, List<Message>> windowStream) {
// Create report for the whole window   
Report report = new Report(nameFromInterval());
    // Loop on the finite iterable that represents the window content
    for (WindowStreamEntry<UserId, List<Message>> entry: windowStream) {
        report.addLine(entry.getKey(), entry.getValue());
    }
    report.close();
}

结果将像这样分组(每个报告都是对我的回调的调用:void process(...))并且每个窗口的提交将在整个窗口被处理时提交:

Report 1:
    KEY [UserId(85077691)@1531502750000/1531502760000] Message [Message(userId=UserId(85077691),message="h")]

Report 2:
    KEY [UserId(82770583)@1531502760000/1531502770000] Message [Message(userId=UserId(82770583),message="a"),Message(userId=UserId(82770583),message="b"),Message(userId=UserId(82770583),message="d")]
    KEY [UserId(77082590)@1531502760000/1531502770000] Message [Message(userId=UserId(77082590),message="g")]
    KEY [UserId(73176289)@1531502760000/1531502770000] Message [Message(userId=UserId(73176289),message="r"),Message(userId=UserId(73176289),message="q")]
    KEY [UserId(92077080)@1531502760000/1531502770000] Message [Message(userId=UserId(92077080),message="w")]
    KEY [UserId(78530050)@1531502760000/1531502770000] Message [Message(userId=UserId(78530050),message="t")]
    KEY [UserId(64640536)@1531502760000/1531502770000] Message [Message(userId=UserId(64640536),message="y")]

Report 3
    KEY [UserId(79117307)@1531502780000/1531502790000] Message [Message(userId=UserId(79117307),message="e")]

【问题讨论】:

  • 如果你想得到“原始记录”,你可以实现一个窗口聚合,返回一个List&lt;KeyValue&gt;作为结果类型,然后应用实际计算。除了使用 DSL,还可以选择使用处理器 API。
  • 谢谢马蒂亚斯。我试图做到这一点,但不幸的是窗户很大,所以把它们作为一个列表可能不安全。最好获得一些迭代器,它在窗口中的对象“列表”上进行迭代,以块的形式加载它们(将窗口视为有限主题)
  • 我明白了。这不会开箱即用。也许您需要使用处理器 API 实现自定义窗口运算符。

标签: java apache-kafka apache-kafka-streams windowing


【解决方案1】:

我也有同样的疑问。我已经与库的开发人员进行了交谈,他们说这是一个非常常见的请求,但尚未实现。它很快就会发布。

您可以在此处找到更多信息: https://cwiki.apache.org/confluence/display/KAFKA/KIP-328%3A+Ability+to+suppress+updates+for+KTables

【讨论】:

  • 谢谢@realBigfoot。我不确定它是否适合我读到的问题:-允许在kafka流中附加回调,当窗口过期时触发issues.apache.org/jira/browse/KAFKA-6556-对于我们的新抑制运算符来支持窗口最终结果,我们需要定义一个窗口结果实际上是最终结果的点!这听起来很接近我的需要,但我没有看到“立即”处理的窗口示例。我想要一个窗口流调用回调传递一个对象(一个可迭代的?)表示整个窗口中的聚合结果
  • 我刚刚更新了帖子,添加了我希望如何对数据进行分组的说明。
  • 与一位开发人员交谈过。我在这里引用他对我说的话:“在这种情况下,一种解决方法是在您知道不会在 punctuation 函数中对该存储应用更多更新后直接查询状态存储:请注意,标点符号是一项功能这仅在处理器 API 中可用,但您始终可以通过调用 KStream#process() / transform() 将这种较低级别的实现添加到您的 DSL 拓扑中。"
  • @Bruno 似乎您正在寻找的改进是可用的。我只是用它来实现类似的 (.suppress(Suppressed.untilWindowCloses(unbounded())))。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-08-21
  • 1970-01-01
  • 2018-11-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多