【问题标题】:Kafka Avro Serializer: org.apache.avro.AvroRuntimeException: not openKafka Avro 序列化程序:org.apache.avro.AvroRuntimeException:未打开
【发布时间】:2023-03-03 03:08:02
【问题描述】:

我正在使用带有 Avro Serializer 的 Apache Kafka,使用特定格式。我正在尝试创建自己的自定义类并用作 kafka 消息值。但是当我尝试发送消息时,我收到以下异常:

Exception in thread "main" org.apache.avro.AvroRuntimeException: not open
    at org.apache.avro.file.DataFileWriter.assertOpen(DataFileWriter.java:82)
    at org.apache.avro.file.DataFileWriter.append(DataFileWriter.java:287)
    at com.harmeetsingh13.java.producers.avroserializer.AvroProducer.main(AvroProducer.java:57)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at com.intellij.rt.execution.application.AppMain.main(AppMain.java:147)

我的 Avro Schema 文件如下:

{
    "namespace": "customer.avro",
    "type": "record",
    "name": "Customer",
    "fields": [{
        "name": "id",
        "type": "int"
    }, {
        "name": "name",
        "type": "string"
    }]
}

客户类别:

public class Customer {
    public int id;
    public String name;

    public Customer() {
    }

    public Customer(int id, String name) {
        this.id = id;
        this.name = name;
    }

    public int getId() {
        return id;
    }

    public void setId(int id) {
        this.id = id;
    }

    public String getName() {
        return name;
    }

    public void setName(String name) {
        this.name = name;
    }
}

使用 Avro 进行数据序列化:

public static void fireAndForget(ProducerRecord<String, DataFileWriter> record) {
        kafkaProducer.send(record);
    }

Customer customer1 = new Customer(1001, "James");

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

SpecificDatumWriter<Customer> writer = new SpecificDatumWriter<>(schema);
DataFileWriter<Customer> dataFileWriter = new DataFileWriter<>(writer);
dataFileWriter.append(customer1);
dataFileWriter.close();

ProducerRecord<String, DataFileWriter> record1 = new ProducerRecord<>("CustomerCountry",
        "Customer One", dataFileWriter
);
fireAndForget(record1);

我想使用SpecificDatumWriter writer 而不是通用的。这个错误与什么有关?

【问题讨论】:

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


    【解决方案1】:

    Kafka 收到一个要序列化的键值对,您传递给它的是 DataFileWriter,这不是您要序列化的值,这是行不通的。

    您需要做的是通过BinaryEncoderByteArrayOutputStream 使用序列化的avro 创建一个字节数组,然后将其传递给ProducerRecord&lt;String, byte[]&gt;

    SpecificDatumWriter<Customer> writer = new SpecificDatumWriter<>(schema);
    ByteArrayOutputStream os = new ByteArrayOutputStream();
    
    try {
      BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(os, null);
      writer.write(customer1, encoder);
      e.flush();
    
      byte[] avroBytes = os.toByteArray();
      ProducerRecord<String, byte[]> record1 = 
        new ProducerRecord<>("CustomerCountry", "Customer One", avroBytes); 
    
      kafkaProducer.send(record1);
    } finally {
      os.close();
    }
    

    【讨论】:

    • 嘿@Yuval,根据这个解决方案,我得到Exception in thread "main" java.lang.ClassCastException: com.harmeetsingh13.java.producers.avroserializer.Customer cannot be cast to org.apache.avro.generic.IndexedRecord Exception
    • @HarmeetSinghTaara 好吧,这是另一个问题,我建议您发布一个不同的问题,其中包括 Customer 的定义和完整的堆栈跟踪。似乎您没有实现要序列化的类型所需的 ISpecificRecord 接口。否则,您可以改用GenericDatumWriter
    • 好的@Yuval,但是根据这个解决方案,我们需要将我们的数据转换为字节,然后我们将消息发送到kafka,是否可以将我们的架构自动指定为avro和avro处理所有?
    • @HarmeetSinghTaara 不,不幸的是,Java 的 Avro 序列化程序不是这样工作的。
    • @HarmeetSinghTaara 酷。
    猜你喜欢
    • 2021-03-11
    • 1970-01-01
    • 2018-08-11
    • 2019-11-18
    • 2019-07-30
    • 2015-08-01
    • 2021-08-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多