多Avro记录字节数组反序列化报错:ArrayList无法转为SpecificRecordBase
问题解决方案
问题根源在于你的反序列化器泛型参数T原本指向单个Avro对象(SpecificRecordBase子类),但修改后试图返回List<GenericRecord>并强制转换为T,这必然导致类型转换异常。以下是具体修正方案:
核心调整思路
- 明确反序列化器的返回类型为
List<T>,而非单个T对象 - 匹配泛型类型与Avro具体类,避免
GenericRecord和SpecificRecordBase的类型不兼容 - 消除不安全的强制类型转换,保证类型安全
修正后的完整代码
import org.apache.avro.specific.SpecificDatumReader; import org.apache.avro.specific.SpecificRecordBase; import org.apache.kafka.common.serialization.Deserializer; import org.apache.avro.io.DecoderFactory; import org.apache.avro.io.BinaryDecoder; import org.apache.avro.io.DatumReader; import java.io.ByteArrayInputStream; import java.io.EOFException; import java.util.ArrayList; import java.util.List; import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import javax.xml.bind.DatatypeConverter; public class MultiAvroDeserializer<T extends SpecificRecordBase> implements Deserializer<List<T>> { private static final Logger LOGGER = LoggerFactory.getLogger(MultiAvroDeserializer.class); private final Class<T> targetType; // 构造函数传入目标Avro类的Class对象 public MultiAvroDeserializer(Class<T> targetType) { this.targetType = targetType; } @Override public void configure(Map<String, ?> configs, boolean isKey) { // 按需添加配置逻辑 } @Override public List<T> deserialize(String topic, byte[] data) { if (data == null || data.length == 0) { return new ArrayList<>(); // 空数据返回空列表 } LOGGER.debug("data='{}'", DatatypeConverter.printHexBinary(data)); ByteArrayInputStream in = new ByteArrayInputStream(data); DatumReader<T> datumReader = new SpecificDatumReader<>(targetType.getSchema()); BinaryDecoder decoder = DecoderFactory.get().directBinaryDecoder(in, null); List<T> records = new ArrayList<>(); try { while (true) { try { T record = datumReader.read(null, decoder); records.add(record); } catch (EOFException e) { break; // 流读取完毕,退出循环 } } LOGGER.info("deserialized data='{}'", records); return records; } catch (Exception ex) { throw new org.apache.kafka.common.errors.SerializationException( "Can't deserialize data from topic '" + topic + "'", ex); } } @Override public void close() { // 按需添加资源清理逻辑 } }
关键修改点说明
- 泛型定义改为
Deserializer<List<T>>,明确返回类型为Avro对象列表 - 使用
Class<T>构造参数获取目标Avro类的Schema,替代反射实例化的潜在风险 - 将记录列表类型从
List<GenericRecord>改为List<T>,完全匹配泛型类型,消除类型转换异常 - 空数据统一返回空列表,保持行为一致性
使用示例
配置Kafka消费者时,指定该反序列化器并传入你的Avro生成类:
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, new MultiAvroDeserializer<>(YourAvroClass.class));
内容的提问来源于stack exchange,提问作者rockyboom
相关产品推荐
相关产品推荐

