【问题标题】:Kafka Connect export multiple event types from same topicKafka Connect 从同一主题导出多种事件类型
【发布时间】:2018-11-02 11:23:48
【问题描述】:

我正在尝试使用一项新功能 (https://www.confluent.io/blog/put-several-event-types-kafka-topic/) 来存储两个 同一主题的不同类型的事件。实际上我使用的是 Confluent 4.1.0 版并在下面设置这些属性 实现这一目标

properties.put(KafkaAvroSerializerConfig.VALUE_SUBJECT_NAME_STRATEGY,TopicRecordNameStrategy.class.getName());
properties.put("value.multi.type", true);

数据写入主题没有问题,并且可以从 Kafka Streams 应用程序中看到作为通用 Avro 记录。还 在 Kafka Schema 注册表中,为在该特定主题上托管的每个事件创建了两个新条目。

我面临的问题是我无法使用 Kafka Connect 从该主题导出这些数据。在最简单的情况下 我使用如下文件接收器连接器

{
  "name": "sink-connector",
  "config": {
      "topics": "source-topic",
      "connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
      "tasks.max": 1,
      "key.converter": "org.apache.kafka.connect.storage.StringConverter",
      "key.converter.schema.registry.url":"http://kafka-schema-registry:8081",
      "value.converter":"io.confluent.connect.avro.AvroConverter",
      "value.converter.schema.registry.url":"http://kafka-schema-registry:8081",
      "value.subject.name.strategy":"io.confluent.kafka.serializers.subject.TopicRecordNameStrategy",
      "file": "/tmp/sink-file.txt"
    }
}

我从连接器收到一个错误,这似乎是某种基于 AvroConverter 的序列化错误,例如 这里显示的那个

org.apache.kafka.connect.errors.DataException: source-topic
    at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:95)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:468)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:301)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:205)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:173)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:170)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:214)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)
Caused by: org.apache.kafka.common.errors.SerializationException: Error retrieving Avro schema for id 2
Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Subject not found.; error code: 40401
    at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:202)
    at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:229)
    at io.confluent.kafka.schemaregistry.client.rest.RestService.lookUpSubjectVersion(RestService.java:296)
    at io.confluent.kafka.schemaregistry.client.rest.RestService.lookUpSubjectVersion(RestService.java:284)
    at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getVersionFromRegistry(CachedSchemaRegistryClient.java:125)
    at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getVersion(CachedSchemaRegistryClient.java:236)
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:152)
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserializeWithSchemaAndVersion(AbstractKafkaAvroDeserializer.java:194)
    at io.confluent.connect.avro.AvroConverter$Deserializer.deserialize(AvroConverter.java:120)
    at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:83)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:468)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:301)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:205)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:173)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:170)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:214)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)

请注意,架构注册表有一个 ID 为 2 的 Avro 架构和另一个架构 ID 为 3 的 Avro 架构,它们描述了托管在 同一个话题。使用 JDBC 连接器时也会出现同样的问题。

那么我该如何处理这种情况,以便从我的 Kafka 集群中导出数据 到外部系统。我的配置是否遗漏了什么?是否可以有一个包含多种事件类型的主题 并通过 Kafka Connect 导出?

【问题讨论】:

    标签: apache-kafka apache-kafka-connect confluent-schema-registry


    【解决方案1】:

    找到了解决办法。我的代码将键作为字符串传递,将值作为 avro 传递。阅读时 Hive-sink 尝试查找密钥的 avro 模式但无法找到它。 添加属性 key.converter=org.apache.kafka.connect.storage.StringConverter key.converter.schema.registry.url=http://localhost:8081 帮助解决了这个问题。

    【讨论】:

    • 出于兴趣,connect 创建的 Avro 文件的架构是什么? union[typeA, typeB]?或者也许:record wrapper { { null, typeA} a, { null, typeB } b }?
    猜你喜欢
    • 2023-03-07
    • 2021-06-08
    • 1970-01-01
    • 1970-01-01
    • 2021-07-27
    • 2019-05-27
    • 2022-06-11
    • 2020-01-02
    • 2019-08-13
    相关资源
    最近更新 更多