【问题标题】:Configure Apache Kafka sink jdbc connector配置 Apache Kafka 接收器 jdbc 连接器
【发布时间】:2019-11-15 13:41:22
【问题描述】:

我想将发送到主题的数据发送到 postgresql 数据库。所以我关注this guide 并像这样配置了属性文件:

name=transaction-sink
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
tasks.max=1
topics=transactions
connection.url=jdbc:postgresql://localhost:5432/db
connection.user=db-user
connection.password=
auto.create=true
insert.mode=insert
table.name.format=transaction
pk.mode=none

我用

开始连接器
./bin/connect-standalone etc/schema-registry/connect-avro-standalone.properties etc/kafka-connect-jdbc/sink-quickstart-postgresql.properties

sink-connector 已创建,但由于此错误而无法启动:

Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing Avro message for id -1
Caused by: org.apache.kafka.common.errors.SerializationException: Unknown magic byte!

架构采用 avro 格式并已注册,我可以向主题发送(生成)消息并从中读取(使用)。但我似乎无法将其发送到数据库。

这是我的./etc/schema-registry/connect-avro-standalone.properties

key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://localhost:8081
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://localhost:8081

这是使用 java-api 提供主题的生产者:

properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
properties.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081");

try (KafkaProducer<String, Transaction> producer = new KafkaProducer<>(properties)) {
    Transaction transaction = new Transaction();
    transaction.setFoo("foo");
    transaction.setBar("bar");
    UUID uuid = UUID.randomUUID();
    final ProducerRecord<String, Transaction> record = new ProducerRecord<>(TOPIC, uuid.toString(), transaction);
    producer.send(record);
}

我正在验证数据是否正确序列化和反序列化使用

./bin/kafka-avro-console-consumer --bootstrap-server localhost:9092 \
    --property schema.registry.url=http://localhost:8081 \
    --topic transactions \
    --from-beginning --max-messages 1

数据库已启动并正在运行。

【问题讨论】:

  • 您为 Producer 设置的 KEY_SERIALIZER_CLASS_CONFIG 是什么?
  • 更新了我的问题。 StringSerializer.class.
  • Bingo :) 我已经更新了我的答案。
  • 你太棒了。 :-) 有用。我必须转换消息,因为我使用 BigDecimal,它在 avro 中是字符串类型,java-class java.math.BigDecimal 来规避 STRUCT 类型没有映射到 SQL 错误。我会看一些链接。
  • 更改 key.converter 后,我尝试将主题下沉数据放到 postgresql 表中。一个看似简单的任务,因为我没有嵌套数据。但它会因错误“不支持的源数据类型:STRUCT”而停止。我看过docs.confluent.io/3.1.1/connect/connect-jdbc/docs/…stackoverflow.com/questions/44385722/…,但后者似乎无关。属性位于顶部。我可以编辑错误吗?

标签: apache-kafka apache-kafka-connect


【解决方案1】:

这是不正确的:

未知的魔术字节可能是由于 id 字段不属于架构的一部分

该错误意味着有关该主题的消息未使用 Schema Registry Avro 序列化程序进行序列化。

您如何为该主题提供数据?

也许所有消息都有问题,也许只有一些 - 但默认情况下,这将停止 Kafka Connect 任务。

你可以设置

"errors.tolerance":"all",

让它忽略无法反序列化的消息。但是,如果它们都没有正确地 Avro 序列化,这将无济于事,您需要正确地序列化它们,或者选择不同的转换器(例如,如果它们实际上是 JSON,请使用 JSONConverter)。

这些参考资料应该对您有更多帮助:


编辑:

如果您使用StringSerializer 序列化密钥,那么您需要在您的连接配置中使用它:

key.converter=org.apache.kafka.connect.storage.StringConverter

您可以在 worker 中设置它(全局属性,适用于您在其上运行的所有连接器),或者仅针对此连接器(即,将其放在连接器属性本身中,它将覆盖 worker 设置)

【讨论】:

  • 我正在使用 java 流 api 为主题生成数据,模式位于 avro 文件中,并且该类是使用“mvn generate-sources”命令创建的。所以我可以向主题发送消息并阅读它们。所以现在我正在尝试将它发送到 postgres。
  • 将 errors.tolerance=all 添加到属性文件时,错误消失但没有消息发送到数据库。生产者正在发送到主题(通过网络界面检查主题)。
  • 您真的在使用 Schema Registry Avro 序列化程序吗?如果不是,则 Kafka Connect 无法读取您的消息。您必须使用序列化程序对它们进行序列化。 docs.confluent.io/current/schema-registry/…
  • 感谢您的宝贵时间,非常感谢。是的,我使用 avro-serializer。我用生产者代码更新了问题。
猜你喜欢
  • 2019-06-17
  • 2020-01-11
  • 2018-05-01
  • 2021-05-07
  • 2018-02-06
  • 2020-08-08
  • 2019-08-04
  • 2021-03-01
  • 2017-10-12
相关资源
最近更新 更多