Avro嵌套数组decimal反序列化HeapByteBuffer转BigDecimal报错
在基于Confluent Schema Registry与Avro搭建数据链路时,通过JDBC连接器将数据接入Kafka,连接器配置了SMT(单消息转换)用于生成合规的Avro Schema;消费端使用SpecificAvroSerde进行消息反序列化时抛出异常。
此前落地的大量同类场景均运行正常,证明数据接入、Avro Schema生成、流处理消费Avro数据的整体链路方案可行。本次场景唯一差异为待消费记录包含主从结构的数组类型字段。
使用GenericRecord读取同一条消息可正常获取所有字段,证明Avro消息本身序列化逻辑正确。
{ "namespace": "io.confluent.base.model", "type": "record", "name": "Test1", "fields": [ { "name": "opt_identifier", "type": [ "null", "string" ],"default": null }, { "name": "opt_amount", "type": [ "null", { "type":"bytes", "logicalType":"decimal", "precision":31, "scale":8 }], "default": null}, { "name": "arr_field", "type": ["null", { "type": "array", "items": { "name": "TestTest1", "type": "record", "fields": [ { "name": "opt_identifier_", "type": [ "null", "string" ],"default": null }, { "name": "opt_amount_", "type": [ "null", { "type":"bytes", "logicalType":"decimal", "precision":31, "scale":8 }], "default": null} ] }, "default": [] }], "default": null} ] }
上述Schema通过avro-maven-plugin编译生成对应Java类,连接器与消费者端使用的Avro相关依赖jar包版本完全一致。
org.apache.kafka.common.errors.SerializationException: Error deserializing Avro message for id 79 at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer$DeserializationContext.read(AbstractKafkaAvroDeserializer.java:409) ~[kafka-avro-serializer-7.0.1.jar:na] at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:114) ~[kafka-avro-serializer-7.0.1.jar:na] at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:88) ~[kafka-avro-serializer-7.0.1.jar:na] at io.confluent.kafka.serializers.KafkaAvroDeserializer.deserialize(KafkaAvroDeserializer.java:55) ~[kafka-avro-serializer-7.0.1.jar:na] at io.confluent.kafka.streams.serdes.avro.SpecificAvroDeserializer.deserialize(SpecificAvroDeserializer.java:66) ~[kafka-streams-avro-serde-7.0.1.jar:na] at io.confluent.kafka.streams.serdes.avro.SpecificAvroDeserializer.deserialize(SpecificAvroDeserializer.java:38) ~[kafka-streams-avro-serde-7.0.1.jar:na] at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:60) ~[kafka-clients-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.SourceNode.deserializeValue(SourceNode.java:58) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.RecordDeserializer.deserialize(RecordDeserializer.java:66) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.RecordQueue.updateHead(RecordQueue.java:176) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.RecordQueue.addRawRecords(RecordQueue.java:112) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.PartitionGroup.addRawRecords(PartitionGroup.java:304) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.StreamTask.addRecords(StreamTask.java:960) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.TaskManager.addRecordsToTasks(TaskManager.java:1000) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.StreamThread.pollPhase(StreamThread.java:914) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:720) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:583) ~[kafka-streams-3.0.0.jar:na] at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:555) ~[kafka-streams-3.0.0.jar:na] Caused by: java.lang.ClassCastException: class java.nio.HeapByteBuffer cannot be cast to class java.math.BigDecimal (java.nio.HeapByteBuffer and java.math.BigDecimal are in module java.base of loader 'bootstrap') at io.confluent.base.model.TestTest1.put(TestTest1.java:416) ~[classes/:na] at org.apache.avro.generic.GenericData.setField(GenericData.java:818) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.specific.SpecificDatumReader.readField(SpecificDatumReader.java:139) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:247) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.specific.SpecificDatumReader.readRecord(SpecificDatumReader.java:123) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:179) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readArray(GenericDatumReader.java:298) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:183) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.specific.SpecificDatumReader.readField(SpecificDatumReader.java:136) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:247) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.specific.SpecificDatumReader.readRecord(SpecificDatumReader.java:123) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:179) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:160) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:187) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:160) ~[avro-1.10.1.jar:1.10.1] at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:153) ~[avro-1.10.1.jar:1.10.1] at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer$DeserializationContext.read(AbstractKafkaAvroDeserializer.java:400) ~[kafka-avro-serializer-7.0.1.jar:na] ... 17 common frames omitted
- 问题与Avro逻辑类型(logical type)的转换处理相关
- 主记录层级的decimal逻辑类型字段(如
opt_amount)可正常完成反序列化 - 嵌套数组内子记录
TestTest1的opt_amount_字段抛出类型转换异常,怀疑嵌套的明细子记录TestTest1未应用与主记录Test1相同的逻辑类型转换逻辑
这是Confluent 7.0.x版本SpecificAvroSerde的已知缺陷:SpecificDatumReader初始化时,只会为顶层Schema注册decimal等逻辑类型对应的转换器,不会递归遍历数组元素、Map值、嵌套子记录等深层Schema结构加载对应的Conversion实例。
读取顶层decimal字段时,转换器正常生效,将原始ByteBuffer转为BigDecimal;读取数组内嵌套子记录的decimal字段时,没有对应的转换器执行类型转换,直接返回原始HeapByteBuffer,子记录的put方法尝试将ByteBuffer赋值给BigDecimal类型的属性,就会抛出ClassCastException。GenericRecord读取正常是因为GenericDatumReader默认全局加载所有逻辑类型转换规则,不受该缺陷影响。
- 优先选择升级Confluent平台及相关Avro序列化依赖到7.1.0及以上版本,官方已修复该嵌套结构逻辑类型转换器递归加载的问题,升级后无需修改业务代码即可正常运行。
- 若暂时无法升级版本,可自定义
SpecificAvroSerde实现,重写反序列化阶段创建SpecificDatumReader的逻辑,递归遍历全量Schema结构(包括数组items、Map values、所有嵌套record字段),为所有携带logicalType的字段手动注册BigDecimalConversion等对应逻辑类型的转换器。 - 临时规避可在消费端先使用
GenericAvroSerde读取消息,手动将嵌套字段的ByteBuffer按decimal规则转为BigDecimal,再组装到SpecificRecord对象中,该方案代码侵入性较高,仅适合短期应急使用。
内容的提问来源于stack exchange,提问作者Thomas

