【问题标题】:kstreams grouping on two fields to get the countkstreams 在两个字段上分组以获取计数
【发布时间】:2019-02-08 23:53:58
【问题描述】:

我们能否按两个字段(一个是键,另一个是值)分组并获取 kstreams 中的计数。

我想为每个 pid(key) 获取不同的 userid(value) 计数。groupByKey 不会给出不同的 userid。 我尝试使用 groupBy 而不是 groupByKey 但看到语法错误。有人可以帮忙吗?

   KStream<Integer, Integer> stream = events.map((key, value) -> new KeyValue<Integer, Integer>(value.getpid(), value.getUserId()));

   KGroupedStream<Integer, Integer> groupedStream = stream.groupByKey(Grouped.with(Serdes.Integer(), Serdes.Integer());

【问题讨论】:

  • 请说明您遇到的错误。
  • 我尝试在上面的 kgroupedstream 中将 groupbykey 更改为 groupby 并且错误是 kstream 不能应用于(org.apache.kafka.streams.kstream.Grouped&lt;java.lang.integer,java.lang.integer&gt;)。分组两个字段的正确方法是什么
  • 这是因为groupBy没有接受Grouped的重载方法

标签: apache-kafka apache-kafka-streams


【解决方案1】:

如果你想通过user-id和pid来计数,你可以将两者都作为Pojo放入key中:

KStream<UserPid, Integer> stream =
    events.selectKey((key, value) -> new UserPid(value.getpid(), value.getUserId()));
KGroupedStream<Integer, Integer> groupedStream =
    stream.groupByKey(Grouped.with(new UserPidSerde(), Serdes.Integer());

需要创建对应的POJO类UserPid和serde类UserPidSerde extends Serde&lt;UserPid&gt;

【讨论】:

  • 感谢 Matthias,我在 docs.confluent.io/4.0.0/streams/developer-guide/… 看到了一些 serdes 示例,我必须实现自定义 serde 才能使用 UserPidSerde?​​span>
  • 是的。您需要为 UserPid 类创建一个自定义 serde。
  • 在创建自定义 serde UserPidSerde 后,我看到错误“无法解析方法 'groupByKey(org.apache.kafka.streams.kstream.Grouped)' 这行 `KGroupedStream groupedStream = stream.groupByKey(Grouped.with(new UserPidSerde(), Serdes.Integer()));
  • 我的最终结果也应该是 pid ,count(distinct userid)。如果我使用 UserPid POJO,如何做到这一点?
  • 好吧,我猜KGroupedStream&lt;Integer, Integer&gt; groupedStream 应该是KGroupedStream&lt; UserPid, Integer&gt; groupedStream
【解决方案2】:

由于每个 pid(key) 需要不同的 user(value) 计数,因此您需要首先使用 groupByKey,它将所有 users 与相同的 pid 分组。然后你需要聚合形成setuser(以获得唯一用户)。之后只需获取 set 的大小,您将获得每个 pid 的不同用户数。

KStream<Integer, Integer> stream = events.map((key, value) -> new KeyValue<Integer, Integer>(value.getpid(), value.getUserId()));
KStream<Integer, Integer> output = stream.groupByKey().
            aggregate((Initializer<Set<Integer>>) HashSet::new,
                    (k, v, current) -> {current.add(v); return current;}).mapValues(Set::size).toStream();

【讨论】:

    猜你喜欢
    • 2023-04-05
    • 1970-01-01
    • 1970-01-01
    • 2014-12-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-10-14
    • 1970-01-01
    相关资源
    最近更新 更多