Kafka JDBC Sink Connector处理Delete操作失败问题求助
问题现象
通过Debezium Postgres源连接器同步Postgres数据到Kafka,再用Kafka JDBC Sink同步至另一台Postgres服务器,插入、更新操作正常,但源库执行删除后,JDBC Sink尝试插入该行数据导致失败。已尝试设置transforms.unwrap.delete.handling.mode为"rewrite"和"drop",均无法解决。
连接器配置
源连接器(Debezium Postgres)
{ "name": "ksqldb-connector-actions", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "plugin.name": "pgoutput", "database.hostname": "ipadress", "database.port": "5432", "database.user": "db", "database.password": "*********", "database.dbname": "config", "database.server.name": "postgres", "topic.prefix":"kcon", "table.include.list": "dbo.actions", "slot.name" : "slot_actions_connector", "transforms":"unwrap", "transforms.unwrap.type":"io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones":"false", "transforms.unwrap.delete.handling.mode":"rewrite", "transforms.unwrap.add.fields":"table,lsn" } }
Sink连接器(Kafka JDBC)
{ "name": "jdbc-sink", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "kcon.dbo.actions", "connection.url": "jdbc:postgresql://ipadress:5432/config", "connection.user": "wft", "connection.password": "*******", "insert.mode": "upsert", "delete.enabled": "true", "table.name.format":"dbo.actions_etl_kafka", "pk.mode":"record_key", "pk.fields": "action_id", "db.timezone":"Asia/Kolkata", "auto.create":"true", "auto.evolve":"true", "errors.tolerance": "all", "errors.log.enable": "true", "errors.log.include.messages": "true", "transforms": "flatten", "transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Key", "transforms.flatten.delimiter": "_", "input.data.format": "AVRO", "key.converter":"io.confluent.connect.avro.AvroConverter", "value.converter":"io.confluent.connect.avro.AvroConverter", "key.converter.schemas.enable":"true", "value.converter.schemas.enable": "true", "key.converter.schema.registry.url":"http://schema-registry-ksql:8081", "value.converter.schema.registry.url":"http://schema-registry-ksql:8081" } }
解决步骤
1. 确认rewrite模式下的消息结构
当delete.handling.mode设为rewrite时,Debezium会将删除操作转换为包含主键、其余字段为null且带有__deleted: true标识的消息。先验证Kafka主题中的消息是否符合这个结构:
kafka-avro-console-consumer --bootstrap-server <kafka地址>:9092 --topic kcon.dbo.actions --from-beginning --property print.key=true
2. 修正源连接器配置
确保__deleted字段被正确包含在消息中(Debezium 1.9+版本在rewrite模式下自动添加,也可显式配置):
"transforms.unwrap.add.fields":"table,lsn,__deleted"
同时保持transforms.unwrap.drop.tombstones=false,保证重写后的删除消息能发送到Kafka。
3. 调整Sink连接器配置
- 移除不必要的Key Flatten转换:如果源库主键
action_id在Kafka消息Key中已是扁平结构,当前的Flatten$Key转换可能破坏主键识别,尝试删除Sink的transforms相关配置。 - 确认删除逻辑触发条件:JDBC Sink的
delete.enabled=true需要配合正确的主键映射和__deleted字段识别。若消息中存在__deleted: true,Sink会自动执行删除;若未识别到,会尝试插入null字段导致失败。
4. 验证主键映射
确认源库表dbo.actions的主键确实是action_id,且Kafka消息的Key中包含该字段。Sink的pk.mode=record_key和pk.fields=action_id必须与消息Key结构匹配,否则无法定位要删除的行。
5. drop模式的正确使用(不推荐用于删除同步场景)
若选择delete.handling.mode=drop,需将transforms.unwrap.drop.tombstones设为true,此时Debezium会直接丢弃墓碑消息,但Sink无法主动执行删除,仅适用于不需要同步删除的场景。
内容的提问来源于stack exchange,提问作者Alphonse

