Debezium Source与JDBC Sink Connector删除操作失效求助
问题
配置Debezium MySQL Source Connector与Confluent JDBC Sink Connector后,能正常同步MySQL到Kafka的增改数据,且Kafka可将数据写入另一MySQL,但删除源表数据时,Kafka主题未生成对应墓碑消息,目标表也未删除对应记录。
Source Connector配置
{ "name": "smartdevsignupconnector111", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "true", "value.converter.schemas.enable": "true", "database.hostname": "mysql1", "database.port": "3306", "database.user": "clusterAdmin", "database.password": "RUNSman001", "database.server.id": "184055", "database.server.name": "smartdevdbserver1", "database.include.list": "signup_db", "schema.history.internal.kafka.topic": "schema-changes.signup_db", "schema.history.internal.kafka.bootstrap.servers": "kafka1:9092", "table.include.list": "signup_db.users", "column.exclude.list": "signup_db.users.fullName, signup_db.users.address, signup_db.users.phoneNo, signup_db.users.gender, signup_db.users.userRole, signup_db.users.reason_for_inactive, signup_db.users.firstvisit, signup_db.users.last_changed_PW, signup_db.users.regDate", "snapshot.mode": "when_needed", "topic.creation.enable": "true", "topic.prefix": "smartdevdbserver1", "topic.creation.default.replication.factor": "1", "topic.creation.default.partitions": "1", "transforms": "unwrap,dropTopicPrefix", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.dropTopicPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.dropTopicPrefix.regex": "smartdevdbserver1.signup_db.(.*)", "transforms.dropTopicPrefix.replacement": "$1", "include.schema.changes": "true" } }
Sink Connector配置
{ "name": "resetpassword-sink-connector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "true", "value.converter.schemas.enable": "true", "topics": "users", "connection.url": "jdbc:mysql://rpwd_mysql:3306/rpwd_db", "connection.user": "rpwd_user", "connection.password": "*RUNSman001*", "table.name.format": "users", "fields.whitelist": "id,email,password,User_status,auth_token", "transforms": "unwrap", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "false", "auto.create": "true", "insert.mode": "upsert", "delete.enabled": "true", "pk.fields": "id", "pk.mode": "record_key" } }
执行源表删除操作:
mysql> delete from users where email = 'testing14@firstclicklimited.com'; Query OK, 1 row affected (0.11 sec)
使用--property print.key=true消费Kafka主题,未发现对应删除操作的墓碑消息,主题中原记录仍存在,目标表也未删除对应数据。
排查与解决要点
Source Connector需保留墓碑消息:当前Source的
ExtractNewRecordState转换未配置drop.tombstones=false,默认会丢弃墓碑消息。需在Source的transforms.unwrap配置中添加:"transforms.unwrap.drop.tombstones": "false"确保删除操作生成的墓碑消息能被发送到Kafka。
确认源表主键存在:Debezium生成墓碑消息依赖表的主键,检查源表
signup_db.users是否将id设为主键(Sink配置指定pk.fields=id,源表需保持一致),无主键则无法生成有效墓碑消息。检查MySQL binlog配置:确保源MySQL的binlog格式为
ROW,且binlog_row_image设置为FULL。Debezium依赖row格式binlog捕获删除操作,格式错误或row image不全会导致无法捕获删除事件。移除Sink Connector的重复Unwrap转换:Source已完成
ExtractNewRecordState转换,Sink重复配置该转换可能破坏墓碑消息结构。直接删除Sink中的transforms相关配置即可。验证Kafka消费方式:消费时需添加
--from-beginning参数,避免错过最新的删除事件;墓碑消息的特征是key存在,value为null,通过print.key=true可直观判断是否生成了墓碑消息。
内容的提问来源于stack exchange,提问作者eedideyahoocom

