.NET生产Kafka消息后,Java Kafka Streams反序列化失败如何解决?
问题描述
我用.NET的Confluent.Kafka库向名为ITEM_PRICES的主题生产键值对消息,键为string类型,值为float类型。
.NET生产者配置及生产代码
生产者序列化器配置:
_producer = new ProducerBuilder<string, float>(producerConfig) .SetKeySerializer(Serializers.Utf8) .SetValueSerializer(Serializers.Single) .Build();
生产消息的代码:
_producer.Produce(ITEM_PRICES, new Confluent.Kafka.Message<string, float> { Key= "ITEMNAME", Value = Convert.ToSingle(itemPrice)});
Java Kafka Streams消费尝试及错误
我在Java的Kafka Streams应用中消费这些消息,尝试的代码如下:
KStream<String, Float> itemPriceStream = streamsBuilder .stream("ITEM_PRICES", Consumed.with(Serdes.String(), Serdes.Float())); itemPriceStream .peek((key, value) -> System.out.println("Key: " + key + ", Value: " + value));
但运行时出现错误:Size of data received by Deserializer is not 4。
我的疑问:是否需要创建自定义反序列化器来读取这些消息?该如何实现?我尝试使用内置的String和Float反序列化器,但无法正常工作。
解决方案
需要自定义反序列化器,核心原因是.NET与Java的float序列化字节序不匹配:
- .NET的
Serializers.Single默认采用小端字节序序列化float值 - Java的
Serdes.Float()默认采用大端字节序(网络标准字节序)反序列化,导致字节解析逻辑不兼容,触发长度或解析异常
自定义Float反序列化器实现
创建适配.NET小端字节序的Float反序列化器:
import org.apache.kafka.common.serialization.Deserializer; import java.nio.ByteBuffer; import java.nio.ByteOrder; import java.util.Map; public class DotNetFloatDeserializer implements Deserializer<Float> { @Override public void configure(Map<String, ?> configs, boolean isKey) { // 无需额外配置 } @Override public Float deserialize(String topic, byte[] data) { if (data == null || data.length != 4) { throw new IllegalArgumentException("Expected 4 bytes for float, got " + (data == null ? 0 : data.length)); } // 按照小端字节序解析字节数组 ByteBuffer buffer = ByteBuffer.wrap(data).order(ByteOrder.LITTLE_ENDIAN); return buffer.getFloat(); } @Override public void close() { // 无需资源释放操作 } }
在Kafka Streams中使用自定义反序列化器
修改消费代码,替换内置Float Serde为自定义实现:
// 构建适配.NET的Float Serde Serde<Float> dotNetFloatSerde = Serdes.serdeFrom(new DotNetFloatDeserializer(), Serdes.Float().serializer()); KStream<String, Float> itemPriceStream = streamsBuilder .stream("ITEM_PRICES", Consumed.with(Serdes.String(), dotNetFloatSerde)); itemPriceStream.peek((key, value) -> System.out.println("Key: " + key + ", Value: " + value));
额外注意事项
- 如果后续需要从Java生产消息给.NET消费,需对应实现小端字节序的Float序列化器
- 生产侧要确保float值序列化后固定为4字节,避免因异常数据导致消费端报错
内容的提问来源于stack exchange,提问作者Katleho mokoena
相关产品推荐
相关产品推荐

