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

无Schema下Avro记录反序列化:Schema嵌入与提取方案咨询

Avro Schema嵌入记录与提取方案

1. 发送端可将Schema嵌入Avro记录,实现方式如下

你可以通过自定义包裹结构的方式,将业务Schema和编码后的业务数据打包成一个新的Avro记录发送,这种方式直观且无需依赖Schema Registry,适合接收端无预存Schema的场景。

示例代码(Java)

首先定义用于包裹的通用Schema:

{
  "type": "record",
  "name": "SchemaWrappedRecord",
  "fields": [
    {"name": "originalSchema", "type": "string"},
    {"name": "data", "type": "bytes"}
  ]
}

发送端实现代码:

import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericDatumWriter;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.io.BinaryEncoder;
import org.apache.avro.io.DatumWriter;
import org.apache.avro.io.EncoderFactory;
import java.io.ByteArrayOutputStream;

public class AvroSender {
    public static void main(String[] args) throws Exception {
        // 1. 定义业务数据的原始Schema
        String userSchemaStr = "{\"type\":\"record\",\"name\":\"User\",\"fields\":[{\"name\":\"id\",\"type\":\"int\"},{\"name\":\"name\",\"type\":\"string\"}]}";
        Schema userSchema = new Schema.Parser().parse(userSchemaStr);

        // 2. 构造业务数据记录
        GenericRecord userRecord = new GenericData.Record(userSchema);
        userRecord.put("id", 1001);
        userRecord.put("name", "Alice");

        // 3. 将业务数据编码为字节数组
        ByteArrayOutputStream dataOut = new ByteArrayOutputStream();
        DatumWriter<GenericRecord> dataWriter = new GenericDatumWriter<>(userSchema);
        BinaryEncoder dataEncoder = EncoderFactory.get().binaryEncoder(dataOut, null);
        dataWriter.write(userRecord, dataEncoder);
        dataEncoder.flush();
        byte[] userDataBytes = dataOut.toByteArray();

        // 4. 解析包裹用的Schema
        Schema wrapperSchema = new Schema.Parser().parse(
            "{\"type\":\"record\",\"name\":\"SchemaWrappedRecord\",\"fields\":[{\"name\":\"originalSchema\",\"type\":\"string\"},{\"name\":\"data\",\"type\":\"bytes\"}]}"
        );

        // 5. 构造包裹记录,嵌入原始Schema和业务数据
        GenericRecord wrappedRecord = new GenericData.Record(wrapperSchema);
        wrappedRecord.put("originalSchema", userSchemaStr);
        wrappedRecord.put("data", userDataBytes);

        // 6. 编码包裹记录为最终发送的字节流
        ByteArrayOutputStream finalOut = new ByteArrayOutputStream();
        DatumWriter<GenericRecord> wrapperWriter = new GenericDatumWriter<>(wrapperSchema);
        BinaryEncoder finalEncoder = EncoderFactory.get().binaryEncoder(finalOut, null);
        wrapperWriter.write(wrappedRecord, finalEncoder);
        finalEncoder.flush();
        byte[] finalBytes = finalOut.toByteArray();

        // 此处将finalBytes发送至接收端(如Kafka、HTTP等)
    }
}

2. 接收端可从记录中提取Schema并解析数据

接收端只需先解码包裹记录,提取出原始Schema字符串,再用该Schema解码业务数据即可。

示例代码(Java)

import org.apache.avro.Schema;
import org.apache.avro.generic.GenericDatumReader;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.io.BinaryDecoder;
import org.apache.avro.io.DecoderFactory;
import java.io.ByteArrayInputStream;

public class AvroReceiver {
    public static void main(String[] args) throws Exception {
        // 替换为实际接收到的字节数组
        byte[] receivedBytes = ...;

        // 1. 解析包裹用的Schema(与发送端一致)
        Schema wrapperSchema = new Schema.Parser().parse(
            "{\"type\":\"record\",\"name\":\"SchemaWrappedRecord\",\"fields\":[{\"name\":\"originalSchema\",\"type\":\"string\"},{\"name\":\"data\",\"type\":\"bytes\"}]}"
        );

        // 2. 解码包裹记录
        GenericDatumReader<GenericRecord> wrapperReader = new GenericDatumReader<>(wrapperSchema);
        ByteArrayInputStream in = new ByteArrayInputStream(receivedBytes);
        BinaryDecoder decoder = DecoderFactory.get().binaryDecoder(in, null);
        GenericRecord wrappedRecord = wrapperReader.read(null, decoder);

        // 3. 提取并解析原始业务Schema
        String originalSchemaStr = wrappedRecord.get("originalSchema").toString();
        Schema originalSchema = new Schema.Parser().parse(originalSchemaStr);

        // 4. 解码业务数据
        byte[] dataBytes = (byte[]) wrappedRecord.get("data");
        GenericDatumReader<GenericRecord> dataReader = new GenericDatumReader<>(originalSchema);
        ByteArrayInputStream dataIn = new ByteArrayInputStream(dataBytes);
        BinaryDecoder dataDecoder = DecoderFactory.get().binaryDecoder(dataIn, null);
        GenericRecord userRecord = dataReader.read(null, dataDecoder);

        // 输出解析结果
        System.out.println("用户ID: " + userRecord.get("id"));
        System.out.println("用户名: " + userRecord.get("name"));
    }
}

补充说明

如果你的场景允许依赖Schema Registry,也可以使用Avro官方的Single Object Encoding:发送端在数据前添加魔术字节(0xC3)和Schema的CRC-64指纹,接收端通过指纹从Registry获取Schema。但如果接收端完全无Schema来源,自定义包裹结构是更直接的方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:20:34