Parquet转Kafka SourceRecord时Avro转Connect Schema遇名称匹配错误
解决AvroData转换Decimal类型时的名称不匹配异常
问题根源
Guidewire生成的Parquet文件对应的Avro Schema中,decimal逻辑类型的fixed类型名称被设置为了字段名(如示例中的cedingrate),但Confluent的AvroData类处理decimal类型时,会强制要求fixed类型的名称为org.apache.kafka.connect.data.Decimal,两者不一致导致抛出Mismatched names异常。
可行解决办法
1. 提前修改Avro Schema再执行转换
遍历原Avro Schema的所有字段,找到带decimal逻辑类型的fixed定义,将其名称替换为Kafka Connect期望的org.apache.kafka.connect.data.Decimal,再用修改后的Schema调用toConnectData方法。示例代码如下:
// 原Avro Schema和GenericRecord Schema originalSchema = genericRecord.getSchema(); GenericRecord record = ...; // 构建修改后的Schema SchemaBuilder modifiedSchemaBuilder = SchemaBuilder.record(originalSchema.getName()) .namespace(originalSchema.getNamespace()); for (Field field : originalSchema.getFields()) { Schema fieldSchema = field.schema(); if (fieldSchema.getType() == Schema.Type.UNION) { List<Schema> adjustedUnionTypes = new ArrayList<>(); for (Schema unionType : fieldSchema.getTypes()) { // 检测是否为带decimal逻辑类型的fixed类型 if (unionType.getType() == Schema.Type.FIXED && unionType.getLogicalType() instanceof DecimalLogicalType) { DecimalLogicalType decimalType = (DecimalLogicalType) unionType.getLogicalType(); // 重建fixed类型,替换名称为Kafka Connect要求的名称 Schema adjustedFixedSchema = SchemaBuilder.fixed("org.apache.kafka.connect.data.Decimal") .size(unionType.getFixedSize()) .precision(decimalType.getPrecision()) .scale(decimalType.getScale()) .build(); adjustedUnionTypes.add(adjustedFixedSchema); } else { adjustedUnionTypes.add(unionType); } } modifiedSchemaBuilder.field(field.name(), Schema.createUnion(adjustedUnionTypes)); } else { modifiedSchemaBuilder.field(field.name(), fieldSchema); } } Schema modifiedSchema = modifiedSchemaBuilder.endRecord(); // 使用修改后的Schema执行转换 AvroData avroData = new AvroData(1000); ConnectData connectData = avroData.toConnectData(modifiedSchema, record);
2. 自定义AvroData子类覆盖类型转换逻辑
继承AvroData类,重写toConnectSchema方法中处理fixed类型的逻辑,跳过名称检查或强制使用兼容的名称。示例代码:
public class CustomAvroData extends AvroData { public CustomAvroData(int maxSchemasPerSubject) { super(maxSchemasPerSubject); } @Override protected Schema toConnectSchema(Schema avroSchema, String fieldName) { // 处理decimal类型的fixed,跳过默认的名称校验逻辑 if (avroSchema.getType() == Schema.Type.FIXED && avroSchema.getLogicalType() instanceof DecimalLogicalType) { DecimalLogicalType decimalLogicalType = (DecimalLogicalType) avroSchema.getLogicalType(); return org.apache.kafka.connect.data.SchemaBuilder.decimal(decimalLogicalType.getPrecision(), decimalLogicalType.getScale()) .build(); } return super.toConnectSchema(avroSchema, fieldName); } }
使用时直接实例化自定义的CustomAvroData类即可完成转换。
3. 调整Guidewire生成Parquet的配置(若可行)
如果有权限修改Guidewire的云数据访问框架配置,可尝试调整其Avro Schema生成规则,让decimal类型对应的fixed名称直接使用org.apache.kafka.connect.data.Decimal,从根源避免名称不匹配问题。具体配置项需参考Guidewire官方文档。
内容的提问来源于stack exchange,提问作者kalosh
相关产品推荐
相关产品推荐

