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

使用ByteArrayDeserializer反序列化Kafka Avro消息为Generic/SpecificRecord遇错

问题:使用ByteArrayDeserializer反序列化Confluent Kafka Avro消息失败

背景

我通过Java Spring Kafka维护事件流,Kafka Topic发送的事件由Schema Registry管理。生产者配置如下,消息发送无问题:

props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);

正常消费配置

使用以下配置的Kafka消费者可正常接收并反序列化消息:

props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, io.confluent.kafka.serializers.KafkaAvroDeserializer.class);

我已获取Schema Registry生成的Java类,上述配置下反序列化完全正常。

问题场景

当前有一个消费者必须使用如下配置接收消息:

props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class);

无论采用哪种实现方式,都会抛出以下异常:

org.apache.avro.InvalidAvroMagicException: Not an Avro data file.
org.apache.avro.AvroRuntimeException: Malformed data. Length is negative: -55

我能获取生产者发送消息所用的Schema(即Schema Registry生成的对应Java类),但不管带Schema还是不带Schema反序列化,都会出现相同错误。

尝试过的失败实现

以下是几种尝试过但未成功的实现示例:

尝试1:反序列化为特定对象集合

// event 是 byte[] 类型
LOG.info("Recieved Event: {}", event);
LOG.debug("data='{}'", DatatypeConverter.printHexBinary(event));
ByteArrayInputStream in = new ByteArrayInputStream(event);
DatumReader<SpecificSchemaObject> userDatumReader = new SpecificDatumReader<>(SpecificSchemaObject.getClassSchema());
BinaryDecoder decoder = DecoderFactory.get().directBinaryDecoder(in, null);
List<SpecificSchemaObject> records = new ArrayList<SpecificSchemaObject>();
SpecificSchemaObject result = null;
try {
  while (true) {
    try {
      SpecificSchemaObject record = userDatumReader.read(null, decoder);
      records.add(record);
    } catch (EOFException eof) {
      break;
    }
  }
  result = (SpecificSchemaObject) userDatumReader.read(null, decoder);
  LOG.info("deserialized data='{}'", records);
  result = (SpecificSchemaObject) records;
} catch (IOException ioe) {

}

尝试2:反序列化为GenericRecord

Schema schema = WebsiteAdminEvent.getClassSchema();
DatumReader<GenericRecord> datumReader = new GenericDatumReader<GenericRecord>(schema);
Decoder decoder = DecoderFactory.get().directBinaryDecoder(new ByteArrayInputStream(event), null);
GenericRecord gr = null;
try {
  gr = datumReader.read(null, decoder);
} catch (IOException e) {
  throw new RuntimeException(e);
}
LOG.info("got record: ", gr);

尝试3:用DataFileReader读取

// event 是 byte[] 类型
List<WebsiteAdminEvent> listOfRecords = new ArrayList<>();
DatumReader<WebsiteAdminEvent> reader = new GenericDatumReader<>();
DataFileReader<WebsiteAdminEvent> fileReader =
    null;
try {
  fileReader = new DataFileReader<>(new SeekableByteArrayInput(event), reader);

  WebsiteAdminEvent record = null;

  while (fileReader.hasNext()) {
    listOfRecords.add(fileReader.next(record));
  }
} catch (IOException e) {
  throw new RuntimeException(e);
}

内容的提问来源于stack exchange,提问作者antonig

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 15:07:09