Debezium SQL Server源连接器非实时、每日单次推送CDC事件问题排查
Debezium SQL Server Source 连接器固定时间批量推送CDC事件排查
异常现象
- 连接器无法实时流式传输CDC捕获的变更数据,每日固定在
19:07:29一次性推送所有累积事件 - 事件时间差特征明显:
事件样例如下:"source":{ "version":"1.8.1.Final", "connector":"sqlserver", "name":"dataroom-rds", "ts_ms":1655275673757, "snapshot":"false", "db":"dbDataRoom", "sequence":null, "schema":"dbo", "table":"tblFolder", "change_lsn":"0056e213:00000e55:000a", "commit_lsn":"0056e213:00000ecd:0014", "event_serial_no":1 }, "op":"c", "ts_ms":1655284049349, "transaction":nullsource.ts_ms对应变更在数据库的提交时间:2022-06-15 16:47:53.757- 顶层
ts_ms对应Debezium处理事件的时间:2022-06-15 19:07:29.349 - 两者延迟超2小时,所有批量推送事件的处理时间均集中在19:07:29附近
当前连接器配置如下:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnector metadata: name: mssql-connector-dbdataroom-raw namespace: kafka labels: strimzi.io/cluster: kafka-connect-cluster spec: class: io.debezium.connector.sqlserver.SqlServerConnector tasksMax: 1 config: database.server.id: "1" database.server.name: "dataroom-rds" database.hostname: "${secrets:kafka/dataroom-rds-creds:host}" database.port: "1433" database.user: "${secrets:kafka/dataroom-rds-creds:username}" database.password: "${secrets:kafka/dataroom-rds-creds:password}" database.dbname: "${secrets:kafka/dataroom-rds-creds:dbname}" table.include.list: "dbo.tblFolder,dbo.tblDocument,dbo.tblAllowedSecurityGroupFolder,dbo.tblAllowedSecurityGroupFolderHistory" snapshot.isolation.mode: "read_committed" snapshot.lock.timeout.ms: "-1" poll.interval.ms: "2000" database.history.kafka.bootstrap.servers: KUSTOMIZE_REPLACEMENT database.history.kafka.topic: "schema-changes.dbdataroom" database.history.producer.security.protocol: "SASL_SSL" database.history.producer.sasl.mechanism: "SCRAM-SHA-512" database.history.producer.sasl.jaas.config: "org.apache.kafka.common.security.scram.ScramLoginModule required username=${secrets:kafka/kafka-connect-msk-secrets:msk_sasl_user} password=${secrets:kafka/kafka-connect-msk-secrets:msk_sasl_password} ;" database.history.consumer.security.protocol: "SASL_SSL" database.history.consumer.sasl.mechanism: "SCRAM-SHA-512" database.history.consumer.sasl.jaas.config: "org.apache.kafka.common.security.scram.ScramLoginModule required username=${secrets:kafka/kafka-connect-msk-secrets:msk_sasl_user} password=${secrets:kafka/kafka-connect-msk-secrets:msk_sasl_password} ;" transforms: "changeTopicCase" transforms.changeTopicCase.type: "com.github.jcustenborder.kafka.connect.transform.common.ChangeTopicCase" transforms.changeTopicCase.from: "UPPER_UNDERSCORE" transforms.changeTopicCase.to: "LOWER_UNDERSCORE"
根因定位
该异常与Debezium端配置无关,核心问题出在SQL Server侧CDC作业配置,判断依据:
- 现有配置
poll.interval.ms设为2000ms,即每2秒轮询一次CDC变更,不存在轮询间隔过长问题 - 事件集中在固定时间点推送,延迟无随机波动,排除网络抖动、Kafka集群性能不足这类不稳定因素
- SQL Server CDC能力依赖两个系统代理作业实现:
- 捕获作业:扫描事务日志,将变更写入CDC系统表,供Debezium等消费端读取
- 清理作业:定期删除过期CDC记录,避免系统表膨胀
该场景下的批量推送问题,是因为CDC捕获作业被配置为每日19:07:29仅执行一次,而非默认的持续运行模式。所有变更先堆积在事务日志中,直到捕获作业触发时才批量写入CDC表,Debezium此时才能拉取到累积的所有事件。
验证方式:在目标SQL Server实例执行如下查询,检查CDC作业调度配置:
USE msdb; GO SELECT j.name job_name, s.active_start_time, s.freq_type, s.freq_subday_type, s.freq_subday_interval FROM sysjobs j JOIN sysjobschedules js ON j.job_id = js.job_id JOIN sysschedules s ON js.schedule_id = s.schedule_id WHERE j.name LIKE 'cdc.%%capture'; GO如果返回结果中捕获作业的
active_start_time值为190729,且freq_subday_type为1(仅在指定时间执行一次),即可确认根因。
存在两个小概率诱因:
- 数据库配置了每日19:00左右执行的事务日志备份任务,日志截断后CDC作业才能读取到对应日志段的变更
- Debezium连接的是Always On可用性组的辅助副本,副本日志同步存在固定延迟,每日19:07左右才完成当日日志的Apply操作
解决方案
核心修复:调整CDC捕获作业调度(90%以上场景适用)
- 打开SSMS连接目标SQL Server实例,展开「SQL Server 代理」-「作业」目录
- 找到名称格式为
cdc.[数据库名]_capture的捕获作业(该场景下为cdc.dbDataRoom_capture),右键打开属性面板 - 切换到「计划」页签,删除原有每日19:07执行一次的调度,新增持续运行调度:
- 计划类型设为「重复执行」,频率为每天
- 重复频率设为按秒间隔,间隔值设为5-10秒(与SQL Server CDC默认配置一致)
- 计划生效时间设为00:00:00,无结束时间
- 手动启动一次捕获作业,观察Debezium事件延迟是否降到秒级,是否开始实时推送
也可通过T-SQL直接修改调度配置:
-- 新建持续运行的调度计划,每5秒执行一次 EXEC sp_add_schedule @schedule_name = N'CDC_Capture_Continuous', @freq_type = 4, @freq_interval = 1, @freq_subday_type = 2, @freq_subday_interval = 5, @active_start_time = 000000; GO -- 将新计划绑定到CDC捕获作业,作业名请替换为实例上的实际名称 EXEC sp_attach_schedule @job_name = N'cdc.dbDataRoom_capture', @schedule_name = N'CDC_Capture_Continuous'; GO
小概率场景处理
- 若为事务日志备份策略导致:将事务日志备份间隔缩短到15-30分钟,避免每日仅备份一次导致的日志段长时间不释放
- 若为Always On副本延迟导致:将Debezium连接地址切换为主副本,或调整副本同步模式为同步提交,检查副本间网络带宽是否存在瓶颈
配置优化建议
可在连接器config段补充如下参数,进一步降低传输延迟:
max.batch.size: "4096" offset.flush.interval.ms: "1000"
内容的提问来源于stack exchange,提问作者hassanbsalimi
相关产品推荐
相关产品推荐

