使用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
相关产品推荐
相关产品推荐

