【问题标题】:Apache Kafka 1.0.0 Streams API Multiple Multilevel groupbyApache Kafka 1.0.0 Streams API 多级多级分组
【发布时间】:2018-02-22 12:11:31
【问题描述】:

如何在 Kafka Streams API 中使用具有多重约束的 .groupby。与下面的 Java 8 Streams API 示例相同

public void twoLevelGrouping(List<Person> persons) {
     final Map<String, Map<String, List<Person>>> personsByCountryAndCity = persons.stream().collect(
         groupingBy(Person::getCountry,
            groupingBy(Person::getCity)
        )
    );
    System.out.println("Persons living in London: " + personsByCountryAndCity.get("UK").get("London").size());
}

【问题讨论】:

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


    【解决方案1】:

    您可以通过将要分组的所有属性/字段放入键中来指定组合键。

    KTable table = stream.selectKey((k, v,) -> k::getCountry + "-" + k::getCity)
                         .groupByKey()
                         .aggregate(...); // or maybe .reduce()
    

    我只是假设国家和城市都是String。您使用交互式查询来查询商店

    store.get("UK-London");
    

    https://docs.confluent.io/current/streams/developer-guide/interactive-queries.html

    【讨论】:

    • 我不认为这是正确的方法,但是,我同意它可以完成这项工作。有没有更好的方法呢??
    • 答案是草图。最好引入一些包含两个字段的 POJO 类型,一个用于国家,一个用于城市(或者通常每个分组属性一个字段)并在selectKey() 中发出这种新类型——您还需要编写您需要在 groupByKeyaggregate() 中设置的类型的自定义 Serdes。聚合函数将获取 this POJO 作为输入类型来进行聚合。当然,POJO 只是其中一种方式,你也可以使用 Tuple 类型或任何其他组合类型。
    • 我已经使用了自定义 POJO,以及带有 GSON 的自定义 serde。你能解释一下吗?
    • 不确定我应该解释什么...与其将两个属性放入单个字符串中,不如发出包装两个属性的 POJO。这将自动按两个属性分组。
    • 对不起,我是菜鸟,我真的可以举个例子!
    猜你喜欢
    • 2017-10-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-11
    • 2021-10-14
    相关资源
    最近更新 更多