Debezium MSSQL添加列后增量快照失败,如何恢复连接器?
解决Debezium MSSQL连接器新增列后增量快照失败问题
问题现象
给源表新增Column3后执行增量快照,连接器触发如下错误:
[2023-03-27 01:20:06,176] ERROR [source|task-0] Producer failure (io.debezium.pipeline.ErrorHandler:35) org.apache.kafka.connect.errors.ConnectException: Error while processing event at offset {transaction_id=null, event_serial_no=1, incremental_snapshot_maximum_key=0123456, commit_lsn=0123456:0123456:0003, change_lsn=0123456:0123456:0002, incremental_snapshot_collections=SOURCE.dbo.TaxIdentification, incremental_snapshot_primary_key=0123456} at io.debezium.pipeline.EventDispatcher.dispatchDataChangeEvent(EventDispatcher.java:246) at io.debezium.connector.sqlserver.SqlServerStreamingChangeEventSource.lambda$executeIteration$1(SqlServerStreamingChangeEventSource.java:290) at io.debezium.jdbc.JdbcConnection.prepareQuery(JdbcConnection.java:606) at io.debezium.connector.sqlserver.SqlServerConnection.getChangesForTables(SqlServerConnection.java:329) ... Caused by: java.lang.IllegalArgumentException: Column 'Column3' not found in result set 'Column1, Column2, Column3' for table 'Database.dbo.Table', columns: { Column1 int(10, 0) NOT NULL Column2 int(10, 0) NOT NULL } primary key: [Column1] default charset: null comment: null . This might be caused by DBZ-4350 at io.debezium.util.ColumnUtils.toArray(ColumnUtils.java:57) at io.debezium.pipeline.source.snapshot.incremental.AbstractIncrementalSnapshotChangeEventSource.lambda$readChunk$2(AbstractIncrementalSnapshotChangeEventSource.java:299)
已尝试以下方法但无效:
- 为该表创建新的CDC捕获实例(保留旧实例)
- 清理
debezium_signal表后暂停并重启连接器
解决方案
1. 清除遗留的增量快照offset状态
连接器的offset中留存了未完成的快照任务状态,这是重启后仍报错的核心原因:
- 暂停目标连接器
- 调用Kafka Connect API获取当前offset:
curl -X GET http://<connect-host>:<port>/connectors/<connector-name>/offsets - 在返回的JSON数据中,找到包含
incremental_snapshot_collections、incremental_snapshot_primary_key等快照相关字段的条目,将其完全删除 - 提交修改后的offset:
(分布式模式下需确保所有Connect节点的offset同步)curl -X PUT -H "Content-Type: application/json" \ http://<connect-host>:<port>/connectors/<connector-name>/offsets \ -d '{"<your-remaining-offset-key>": "<updated-offset-value>"}'
2. 强制刷新表元数据
Debezium缓存了旧的表结构,需要触发元数据刷新:
- 在数据库的
debezium_signal表中插入刷新信号:INSERT INTO debezium_signal (id, type, data) VALUES ('schema-refresh-' + NEWID(), 'execute-snapshot', '{"data-collections": ["Database.dbo.Table"], "type": "schema-refresh"}'); - 重启连接器,使其加载最新的表结构
3. 重新触发增量快照
待连接器启动完成且元数据刷新后,插入增量快照触发信号:
INSERT INTO debezium_signal (id, type, data) VALUES ('inc-snapshot-' + NEWID(), 'execute-snapshot', '{"data-collections": ["Database.dbo.Table"], "type": "incremental"}');
4. 验证结果
启动连接器后,检查日志是否不再出现列不存在的错误,同时确认Kafka主题中生成的快照数据包含新增的Column3。
后续预防建议
- 执行表结构变更前,先暂停所有正在进行的增量快照任务
- 升级Debezium至最新稳定版本,DBZ-4350已在后续版本中修复
内容的提问来源于stack exchange,提问作者jnnnnn
相关产品推荐
相关产品推荐

