如何配置Kafka Connect JDBC Sink实现MySQL表的插入与更新?
解决方案:通过分离Kafka主题与连接器区分MySQL表的插入和更新操作
是的,完全可以通过配置独立的Kafka主题和对应JDBC Sink连接器来区分同一MySQL表上的插入与更新操作,避免两者的冲突。具体实现步骤如下:
1. 创建两个独立的Kafka主题
为插入和更新操作分别创建专属主题:
table-inserts:用于接收不需要指定主键的插入事件table-updates:用于接收包含主键id的更新事件
调整消息生产者逻辑,将插入消息发送到table-inserts,更新消息发送到table-updates即可。
2. 配置插入专用连接器
修改原插入连接器的topics配置,指向table-inserts主题,其余配置保持不变(插入操作无需指定自增主键):
{ "name": "sink-connector-insert", "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "connection.url": "jdbc:mysql://server:3306/database?serverTimezone=UTC&useLegacyDatetimeCode=false", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "https://registry:8081", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "https://registry:8081", "topics": "table-inserts", // 修改为插入专属主题 "connection.user": "user", "connection.password": "XXXXXXXXXX", "ssl.mode": "prefer", "insert.mode": "insert" }
3. 配置更新专用连接器
修改原更新连接器的topics配置,指向table-updates主题,其余配置保持不变(更新操作需要主键匹配):
{ "name": "sink-connector-update", "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "connection.url": "jdbc:mysql://server:3306/database?serverTimezone=UTC&useLegacyDatetimeCode=false", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "https://registry:8081", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "https://registry:8081", "topics": "table-updates", // 修改为更新专属主题 "connection.user": "user", "connection.password": "XXXXXXXXXX", "ssl.mode": "prefer", "insert.mode": "update", "pk.mode": "record_value", "pk.fields": "id" }
核心原理
- 插入连接器仅处理
table-inserts主题的消息,这些消息不含主键,不会触发主键冲突,同时利用MySQL自增主键生成ID - 更新连接器仅处理
table-updates主题的消息,这些消息包含主键id,可以正常匹配到目标记录执行更新,不会出现缺少主键的错误
这样就彻底隔离了插入和更新操作的数据流,解决了原配置中两个连接器监听同一主题导致的冲突问题。
内容的提问来源于stack exchange,提问作者andand
相关产品推荐
相关产品推荐

