You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在KStream中按双字段分组获取各PID对应的去重用户ID计数?

解决KStream中按PID分组统计去重用户ID数量的问题

我来帮你搞定这个需求——要统计每个PID对应的去重用户ID数量,groupByKey确实做不到自动去重,因为它只会按Key(PID)分组,把所有对应的Value(userId)都保留下来,重复的userId会被多次计数。另外你提到的groupBy语法错误,先看你给出的代码:

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

这里最后少了一个闭合括号,正确的写法应该是:

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

不过就算修正了语法,还是解决不了去重问题,下面给你两种可行的实现方案:

方案一:先去重再统计(无需自定义Serde)

这个思路是先把PID+userId作为复合Key,确保每个组合只出现一次,之后再拆分Key回到PID,统计每个PID对应的组合数量,也就是去重后的用户数:

// 第一步:生成PID+userId的复合Key,实现去重
KStream<String, Integer> deduplicatedUserStream = stream
    // 将Key转换为"PID:userId"的格式,确保每个用户在同一个PID下只保留一条记录
    .map((pid, userId) -> new KeyValue<>(pid + ":" + userId, userId))
    .groupByKey(Grouped.with(Serdes.String(), Serdes.Integer()))
    .count() // 每个复合Key计数为1,达到去重效果
    .toStream()
    // 拆分复合Key,将Value转为1,方便后续统计
    .map((compositeKey, count) -> {
        String pidStr = compositeKey.split(":")[0];
        Integer pid = Integer.parseInt(pidStr);
        return new KeyValue<>(pid, 1);
    });

// 第二步:按PID分组,统计去重后的用户数量
KTable<Integer, Long> uniqueUserCountPerPid = deduplicatedUserStream
    .groupByKey(Grouped.with(Serdes.Integer(), Serdes.Integer()))
    .sum();

方案二:使用聚合维护用户集合(更高效)

这种方式直接在聚合过程中维护每个PID对应的用户ID集合,最后返回集合的大小,性能更优,但需要自定义集合的Serde(因为Kafka默认不支持HashSet的序列化/反序列化):

第一步:自定义HashSet的Serde

public class HashSetSerde<T> implements Serde<HashSet<T>> {
    private final Serde<HashSet<T>> innerSerde;

    public HashSetSerde() {
        // 使用JSON序列化/反序列化HashSet
        innerSerde = Serdes.serdeFrom(new JsonSerializer<>(), new JsonDeserializer<>(HashSet.class));
    }

    @Override
    public Serializer<HashSet<T>> serializer() {
        return innerSerde.serializer();
    }

    @Override
    public Deserializer<HashSet<T>> deserializer() {
        return innerSerde.deserializer();
    }
}

第二步:聚合统计去重用户数

// 初始化状态:每个PID对应的空用户集合
Initializer<Set<Integer>> initializer = HashSet::new;

// 聚合逻辑:将新的userId加入集合
Aggregator<Integer, Integer, Set<Integer>> aggregator = (pid, userId, userSet) -> {
    userSet.add(userId);
    return userSet;
};

// 合并逻辑:用于多分区状态合并时,合并两个用户集合
Merger<Integer, Set<Integer>, Set<Integer>> merger = (pid, leftSet, rightSet) -> {
    leftSet.addAll(rightSet);
    return leftSet;
};

// 按PID分组聚合,最后映射为集合的大小(即去重用户数)
KTable<Integer, Integer> uniqueUserCountPerPid = stream
    .groupByKey(Grouped.with(Serdes.Integer(), Serdes.Integer()))
    .aggregate(
        initializer,
        aggregator,
        merger,
        // 指定状态存储的Serde
        Materialized.with(Serdes.Integer(), Serdes.serdeFrom(new HashSetSerde<>()))
    )
    .mapValues(Set::size);

两种方案对比

  • 方案一:实现简单,不需要额外写Serde,适合快速开发;缺点是多了一次分组聚合,性能略逊。
  • 方案二:直接在状态中维护用户集合,减少了一次分组操作,性能更优;但需要自定义Serde,代码量稍多。

内容的提问来源于stack exchange,提问作者user11036105

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 09:31:53