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

