You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka JDBC Sink Connector处理Delete操作失败问题求助

Debezium+Kafka JDBC Sink 删除同步失败问题排查与解决

问题现象

通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 02:06:05