【问题标题】:How to add multiple topics in JDBC Sink Connector configuration and get topics data in multiple target tables?如何在 JDBC Sink Connector 配置中添加多个主题并获取多个目标表中的主题数据?
【发布时间】:2022-01-21 05:13:14
【问题描述】:

以下是我的 JDBC-Sink 连接器配置:

connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
behavior.on.null.values=ignore
table.name.format=kafka_Address_V1, kafka_Attribute_V1
connection.password=***********
topics=Address,Attribute
task.max=3
batch.size=500
value.converter.value.subject.name.strategy=io.confluent.kafka.serializers.subject.RecordNameStrategy
value.converter.schema.registry.url=http://localhost:8081
auto.evolve=true
connection.user=user
name=sink-jdbc-connector
errors.tolerance=all
auto.create=true
value.converter=io.confluent.connect.avro.AvroConverter
connection.url=jdbc:sqlserver://localhost:DB;
insert.mode=upsert
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://localhost:8081
pk.mode=record_value
pk.fields=id

如果我使用此配置,我将在目标数据库中以这种 kafka_Address_V1、kafka_Attribute_V1 格式获取单个表,这是这两者的组合。

请告诉我如何使用 JDBC-Sink 连接器将不同的主题数据存储在不同的表中。

【问题讨论】:

  • topics 单独应该创建两个表。我不认为表格可以有逗号。为什么需要table.name.format?
  • @OneCricketeer 我已经通过删除 table.name.format 尝试过这种方式,但表没有在目标数据库中创建。如果我在源主题接收器连接器中插入或更新任何记录,则读取该记录但它没有将该记录发送到目标数据库,并且在这种情况下,表不是在目标数据库中创建的。请让我知道如何进一步进行,以便在这种情况下实现 JDBC 接收器连接器的预期行为。
  • 当事情没有发送时,您是否从连接中收到错误?我看到你有auto.create=true,所以表应该根据主题名称本身自动创建,所以你真的不需要table.name.format,除了匹配预期的表名。
  • 是的,我出错了,我的主题名称就像 iq.db.topicName 这就是为什么它要寻找 iq db 这是源 DB,我使用了转换,dropPrefix 来获得主题名称。而且我使用了 kafka_${topic},它在目标数据库中为我在主题属性中定义的所有主题创建表。现在它按预期工作。

标签: jdbc apache-kafka apache-kafka-connect mssql-jdbc


【解决方案1】:

根据the docs,table.name.format 采用单个值,并且默认使用主题名称本身。

编辑:由@OneCricketeer 提供,您也可以只使用table.name.format=kafka_${topic}_V1。下面的 SMT 对于更复杂的名称转换很有用。

要实现您想要的,您可以使用RegExRouter Single Message Transform 修改主题,因为它由 Kafka Connect 处理

试试这个:

transforms                             =changeTopicName
transforms.changeTopicName.type        =org.apache.kafka.connect.transforms.RegexRouter
transforms.changeTopicName.regex       =(.*)
transforms.changeTopicName.replacement =kafka_$1_V1

【讨论】:

  • ?‍♂️ 是的。很好的收获。
猜你喜欢
  • 2020-05-02
  • 2019-06-23
  • 2020-07-27
  • 2022-07-05
  • 2021-11-16
  • 2021-11-20
  • 2019-11-24
  • 2013-06-15
  • 1970-01-01
相关资源
最近更新 更多