【问题标题】:Avro Schema with Kafka, ClassCastException?Avro Schema 与 Kafka,ClassCastException?
【发布时间】:2019-09-13 14:26:36
【问题描述】:

我为要发布到 Kafka 主题的记录创建了 Avro 模式。我们实际的 Kafka 记录模式更复杂,但为了简洁起见,我只是附上了相关部分。我们在记录中有多个嵌套子类,但由于某种原因,我在尝试发布记录时遇到以下异常(包名已被隐藏):

java.lang.ClassCastException: aaa.bbb.ccc.ddd.Amount cannot be cast to org.apache.avro.generic.IndexedRecord

class KafkaRecord {

    private Amount amount;

    class Amount {

        String currency;
        long value;

    }

}

这是我定义的 Avro 模式的当前子集。


{
  "type" : "record",
  "name" : "KafkaRecord",
  "namespace" : "com.company.department",
  "fields" : [ {
    "name" : "amount",
    "type" : {
      "type" : "record",
      "name" : "Amount",
      "namespace" : "aaa.bbb.ccc.ddd",
      "fields" : [ {
        "name" : "value",
        "type" : "long"
      }, {
        "name" : "currency",
        "type" : "string"
      } ]
    }
  }
}

我们的对象 (KafkaRecord) 的 JSON 表示如下所示:

{
  "amount": {
    "currency": "GBP",
    "value": 12345
  }
}

我似乎无法弄清楚为什么 Avro 不喜欢这个嵌套记录,而且我不想剥离这些嵌套类,因为这会使 JSON 记录非常难以阅读并且难以管理类.

如果有人能够指出我在这里做错了什么,那就太好了!

【问题讨论】:

  • 您是否使用了 SchemaRegistry(例如:Confluent 的那个)?

标签: java json apache-kafka schema avro


【解决方案1】:

好吧,您不需要编写自己的 Java 文件。听起来像你做的那样,所以错误就像它所说的那样 - Amount class is not an IndexedRecord

例如,如果我采用您的架构并运行

java -jar ~/Downloads/avro-tools-1.8.2.jar compile schema KafkaRecord.avsc .

然后查看文件,我们看到它扩展了一些 Avro java 类。

$ head -n 15 aaa/bbb/ccc/ddd/Amount.java
/**
 * Autogenerated by Avro
 *
 * DO NOT EDIT DIRECTLY
 */
package aaa.bbb.ccc.ddd;

import org.apache.avro.specific.SpecificData;
import org.apache.avro.message.BinaryMessageEncoder;
import org.apache.avro.message.BinaryMessageDecoder;
import org.apache.avro.message.SchemaStore;

@SuppressWarnings("all")
@org.apache.avro.specific.AvroGenerated
public class Amount extends org.apache.avro.specific.SpecificRecordBase implements org.apache.avro.specific.SpecificRecord {

Avro Maven Plugin 指南记录了以编程方式执行此操作的方法。

是的,Avro 可以很好地处理嵌套记录,我个人喜欢 using IDL to more easily create them

@namespace("com.company.department")
protocol KafkaEventProtocol {
  @namespace("aaa.bbb.ccc.ddd")
  record Amount {
    string currency;
    long value;
  }

  record KafkaValue {
    Amount amount;
  }
}

【讨论】:

    猜你喜欢
    • 2019-10-24
    • 2021-11-17
    • 1970-01-01
    • 2016-10-08
    • 2019-12-21
    • 2021-07-01
    • 2021-06-14
    • 1970-01-01
    • 2020-06-26
    相关资源
    最近更新 更多