如何从Java类访问Avro Schema自定义属性mapto并转换Avro文件
访问Avro自定义属性"mapto"并实现字段映射转换
Avro工具生成的Java类只会包含标准的Avro字段属性(name、type、default等),自定义属性不会直接生成到类中,需要通过Avro Schema对象来读取。以下是具体实现步骤和代码示例:
1. 加载原始Avro Schema
你可以通过生成的Java类自带的静态方法获取Schema,或者直接解析原始的.avsc文件:
// 方式1:通过生成的Java类获取Schema Schema productsSchema = Products.getClassSchema(); // 方式2:解析本地Schema文件(如果保留了avsc文件) // Schema.Parser parser = new Schema.Parser(); // Schema productsSchema = parser.parse(new File("products.avsc"));
2. 遍历Schema获取"mapto"属性
需要逐层遍历Schema的字段结构,拿到目标字段的自定义属性:
// 获取顶层的ProductDetails字段 Schema.Field productDetailsField = productsSchema.getField("ProductDetails"); // 处理联合类型:取非null的Record类型(联合类型第一个元素是null,索引为1的是目标Record) Schema productDetailsRecordSchema = productDetailsField.schema().getTypes().get(1); // 遍历子字段,读取mapto属性 for (Schema.Field field : productDetailsRecordSchema.getFields()) { String mapToFieldName = field.getProp("mapto"); System.out.printf("原字段%s -> 目标字段%s%n", field.name(), mapToFieldName); }
3. 构建输出Avro Schema
根据"mapto"的值生成输出用的Schema,确保输出字段名符合需求:
// 构建子Record的输出字段列表 List<Schema.Field> outputDetailFields = new ArrayList<>(); for (Schema.Field field : productDetailsRecordSchema.getFields()) { // 用mapto的值作为新字段名,保留原字段的类型、默认值等属性 Schema.Field outputField = new Schema.Field( field.getProp("mapto"), field.schema(), field.doc(), field.defaultValue(), field.order() ); outputDetailFields.add(outputField); } // 创建子Record的Schema Schema outputDetailRecordSchema = Schema.createRecord( "OutputProductDetailsRecord", null, "com.example.datasets", false, outputDetailFields ); // 构建顶层输出Schema List<Schema.Field> outputRootFields = new ArrayList<>(); Schema.Field outputRootField = new Schema.Field( "ProductDetails", // 顶层字段名可根据需求修改 Schema.createUnion(Schema.createNull(), outputDetailRecordSchema), null, null ); outputRootFields.add(outputRootField); Schema outputProductsSchema = Schema.createRecord( "OutputProducts", null, "com.example.datasets", false, outputRootFields );
4. 实现Avro文件转换
读取输入文件的SpecificRecord(生成的Java类实例),转换为输出的GenericRecord,最后写入新的Avro文件:
try ( // 读取输入Avro文件 DataFileReader<Products> inputReader = new DataFileReader<>( new File("input.avro"), new SpecificDatumReader<>(Products.getClassSchema()) ); // 写入输出Avro文件 DataFileWriter<GenericRecord> outputWriter = new DataFileWriter<>( new GenericDatumWriter<>(outputProductsSchema) ) ) { outputWriter.create(outputProductsSchema, new File("output.avro")); while (inputReader.hasNext()) { Products inputProduct = inputReader.next(); GenericRecord outputProduct = new GenericData.Record(outputProductsSchema); // 处理ProductDetails字段(可能为null) ProductDetailsRecord inputDetails = inputProduct.getProductDetails(); if (inputDetails != null) { GenericRecord outputDetails = new GenericData.Record(outputDetailRecordSchema); // 按mapto映射赋值 for (Schema.Field field : productDetailsRecordSchema.getFields()) { String targetFieldName = field.getProp("mapto"); Object fieldValue = inputDetails.get(field.name()); outputDetails.put(targetFieldName, fieldValue); } outputProduct.put("ProductDetails", outputDetails); } else { outputProduct.put("ProductDetails", null); } outputWriter.append(outputProduct); } } catch (IOException e) { e.printStackTrace(); }
关键说明
- 自定义属性通过
Schema.Field.getProp(String propName)方法获取,Avro会保留Schema中所有自定义键值对 - 联合类型需要注意索引:Avro联合类型中
null通常排在第一个,实际类型的索引为1 - 如果输出也需要生成对应的Java类,可以用avro-tools基于构建好的输出Schema重新编译
内容的提问来源于stack exchange,提问作者AlpsToronto
相关产品推荐
相关产品推荐

