【问题标题】:Kafka: Error serializing Avro message with Schema RegistryKafka:使用 Schema Registry 序列化 Avro 消息时出错
【发布时间】: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


【解决方案1】:

回答我自己的问题。显然,(反)序列化不是自动完成的(通过反射或其他方式),但您必须从 avro 模式文件生成类。如果对某人有帮助,请发布我的pom.xml:

<build>
    <sourceDirectory>src/main/scala</sourceDirectory>
    <testSourceDirectory>src/test/scala</testSourceDirectory>

    <plugins>

        <!--force java 8-->
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-compiler-plugin</artifactId>
            <version>3.6.1</version>
            <configuration>
                <source>1.8</source>
                <target>1.8</target>
            </configuration>
        </plugin>

        <plugin>
            <groupId>net.alchim31.maven</groupId>
            <artifactId>scala-maven-plugin</artifactId>
            <version>3.3.1</version>
        </plugin>

        <plugin>
            <!-- Build an executable JAR -->
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-jar-plugin</artifactId>
            <version>3.0.2</version>
            <configuration>
                <archive>
                    <manifest>
                        <mainClass>Main</mainClass>
                    </manifest>
                </archive>
            </configuration>
        </plugin>

        <plugin>
            <groupId>org.apache.avro</groupId>
            <artifactId>avro-maven-plugin</artifactId>
            <version>${avro.version}</version>
            <executions>
                <execution>
                    <phase>generate-sources</phase>
                    <goals>
                        <goal>schema</goal>
                        <goal>protocol</goal>
                        <goal>idl-protocol</goal>
                    </goals>
                    <configuration>
                        <sourceDirectory>src/main/avro
                        </sourceDirectory>
                    </configuration>

                </execution>
            </executions>
        </plugin>
        <!--force discovery of generated classes-->
        <plugin>
            <groupId>org.codehaus.mojo</groupId>
            <artifactId>build-helper-maven-plugin</artifactId>
            <version>3.0.0</version>
            <executions>
                <execution>
                    <id>add-source</id>
                    <phase>generate-sources</phase>
                    <goals>
                        <goal>add-source</goal>
                    </goals>
                    <configuration>
                        <sources>
                            <source>target/generated-sources/avro</source>
                        </sources>
                    </configuration>
                </execution>
            </executions>
        </plugin>

    </plugins>
</build>

<repositories>
    <repository>
        <id>confluent</id>
        <url>http://packages.confluent.io/maven/</url>
    </repository>
</repositories>

<properties>
    <kafka.version>1.0.0</kafka.version>
    <confluent.version>4.0.0</confluent.version>
    <avro.version>1.8.2</avro.version>
</properties>

<dependencies>

    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka_2.12</artifactId>
        <version>${kafka.version}</version>
    </dependency>

    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>${kafka.version}</version>
    </dependency>

    <dependency>
        <groupId>io.confluent</groupId>
        <artifactId>kafka-avro-serializer</artifactId>
        <version>${confluent.version}</version>
    </dependency>

    <dependency>
        <groupId>org.apache.avro</groupId>
        <artifactId>avro</artifactId>
        <version>${avro.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.avro</groupId>
        <artifactId>avro-maven-plugin</artifactId>
        <version>${avro.version}</version>
    </dependency>

</dependencies>

我使用以下 mvn 命令构建它:

mvn clean:clean avro:schema compiler:compile scala:compile jar:jar

【讨论】:

  • 我在制作记录时遇到了同样的错误。因为我没有 avdl 或 avpr 文件,所以我从 pom 中删除了 protocol 和 idl-protocol 目标是否重要。没有它们,Java 类就可以生成。生成的 POJO 已经实现了SpecificRecord
猜你喜欢
  • 2020-06-26
  • 2018-01-31
  • 2020-08-11
  • 2017-12-02
  • 2021-06-13
  • 2019-10-24
  • 2016-09-16
  • 2019-07-12
  • 1970-01-01
相关资源
最近更新 更多