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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 17:20:31