如何通过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
相关产品推荐
相关产品推荐

