Kafka Streams Twitter词频统计:序列化后计数非Long类型问题
问题排查与解决
问题根源
你遇到的Size of data received by LongDeserializer is not 8错误,本质是Kafka Streams写入输出主题的计数结果不是标准Long类型的二进制数据,而是被序列化为了字符串格式。
原因出在count()操作的Materialized配置上:你的全局默认值Serde设置为StringSerde(props.put(DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass())),而count()返回的计数是Long类型,当你调用count(Materialized.as("WordCount"))时没有显式指定值的Serde,Kafka Streams会使用全局默认的StringSerde来序列化Long类型的计数,最终写入主题的是字符串的字节(比如"1"的字节长度是1,远小于Long的8字节),导致消费者用LongDeserializer反序列化时出现字节长度不匹配。
解决方案
方法一:显式指定Materialized的Value Serde(推荐)
修改count()方法的调用,为状态存储指定Long类型的Value Serde,确保计数结果以正确的Long格式序列化:
.count(Materialized.<String, Long>as("WordCount").withValueSerde(Serdes.Long()))
修改后的完整处理链代码:
textLines .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) .groupBy((key, value) -> value) .count(Materialized.<String, Long>as("WordCount").withValueSerde(Serdes.Long())) .toStream() .to("test-output", Produced.with(Serdes.String(), Serdes.Long()));
方法二:确认消费者命令的反序列化配置
确保你使用kafka-console-consumer消费时,指定了正确的反序列化器,命令示例:
kafka-console-consumer.sh --bootstrap-server ec2-ip:9092 --topic test-output --from-beginning \ --key-deserializer org.apache.kafka.common.serialization.StringDeserializer \ --value-deserializer org.apache.kafka.common.serialization.LongDeserializer
验证
修改代码重新部署Kafka Streams应用后,重新消费test-output主题,即可正常获取Long类型的计数结果。
内容的提问来源于stack exchange,提问作者chris_vee
相关产品推荐
相关产品推荐

