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

使用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即可直接丢弃
  • 完整配置示例:
{
  "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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 08:54:02