【发布时间】:2018-01-20 16:01:27
【问题描述】:
我正在尝试将我的自定义类型的 ProducerRecords 发送到 Kafka,但我收到了错误:
Caused by: org.apache.kafka.common.errors.SerializationException: Error serializing Avro message
Caused by: java.lang.IllegalArgumentException: Unsupported Avro type. Supported types are null, Boolean, Integer, Long, Float, Double, String, byte[] and IndexedRecord
我在 Schema 中设置了架构: 获取
http://localhost:8081/subjects/documentCreations-key/versions/3
回复:
{
"subject": "documentCreations-key",
"version": 3,
"id": 1,
"schema": "\"string\""}
获取
http://localhost:8081/subjects/documentCreations-value/versions/4
回应
{
"subject": "documentCreations-value",
"version": 4,
"id": 23,
"schema": "{\"type\":\"record\",\"name\":\"Document\",\"namespace\":\"com.bade\",\"fields\":[{\"name\":\"name\",\"type\":\"string\"},{\"name\":\"path\",\"type\":\"string\"}]}"
}
这是我的 Scala 课程:
class Document(val name: java.lang.String,
val title: java.lang.String,
val path: java.lang.String)
还有 KafkaProducer 的部分:
class MyKafkaProducer {
val props = new Properties()
props.put("bootstrap.servers", "localhost:9092")
props.put("key.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer")
props.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer")
props.put("schema.registry.url", "http://localhost:8081")
private val producer = new KafkaProducer[java.lang.String, Document](props)
def sendCreateDocumentMessage(document: Document): RecordMetadata = {
val documentRecord = new ProducerRecord[java.lang.String, Document](SharedConfig
.documentCreationsTopic,
document.name, document)
producer.send(documentRecord).get()
}
我错过了什么?我看到我可以为我的班级实施SpecificRecord,但我认为在我一直在阅读的书籍/教程中没有必要这样做。 谢谢!
已编辑:固定类名
【问题讨论】:
-
您显示
MoreDocument,但您发送的是Document... -
已编辑,抱歉。我更改了命名,以免干扰业务逻辑。
-
那么,这段代码的哪一部分将您的案例类转换为 Avro?
-
我认为 KafkaAvroSerializer 可以做到这一点,可能是通过反射和在模式中提供类型。好的,我去看看,谢谢。
标签: scala apache-kafka avro