【问题标题】:Aggregate over multiple partitions in Kafka Streams在 Kafka Streams 中聚合多个分区
【发布时间】:2018-06-03 19:53:35
【问题描述】:

这部分是Aggregation over a specific partition in Apache Kafka Streams的后续行动

假设我有一个名为“事件”的主题,有 3 个分区,我在其上发送字符串 -> 整数数据,如下所示:

(Bob, 3) 在分区 1

(Sally, 4) 在分区 2

(Bob, 2) 在分区 3

...

我想汇总所有分区中的值(在此示例中,只是一个简单的总和),最终得到一个类似于以下内容的 KTable

(莎莉,4 岁)

(鲍勃,5)

正如我在上面链接到的问题的答案中提到的,直接进行这种跨分区聚合是不可能的。但是,回答者提到,如果消息具有相同的键(在这种情况下是正确的),则有可能。这如何实现?

我还希望能够从跨 Kafka Streams 应用程序的每个实例复制的“全局”状态存储中查询这些聚合值。

我的第一个想法是使用GlobalKTable(我相信,根据this page,这应该是我需要的)。但是,此状态存储的更改日志主题与原始“事件”主题具有相同数量的分区,并且只是基于每个分区而不是跨所有分区进行聚合。

这是我的应用程序的精简版 - 不太确定从哪里开始:

final Properties streamsConfig = new Properties();
streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "metrics-aggregator");
streamsConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
streamsConfig.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
streamsConfig.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, CustomDoubleSerde.class.getName());
streamsConfig.put(StreamsConfig.producerPrefix(ProducerConfig.LINGER_MS_CONFIG), 0);
streamsConfig.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 1);

final StreamsBuilder builder = new StreamsBuilder();

KStream<String, Double> eventStream = builder.stream(INCOMING_EVENTS_TOPIC);
KTable<String, Double> aggregatedMetrics = eventStream
        .groupByKey()
        .aggregate(() -> 0d, (key, value, aggregate) -> value + aggregate);

aggregatedMetrics.toStream().print(Printed.<String, Double>toSysOut());
aggregatedMetrics.toStream().to(METRIC_CHANGES_TOPIC);

final KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfig);
streams.cleanUp();
streams.start();

builder.globalTable(METRIC_CHANGES_TOPIC, Materialized.<String, Double, KeyValueStore<Bytes, byte[]>>as(METRICS_STORE_NAME));

Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    streams.close();
}));

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    Kafka Streams 假设输入主题是按键分区的。这个假设不适用于您的情况。因此,您需要告诉 Kafka Streams。

    在您的特定情况下,您可以将 groupByKey 替换为 groupBy()

    KTable<String, Double> aggregatedMetrics = eventStream
        .groupBy((k,v) -> k)
        .aggregate(() -> 0d, (key, value, aggregate) -> value + aggregate);
    

    lambda 是一个不修改键的虚拟变量,但是,它提示 Kafka Streams 在进行聚合之前根据键对数据进行重新分区。

    关于GlobalKTable:这是一种特殊的表,它不是聚合的结果,而是仅从变更日志主题中填充的。看来您的代码已经在做正确的事情了:将聚合结果写入主题并将主题重新读取为GlobalKTable

    【讨论】:

    • 谢谢你,成功了!我原以为groupByKeygroupBy((k,v) -&gt; k) 在功能上是相同的,但显然不是。再次感谢!
    • 不是。 groupByKey 告诉 Kafka Stream 使用当前密钥,如果密钥在读取输入后未更改,则允许避免重新分区步骤。使用groupBy 设置了一个新的key,并且不允许在groupByKey 之前尽可能优化。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-10-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-02
    相关资源
    最近更新 更多