如何使用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
相关产品推荐
相关产品推荐

