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

无需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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 11:10:21