【问题标题】:How to set the key of the JDBC source connector (kafka)?如何设置JDBC源连接器(kafka)的key?
【发布时间】:2020-09-02 07:36:40
【问题描述】:

我正在使用Kafka Source JDBC connector 从 数据库表中读取数据并将其发布到主题test-mysql-petai。

数据库表有2个字段,Id是主键:

+---------+-------------+------+-----+---------+----------------+
| Field   | Type        | Null | Key | Default | Extra          |
+---------+-------------+------+-----+---------+----------------+
| id      | int(11)     | NO   | PRI | NULL    | auto_increment |
| name    | varchar(20) | YES  |     | NULL    |                |
+---------+-------------+------+-----+---------+----------------+

我需要id 字段的值作为主题的键。我尝试向 jdbc 连接器属性添加转换。

JDBCConnector.properties:

name=jdbc-source-connector    
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
tasks.max=1   
connection.url=jdbc:mysql://127.0.0.1:3306/test?user=dins&password=pw&serverTimezone=UTC
table.whitelist=petai 
mode=incrementing
incrementing.column.name=id    
schema.pattern=""    
transforms=createKey,extractInt    
transforms.createKey.type=org.apache.kafka.connect.transforms.ValueToKey    
transforms.createKey.fields=id    
transforms.extractInt.type=org.apache.kafka.connect.transforms.ExtractField$Key    
transforms.extractInt.field=id    
topic.prefix=test-mysql-jdbc-

但是,当我使用消费者读取键和值时,我得到以下信息:

Key = {"schema":{"type":"int32","optional":false},"payload":61} 
Value ={"id":61,"name":"ttt"}

我需要得到以下内容:

Key = 61    
Value ={"id":61,"name":"ttt"}

我做错了什么?任何帮助表示赞赏。

谢谢。

【问题讨论】:

  • 你如何消费消息?并请提供以下属性:key.converter、value.converter、key.converter.schemas.enable、value.converter.schemas.enable
  • @IskuskovAlexander 为了测试这一点,我使用了 kafka 附带的默认 kafka-console-consumer.sh 脚本。 bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-mysql-jdbc-petai --from-beginning 运行脚本时没有传递任何附加参数。我需要添加你提到的那些吗?
  • 尝试设置"key.converter.schemas.enable": "false"
  • 我试过了,还是一样。
  • 它正在工作。 :) 除了你提到的它还需要你在第一次提到的其他属性。 key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false 另外,正如您所说,它只是更改新消息的键,即使我再次运行消费者。我可以知道为什么吗?

标签: mysql apache-kafka-connect


【解决方案1】:

如果您不想在键中包含架构,可以通过设置key.converter.schemas.enable=false 告诉 Kafka Connect。

详细解释请见Kafka Connect Deep Dive – Converters and Serialization Explained by Robin Moffatt。

【讨论】:

    猜你喜欢
    • 2017-10-12
    • 2020-01-15
    • 2019-11-17
    • 2019-12-24
    • 2021-05-07
    • 2018-05-01
    • 2019-11-13
    • 2019-11-15
    • 2018-11-30
    相关资源
    最近更新 更多