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

如何使用Kafka Streams统计GenericRecord中的去重值数量?

嘿,作为Kafka Streams新手,这个统计GenericRecord中值的去重数量的需求其实挺典型的,我来一步步给你拆解怎么做,保证你能快速上手~

核心思路

我们需要借助Kafka Streams的状态存储来跟踪已经出现过的唯一值,然后实时统计这个存储里的条目数量,最后输出成你要的JSON格式。因为GenericRecord通常对应Avro数据,所以还要处理好Avro的序列化/反序列化。

具体实现步骤

1. 配置Kafka Streams环境

首先得初始化Streams的配置,重点要指定Avro的序列化器和Schema Registry地址(因为Avro依赖Schema Registry来管理schema):

import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.common.serialization.Serdes;
import io.confluent.kafka.streams.serdes.avro.GenericAvroSerde;
import java.util.Properties;

Properties props = new Properties();
// 应用ID,必须唯一
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "distinct-value-count-app");
// Kafka集群地址
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 默认Key的序列化器(这里用String,你可以根据实际调整)
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// GenericRecord的序列化器/反序列化器
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, GenericAvroSerde.class);
// Schema Registry地址
props.put("schema.registry.url", "http://localhost:8081");

2. 构建流处理拓扑

接下来我们要创建流、定义状态存储、处理每条记录并统计去重数量:

import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.processor.Transformer;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.state.Stores;
import org.apache.kafka.streams.state.StoreBuilder;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.kstream.Produced;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.avro.generic.GenericRecord;
import java.util.concurrent.TimeUnit;

public class DistinctCountApp {
    public static void main(String[] args) {
        StreamsBuilder builder = new StreamsBuilder();

        // 1. 读取输入主题的Avro数据
        KStream<String, GenericRecord> inputStream = builder.stream("your-input-topic");

        // 2. 定义状态存储:用来保存已经出现过的唯一值(key是目标字段值,value用占位符)
        // 注意:这里的Serde要和你要去重的字段类型匹配,比如字段是String就用Serdes.String()
        StoreBuilder<KeyValueStore<Object, Boolean>> distinctStore =
                Stores.keyValueStoreBuilder(
                        Stores.persistentKeyValueStore("distinct-values-store"),
                        Serdes.Object(), // 示例用Object,实际替换成对应类型的Serde
                        Serdes.Boolean());
        builder.addStateStore(distinctStore);

        // 3. 用Transformer处理每条记录,更新状态并输出统计结果
        KStream<String, String> resultStream = inputStream.transform(
                () -> new Transformer<String, GenericRecord, KeyValue<String, String>>() {
                    private KeyValueStore<Object, Boolean> store;
                    private ProcessorContext context;

                    @Override
                    public void init(ProcessorContext context) {
                        this.context = context;
                        // 绑定状态存储
                        this.store = (KeyValueStore<Object, Boolean>) context.getStateStore("distinct-values-store");
                        // 可选:定时输出统计结果(比如每分钟一次),避免每条记录都输出
                        context.schedule(
                                java.time.Duration.ofMinutes(1),
                                org.apache.kafka.streams.processor.PunctuationType.WALL_CLOCK_TIME,
                                timestamp -> {
                                    long count = countStoreEntries(store);
                                    String outputJson = String.format("{\"Distinct values count\":%d}", count);
                                    context.forward("count-key", outputJson);
                                });
                    }

                    @Override
                    public KeyValue<String, String> transform(String key, GenericRecord value) {
                        // 替换成你要去重的字段名,比如"user_id"或者"product_code"
                        Object targetValue = value.get("your-target-field");
                        // 如果值不为空且未在存储中,就添加进去
                        if (targetValue != null && store.get(targetValue) == null) {
                            store.put(targetValue, Boolean.TRUE);
                        }
                        // 如果不需要每条记录都输出,这里返回null即可,依赖上面的定时任务输出
                        return null;
                    }

                    @Override
                    public void close() {}

                    // 辅助方法:统计状态存储中的条目数量
                    private long countStoreEntries(KeyValueStore<Object, Boolean> store) {
                        long count = 0;
                        var iterator = store.all();
                        while (iterator.hasNext()) {
                            iterator.next();
                            count++;
                        }
                        iterator.close();
                        return count;
                    }
                },
                "distinct-values-store"); // 指定使用的状态存储名称

        // 4. 将结果输出到目标主题
        resultStream.to("your-output-topic", Produced.with(Serdes.String(), Serdes.String()));

        // 启动流应用
        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();

        // 优雅关闭钩子
        Runtime.getRuntime().addShutdownHook(new Thread(() -> streams.close(5, TimeUnit.SECONDS)));
    }
}

关键注意事项

  • Serde匹配:上面代码里的Serdes.Object()是示例,你必须根据要去重的字段类型替换成对应的Serde,比如字段是String就用Serdes.String(),Integer用Serdes.Integer(),否则会出现序列化失败的问题。
  • 状态存储持久化:用persistentKeyValueStore可以保证应用重启后不丢失之前统计的唯一值,如果你不需要持久化,也可以用inMemoryKeyValueStore,但重启后数据会清空。
  • 输出时机:代码里提供了两种输出方式——每条记录处理后输出,或者定时输出。建议用定时输出,尤其是数据量大的时候,能减少输出主题的消息量。
  • Schema Registry依赖:确保你的环境中已经启动了Confluent Schema Registry,并且输入主题的Avro数据schema已经注册(可以配置自动注册:props.put("auto.register.schemas", "true"))。
  • 字段非空判断:处理GenericRecord时一定要加非空判断,避免NullPointerException。

内容的提问来源于stack exchange,提问作者mi.mo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:57:10