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

Kafka如何对包含ByteBuffer的自定义复杂对象实现序列化与反序列化

Kafka带ByteBuffer自定义对象的序列化/反序列化实现方案

1. 自定义序列化器(生产者侧)

序列化核心逻辑是按固定格式拼接字段内容:先存储String字段的字节长度+实际字节,再存储ByteBuffer的剩余可读长度+实际字节,避免反序列化时无法区分字段边界。

import org.apache.kafka.common.serialization.Serializer;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.Map;

public class NvrWSBinaryMessageSerializer implements Serializer<NvrWSBinaryMessage> {
    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        // 无额外配置需求可留空
    }

    @Override
    public byte[] serialize(String topic, NvrWSBinaryMessage data) {
        if (data == null) {
            return null;
        }
        // 处理String字段
        byte[] uuidBytes = data.getMessageUuid().getBytes(StandardCharsets.UTF_8);
        // 处理ByteBuffer字段,用duplicate()避免修改原对象的position/limit指针
        ByteBuffer payloadView = data.getPayload().duplicate();
        int payloadSize = payloadView.remaining();
        // 计算总字节长度:4字节存uuid长度 + uuid字节长度 + 4字节存payload长度 + payload字节长度
        int totalSize = 4 + uuidBytes.length + 4 + payloadSize;
        ByteBuffer outputBuffer = ByteBuffer.allocate(totalSize);
        // 按顺序写入字段
        outputBuffer.putInt(uuidBytes.length);
        outputBuffer.put(uuidBytes);
        outputBuffer.putInt(payloadSize);
        outputBuffer.put(payloadView);
        return outputBuffer.array();
    }

    @Override
    public void close() {
        // 无资源释放需求可留空
    }
}

2. 自定义反序列化器(消费者侧)

和序列化逻辑对应,按顺序读取对应长度的字节,还原为原对象字段即可。

import org.apache.kafka.common.serialization.Deserializer;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.Map;

public class NvrWSBinaryMessageDeserializer implements Deserializer<NvrWSBinaryMessage> {
    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        // 无额外配置需求可留空
    }

    @Override
    public NvrWSBinaryMessage deserialize(String topic, byte[] data) {
        if (data == null || data.length == 0) {
            return null;
        }
        ByteBuffer inputBuffer = ByteBuffer.wrap(data);
        // 还原messageUuid字段
        int uuidLen = inputBuffer.getInt();
        byte[] uuidBytes = new byte[uuidLen];
        inputBuffer.get(uuidBytes);
        String messageUuid = new String(uuidBytes, StandardCharsets.UTF_8);
        // 还原payload字段
        int payloadLen = inputBuffer.getInt();
        ByteBuffer payload = ByteBuffer.allocate(payloadLen);
        inputBuffer.get(payload.array(), 0, payloadLen);
        // 根据你自定义类的构造方法调整入参即可
        return new NvrWSBinaryMessage(messageUuid, payload);
    }

    @Override
    public void close() {
        // 无资源释放需求可留空
    }
}

3. Kafka配置调整

在生产者和消费者的配置中指定自定义的序列化/反序列化器即可:
生产者配置:

value.serializer=替换为你自己的包路径.NvrWSBinaryMessageSerializer

消费者配置:

value.deserializer=替换为你自己的包路径.NvrWSBinaryMessageDeserializer

注意事项

  • 序列化时调用duplicate()获取ByteBuffer的视图读取,不会修改原对象的指针位置,不会影响业务侧后续对该ByteBuffer的操作
  • 若视频分片大小超过默认阈值,需要同步调整Kafka生产者max.request.size、Broker端message.max.bytes、消费者端fetch.max.bytes配置,避免消息被拦截
  • 确保NvrWSBinaryMessage提供包含两个字段的构造方法,用于反序列化时初始化final修饰的payload字段

内容的提问来源于stack exchange,提问作者Jocheinfa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 13:45:04