如何为Pulsar中的Debezium PostgreSQL连接器配置消息过滤?
Pulsar Debezium CDC连接器Filter SMT正确配置方法
核心问题:配置前缀缺失
Pulsar的Debezium CDC连接器要求所有Debezium原生配置(包括单消息转换SMT)必须添加debezium.前缀,这是导致你配置不生效的主要原因。
修正后的Filter SMT配置
把你原来的配置加上前缀,同时注意YAML文件中特殊字符的处理,正确配置如下:
# 其他Pulsar连接器基础配置... debezium.transforms=filter debezium.transforms.filter.type=io.debezium.transforms.Filter debezium.transforms.filter.language=jsr223.groovy # 方式1:用单引号包裹条件避免转义 debezium.transforms.filter.condition='value.op == ''u'' && value.before.id == 2'
或者用YAML多行字符串语法,更直观且无需转义:
debezium.transforms.filter.condition: | value.op == 'u' && value.before.id == 2
配置验证要点
- 查看连接器启动日志,确认
Filter转换是否加载成功,若有报错(比如字段不存在、Groovy语法错误),根据日志调整条件 - 确认消息结构符合预期:只有更新操作(
op='u')的消息才会包含value.before字段,需确保你的业务场景中存在此类消息 - 测试时可先简化条件(比如仅保留
value.op == 'u'),验证生效后再添加复杂判断
内容的提问来源于stack exchange,提问作者Han Cygnus
相关产品推荐
相关产品推荐

