Debezium SQL Server源连接器临时快照不符合预期排查
Debezium SQL Server临时增量快照条件不生效问题排查
问题
执行带WHERE条件last_name='Walker'的临时增量快照时,连接器仍快照dbo.customers全表数据,未过滤出目标行。
用户操作步骤
1. 数据库初始化
CREATE DATABASE testDB; GO USE testDB; EXEC sys.sp_cdc_enable_db; CREATE TABLE customers ( id INTEGER IDENTITY(1001,1) NOT NULL PRIMARY KEY, first_name VARCHAR(255) NOT NULL, last_name VARCHAR(255) NOT NULL, email VARCHAR(255) NOT NULL UNIQUE ); INSERT INTO customers(first_name,last_name,email) VALUES ('Sally','Thomas','sally.thomas@acme.com'); INSERT INTO customers(first_name,last_name,email) VALUES ('George','Bailey','gbailey@foobar.com'); INSERT INTO customers(first_name,last_name,email) VALUES ('Edward','Walker','ed@walker.com'); INSERT INTO customers(first_name,last_name,email) VALUES ('Anne','Kretchmar','annek@noanswer.org'); EXEC sys.sp_cdc_enable_table @source_schema = 'dbo', @source_name = 'customers', @role_name = NULL, @supports_net_changes = 0;
2. 信号表创建与指令插入
CREATE TABLE debezium_signal (id VARCHAR(42) PRIMARY KEY, type VARCHAR(32) NOT NULL, data VARCHAR(2048) NULL); INSERT INTO dbo.debezium_signal (id, type, data) VALUES ('ad-hoc-1','execute-snapshot','{"data-collections": ["dbo.customers"],"type":"incremental","additional-conditions":"last_name=Walker"}');
3. 连接器配置
{ "name": "customer-adhoc", "config": { "connector.class" : "io.debezium.connector.sqlserver.SqlServerConnector", "tasks.max" : "1", "topic.prefix" : "CDC", "database.hostname" : "sqlserver12", "database.port" : "1433", "database.user" : "sa", "database.password" : "Password!", "database.names" : "testDB", "snapshot.mode": "initial", "schema.history.internal.kafka.bootstrap.servers" : "kafka12:9092", "schema.history.internal.kafka.topic": "schema-changes.inventory", "include.schema.changes": "true", "database.encrypt": "false", "table.include.list": "dbo.customers,dbo.debezium_signal", "column.mask.with.0.chars": "testDB.dbo.customers.first_name, testDB.dbo.customers.last_name", "schema.history.internal.store.only.captured.tables.ddl": "true", "schema.history.internal.store.only.captured.databases.ddl": "true", "incremental.snapshot.allow.schema.changes" : "true" , "key.converter.apicurio.registry.auto-register": "true", "key.converter.apicurio.registry.find-latest": "true", "value.converter.apicurio.registry.auto-register": "true", "value.converter.apicurio.registry.find-latest": "true", "schema.name.adjustment.mode": "avro", "value.converter": "io.apicurio.registry.utils.converter.AvroConverter", "key.converter": "io.apicurio.registry.utils.converter.AvroConverter", "value.converter.apicurio.registry.global-id": "io.apicurio.registry.utils.serde.strategy.AutoRegisterIdStrategy", "key.converter.apicurio.registry.global-id": "io.apicurio.registry.utils.serde.strategy.AutoRegisterIdStrategy", "key.converter.apicurio.registry.url": "http://****:8080/apis/registry/v2", "value.converter.apicurio.registry.url": "http://****:8080/apis/registry/v2", "signal.data.collection": "testDB.dbo.debezium_signal", "signal.kafka.topic":"CDC.dbz-signal", "kafka.consumer.offset.commit.enabled": "true", "signal.kafka.groupId": "customer-kafka-signal", "signal.kafka.bootstrap.servers": "kafka12:9092" } }
修复方案
1. 调整快照模式配置
当前snapshot.mode=initial会触发启动时全量快照,直接覆盖临时快照指令。
- 修改为
snapshot.mode=schema_only,避免启动时自动执行全量快照;或者先完成初始快照后,再发送临时快照信号。
2. 修正信号条件的SQL格式
信号数据中的additional-conditions缺少字符串常量的单引号,导致条件逻辑不生效:
- 重新插入信号指令(SQL Server中需用双单引号转义字符串内的单引号):
INSERT INTO dbo.debezium_signal (id, type, data) VALUES ('ad-hoc-2','execute-snapshot','{"data-collections": ["dbo.customers"],"type":"incremental","additional-conditions":"last_name=''Walker''"}');
3. 调整信号写入时机
先启动连接器,等待连接器完成初始化(日志显示已连接数据库、开始监听CDC事件)后,再插入信号表数据,确保连接器能正确捕获快照指令。
4. 验证增量快照配置
确认连接器配置中incremental.snapshot.allow.schema.changes=true已正确设置,且信号表已加入table.include.list(当前配置已满足)。
内容的提问来源于stack exchange,提问作者Dyan
相关产品推荐
相关产品推荐

