无需Schema Registry:Kafka Avro序列化器如何随Payload发送Schema
Kafka Avro序列化器/反序列化器(无Schema Registry)实现方案
因为要支持Schema演进且不使用Schema Registry,你需要自定义Kafka消息的字节流格式,把Writer Schema和Payload打包在一起发送,让消费端能按约定拆分并完成反序列化。以下是具体实现方案:
自定义消息字节流格式
约定固定的字节结构,确保消费端能准确拆分Schema和Payload:
- 魔数(2字节):自定义唯一标识,用于验证消息格式(例如
0xAB 0xCD) - Schema长度(4字节,大端字节序):记录后续Writer Schema的JSON字节长度
- Writer Schema字节:Avro Schema序列化后的JSON字节数组
- Payload字节:用Writer Schema序列化后的业务数据二进制流
Producer端序列化器实现(基于Avro生成类)
利用Avro生成类自带的SCHEMA$字段获取Writer Schema,按约定格式拼接字节流:
import org.apache.avro.Schema; import org.apache.avro.generic.GenericDatumWriter; import org.apache.avro.io.BinaryEncoder; import org.apache.avro.io.EncoderFactory; import org.apache.avro.specific.SpecificRecord; import org.apache.kafka.common.serialization.Serializer; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; public class CustomAvroSerializer<T extends SpecificRecord> implements Serializer<T> { // 自定义魔数,用于标识消息格式 private static final byte[] MAGIC_BYTES = new byte[]{(byte) 0xAB, (byte) 0xCD}; private static final int INT_LENGTH = 4; @Override public byte[] serialize(String topic, T data) { if (data == null) return null; // 获取Avro生成类的Writer Schema Schema writerSchema = data.getSchema(); byte[] schemaJsonBytes = writerSchema.toString().getBytes(StandardCharsets.UTF_8); int schemaLength = schemaJsonBytes.length; try (ByteArrayOutputStream outputStream = new ByteArrayOutputStream()) { // 写入魔数 outputStream.write(MAGIC_BYTES); // 写入Schema长度(大端字节序) outputStream.write(ByteBuffer.allocate(INT_LENGTH).putInt(schemaLength).array()); // 写入Writer Schema的JSON字节 outputStream.write(schemaJsonBytes); // 序列化并写入业务数据Payload GenericDatumWriter<T> datumWriter = new GenericDatumWriter<>(writerSchema); BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(outputStream, null); datumWriter.write(data, encoder); encoder.flush(); return outputStream.toByteArray(); } catch (IOException e) { throw new RuntimeException("Avro序列化失败", e); } } }
Consumer端反序列化器实现(基于Avro生成类)
按约定格式拆分字节流,解析Writer Schema后,结合消费端的Reader Schema(生成类的SCHEMA$)完成反序列化:
import org.apache.avro.Schema; import org.apache.avro.generic.GenericDatumReader; import org.apache.avro.io.BinaryDecoder; import org.apache.avro.io.DecoderFactory; import org.apache.avro.specific.SpecificRecord; import org.apache.kafka.common.serialization.Deserializer; import java.io.ByteArrayInputStream; import java.io.IOException; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.Arrays; public class CustomAvroDeserializer<T extends SpecificRecord> implements Deserializer<T> { private static final byte[] MAGIC_BYTES = new byte[]{(byte) 0xAB, (byte) 0xCD}; private static final int MAGIC_LENGTH = MAGIC_BYTES.length; private static final int INT_LENGTH = 4; private final Class<T> targetClass; public CustomAvroDeserializer(Class<T> targetClass) { this.targetClass = targetClass; } @Override public T deserialize(String topic, byte[] data) { if (data == null) return null; try (ByteArrayInputStream inputStream = new ByteArrayInputStream(data)) { // 验证魔数,确认是自定义Avro消息 byte[] magicBytes = new byte[MAGIC_LENGTH]; if (inputStream.read(magicBytes) != MAGIC_LENGTH || !Arrays.equals(magicBytes, MAGIC_BYTES)) { throw new RuntimeException("无效的Avro消息格式"); } // 读取Schema长度 byte[] schemaLengthBytes = new byte[INT_LENGTH]; if (inputStream.read(schemaLengthBytes) != INT_LENGTH) { throw new RuntimeException("读取Schema长度失败"); } int schemaLength = ByteBuffer.wrap(schemaLengthBytes).getInt(); // 读取并解析Writer Schema byte[] schemaJsonBytes = new byte[schemaLength]; if (inputStream.read(schemaJsonBytes) != schemaLength) { throw new RuntimeException("读取Schema内容失败"); } Schema writerSchema = new Schema.Parser().parse(new String(schemaJsonBytes, StandardCharsets.UTF_8)); // 获取消费端的Reader Schema(Avro生成类自带的SCHEMA$字段) Schema readerSchema = (Schema) targetClass.getField("SCHEMA$").get(null); // 反序列化业务数据Payload GenericDatumReader<T> datumReader = new GenericDatumReader<>(writerSchema, readerSchema); BinaryDecoder decoder = DecoderFactory.get().binaryDecoder(inputStream, null); return datumReader.read(null, decoder); } catch (IOException | NoSuchFieldException | IllegalAccessException e) { throw new RuntimeException("Avro反序列化失败", e); } } }
关键注意事项
- 魔数验证:必须添加魔数校验,避免消费端错误解析非目标格式的消息
- 字节序统一:Schema长度必须用大端字节序(网络字节序),保证跨平台兼容性
- Schema Resolution:Avro会自动处理Writer/Reader Schema的兼容问题,只要符合Avro官方的Schema解析规则
- 性能优化:每次发送Schema会增加消息体积,若Schema变更不频繁,可在Producer端缓存已发送的Schema(Consumer端需同步缓存),减少重复发送
- 生成类使用:Avro生成的类自带
SCHEMA$静态字段,直接调用即可获取对应Schema,无需手动编写
内容的提问来源于stack exchange,提问作者Pintu
相关产品推荐
相关产品推荐

