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
相关产品推荐
相关产品推荐

