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

如何通过Debezium仅捕获PostgreSQL的创建与更新操作?

解决方案:过滤Debezium PostgreSQL连接器的DELETE事件

要仅捕获CREATE(INSERT)和UPDATE操作的变更记录,你可以通过添加Kafka Connect Single Message Transform (SMT) 来过滤掉DELETE事件,具体配置步骤如下:

核心配置添加项

在你的连接器配置中新增以下SMT相关参数,使用Debezium内置的Filter转换来筛选符合要求的事件:

"transforms": "filterDelete",
"transforms.filterDelete.type": "io.debezium.transforms.Filter",
"transforms.filterDelete.language": "jsr223.groovy",
"transforms.filterDelete.condition": "operation in ['c', 'u']"

参数说明

  • transforms: 定义转换的名称(这里命名为filterDelete,可自定义)
  • transforms.filterDelete.type: 指定使用Debezium的Filter转换类
  • transforms.filterDelete.language: 指定使用Groovy作为条件脚本语言(Debezium默认支持)
  • transforms.filterDelete.condition: 过滤条件,其中:
    • operation是Debezium事件中的内置字段,对应操作类型标识:
      • c = CREATE(INSERT操作)
      • u = UPDATE(UPDATE操作)
      • d = DELETE(DELETE操作)
      • r = READ(快照读取,你的配置中snapshot.mode: schema_only不会产生此类事件)
        该条件会仅保留operation为c或u的事件,自动过滤DELETE事件。

完整连接器配置

将上述配置项整合到你的原有配置中,完整配置如下:

{
    "name": "transaction-connector",
    "config": {
        "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
        "database.hostname": "localhost",
        "database.port": "5432",
        "database.user": "postgres",
        "database.password": "123123",
        "database.dbname": "walletgate_rg1",
        "table.include.list": "transaction_ctx.transaction",
        "topic.prefix" : "data",
        "decimal.handling.mode": "double",
        "time.precision.mode": "adaptive_time_microseconds",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "key.converter.schemas.enable": "false",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter.schemas.enable": "false",
        "include.schema.changes": "false",
        "snapshot.mode": "schema_only",
        "plugin.name": "pgoutput",
        "transforms": "filterDelete",
        "transforms.filterDelete.type": "io.debezium.transforms.Filter",
        "transforms.filterDelete.language": "jsr223.groovy",
        "transforms.filterDelete.condition": "operation in ['c', 'u']"
    }
}

生效方式

应用上述配置后,需要重启连接器或重新创建连接器实例,新的过滤规则才会生效。

内容的提问来源于stack exchange,提问作者Murat Dinc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:13:22