使用Debezium SQL Server连接器如何在Debezium层实现数据过滤?
Debezium SQL Server Connector 层数据过滤实现方案
适用场景
你当前使用Debezium SQL Server connector做数据同步,需要在connector层完成数据筛选,仅将符合规则的数据推送至Kafka topic,可以通过以下两种官方原生方案实现,无需额外组件介入。
方案一:粗粒度表/列过滤(无需SMT)
如果只需要按表、列维度做过滤,直接使用connector内置配置即可:
- 表包含配置:
table.include.list,格式为<数据库名>.<模式名>.<表名>,多表用逗号分隔,仅同步列表内的表数据 - 表排除配置:
table.exclude.list,排除指定表的同步,优先级高于包含配置 - 列包含配置:
column.include.list,格式为<数据库名>.<模式名>.<表名>.<列名>,仅同步列表内的列数据 - 列排除配置:
column.exclude.list,排除指定列的同步,优先级高于包含配置
方案二:细粒度行级过滤(使用Filter SMT)
如果需要按行数据的内容做条件过滤,使用Debezium自带的Filter单消息转换组件,通过SpEL表达式定义过滤规则:
- 核心配置项:
- 指定SMT类型为
io.debezium.transforms.Filter - 用
transforms.filter.condition配置过滤表达式,支持读取CDC消息的元数据、行数据内容定义规则 - 用
transforms.filter.null.behavior指定不符合条件的消息处理逻辑,设置为drop即可直接丢弃
- 指定SMT类型为
- 完整配置示例:
{ "connector.class": "io.debezium.connector.sqlserver.SqlServerConnector", "tasks.max": "1", "database.hostname": "你的SQL Server地址", "database.port": "1433", "database.user": "账号", "database.password": "密码", "database.dbname": "数据库名", "database.server.name": "自定义逻辑服务名", "transforms": "row-filter", "transforms.row-filter.type": "io.debezium.transforms.Filter", "transforms.row-filter.language": "spel", // 示例规则:仅同步users表中status=1的新增/更新数据 "transforms.row-filter.condition": "value.payload.after != null && value.payload.source.table == 'users' && value.payload.after.status == 1", "transforms.row-filter.null.behavior": "drop" }
注意:过滤操作在数据解析完成后、写入Kafka前执行,不会影响CDC日志的正常拉取,所有不符合条件的消息会直接在connector层丢弃,不会进入Kafka topic。
内容的提问来源于stack exchange,提问作者Ashish Raj
相关产品推荐
相关产品推荐

