【问题标题】:Apache Kafka Avro Deserialization: Unable to deserialize or decode Specific type message.Apache Kafka Avro 反序列化:无法反序列化或解码特定类型的消息。
【发布时间】:2017-06-29 17:12:27
【问题描述】:

我正在尝试将 Avro Serialize 与 Apache kafka 一起使用来序列化/反序列化消息。我正在创建一个生产者,用于序列化特定类型的消息并将其发送到队列。当消息成功发送到队列时,我们的消费者选择消息并尝试处理,但在尝试时我们面临一个异常,对于特定对象的案例字节。例外情况如下:

[error] (run-main-0) java.lang.ClassCastException: org.apache.avro.generic.GenericData$Record cannot be cast to com.harmeetsingh13.java.avroserializer.Customer
java.lang.ClassCastException: org.apache.avro.generic.GenericData$Record cannot be cast to com.harmeetsingh13.java.avroserializer.Customer
    at com.harmeetsingh13.java.consumers.avrodesrializer.AvroSpecificDeserializer.lambda$infiniteConsumer$0(AvroSpecificDeserializer.java:51)
    at java.lang.Iterable.forEach(Iterable.java:75)
    at com.harmeetsingh13.java.consumers.avrodesrializer.AvroSpecificDeserializer.infiniteConsumer(AvroSpecificDeserializer.java:46)
    at com.harmeetsingh13.java.consumers.avrodesrializer.AvroSpecificDeserializer.main(AvroSpecificDeserializer.java:63)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)

根据异常,我们使用了一些不方便的方式来读取数据,下面是我们的代码:

Kafka生产者代码:

static {
        kafkaProps.put("bootstrap.servers", "localhost:9092");
        kafkaProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
        kafkaProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
        kafkaProps.put("schema.registry.url", "http://localhost:8081");
        kafkaProducer = new KafkaProducer<>(kafkaProps);
    }


public static void main(String[] args) throws InterruptedException, IOException {
        Customer customer1 = new Customer(1002, "Jimmy");

        Parser parser = new Parser();
        Schema schema = parser.parse(AvroSpecificProducer.class
                .getClassLoader().getResourceAsStream("avro/customer.avsc"));

        SpecificDatumWriter<Customer> writer = new SpecificDatumWriter<>(schema);
        try(ByteArrayOutputStream os = new ByteArrayOutputStream()) {
            BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(os, null);
            writer.write(customer1, encoder);
            encoder.flush();

            byte[] avroBytes = os.toByteArray();

            ProducerRecord<String, byte[]> record1 = new ProducerRecord<>("CustomerSpecificCountry",
                    "Customer One 11 ", avroBytes
            );

            asyncSend(record1);
        }

        Thread.sleep(10000);
    }

Kafka 消费者代码:

static {
        kafkaProps.put("bootstrap.servers", "localhost:9092");
        kafkaProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        kafkaProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        kafkaProps.put(ConsumerConfig.GROUP_ID_CONFIG, "CustomerCountryGroup1");
        kafkaProps.put("schema.registry.url", "http://localhost:8081");
    }

    public static void infiniteConsumer() throws IOException {
        try(KafkaConsumer<String, byte[]> kafkaConsumer = new KafkaConsumer<>(kafkaProps)) {
            kafkaConsumer.subscribe(Arrays.asList("CustomerSpecificCountry"));

            while(true) {
                ConsumerRecords<String, byte[]> records = kafkaConsumer.poll(100);
                System.out.println("<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<" + records.count());

                Schema.Parser parser = new Schema.Parser();
                Schema schema = parser.parse(AvroSpecificDeserializer.class
                        .getClassLoader().getResourceAsStream("avro/customer.avsc"));

                records.forEach(record -> {
                    DatumReader<Customer> customerDatumReader = new SpecificDatumReader<>(schema);
                    BinaryDecoder binaryDecoder = DecoderFactory.get().binaryDecoder(record.value(), null);
                    try {
                        System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>");
                        Customer customer = customerDatumReader.read(null, binaryDecoder);
                        System.out.println(customer);
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                });
            }

        }
    }

在控制台中使用消费者,我们可以成功接收到消息。那么将消息解码到我们的 pojo 文件中的方法是什么?

【问题讨论】:

    标签: java scala apache-kafka avro kafka-consumer-api


    【解决方案1】:

    这个问题的解决方法是,使用

    DatumReader<GenericRecord> customerDatumReader = new SpecificDatumReader<>(schema);
    

    而不是

    `DatumReader<Customer> customerDatumReader = new SpecificDatumReader<>(schema);
    

    这个的确切原因,仍然没有找到。这可能是因为 Kafka 不知道消息的结构,我们为消息显式定义了 schema,GenericRecord 用于将任何消息根据 schema 转换为可读的 JSON 格式。创建 JSON 后,我们可以轻松地将其转换为我们的 POJO 类。

    但是,仍然需要找到直接转换为我们的 POJO 类的解决方案。

    【讨论】:

      【解决方案2】:

      在将值传递给 ProduceRecord 之前,您不需要显式地执行 Avro 序列化。序列化程序将为您完成。您的代码如下所示:

      Customer customer1 = new Customer(1002, "Jimmy");
      ProducerRecord<String, Customer> record1 = new ProducerRecord<>("CustomerSpecificCountry", customer1);
          asyncSend(record1);
      }
      

      查看 Confluent 中的示例以获取 simple producer using avro

      【讨论】:

      • 嘿@Javier,首先认为我的问题与消费者有关,而不是与生产者有关。第二:如果您查看示例,JavaSessionize.avro.LogLine 看起来像 avro 类,所以他们可能会为此处理序列化。第三:我使用的是特定类型转换而不是泛型转换。 Avro 只支持 8 种类型,否则我们需要定义整个模式转换。
      • @HarmeetSinghTaara 我建议调查生产者,因为如果您序列化与预期不同的东西,消费者将会失败。我建议使用 Confluent 的 kafka-avro-console-consumer 来阅读该事件,看看它是否符合您的预期。如果这也失败了,那么问题不在于您的消费者。
      • @HarmeetSinghTaara LogLine 是通过this avro schema 自动生成Java 类的this avro schema 的结果。如果您编译项目,您可以查看并与您自己的Customer 类进行比较。我不认为他们在 LogLine 内部进行任何序列化,我很确定这只是一个 POJO。
      • 我正在使用 sbt 插件生成客户类。看看这个:github.com/sbt/sbt-avro
      • 嘿@Javier,我正在尝试遵循 avro 示例,但遇到异常。请查看stackoverflow.com/questions/42200875/…
      猜你喜欢
      • 2019-07-12
      • 2020-12-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-11
      • 2015-01-26
      • 2019-07-30
      相关资源
      最近更新 更多