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

.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:25:54