Kafka 1.0.0中Serialized.with()未用指定Serde默认用ByteArray问题
解决Kafka Streams 1.0.0中groupByKey Serialized.with不生效的问题
我之前在把Kafka Streams从0.10升级到1.0.0的时候,碰到过和你完全一样的问题!当时也是明明指定了Serialized.with(),但框架还是自动回退到ByteArraySerde,导致Long类型Key出现转换报错,折腾了好一阵才找到根源。
问题根源
Kafka 1.0.0的Streams API在链式调用场景下,Java的泛型推断偶尔会失效。虽然你在groupByKey()里传入了指定的Serde,但编译器没能正确识别泛型类型,导致框架没有读取到你配置的序列化规则,最终用了默认的ByteArraySerde。
解决方案:显式指定泛型参数
你只需要在Serialized.with()前显式加上泛型参数,强制编译器绑定正确的Key和Value类型,就能解决这个问题:
KTable<Long, myClass> myKTable = this.streamBuilder .stream(sub_topic, Consumed.with(Serdes.Long(), mySerde)) // 显式声明泛型参数<Long, myClass>,确保Serde类型被正确识别 .groupByKey(Serialized.<Long, myClass>with(Serdes.Long(), mySerde)) .reduce(myReducer, Materialized.as(my_store));
额外的保险措施
为了彻底避免序列化配置不匹配的问题,你还可以在Materialized中也显式指定Serde,让整个处理流程的序列化规则保持一致:
KTable<Long, myClass> myKTable = this.streamBuilder .stream(sub_topic, Consumed.with(Serdes.Long(), mySerde)) .groupByKey(Serialized.<Long, myClass>with(Serdes.Long(), mySerde)) .reduce(myReducer, Materialized.<Long, myClass, KeyValueStore<Bytes, byte[]>>as(my_store) .withKeySerde(Serdes.Long()) .withValueSerde(mySerde) );
我当时就是靠显式指定泛型参数解决了这个问题,你可以试试这个方法,应该能立刻解决你的Long类型Key转换错误。
内容的提问来源于stack exchange,提问作者Fizi
相关产品推荐
相关产品推荐

