Debezium分布式模式重启报错:快照配置冲突问题
Debezium重启MSSQL连接器时的快照异常问题
问题场景
在分布式模式下使用Debezium将MSSQL数据库的CDC事件推送至Kafka主题,启动命令为:
%KAFKA_HOME%\bin\windows\connect-distributed.bat
首次部署连接器后可正常同步表数据到Kafka,但停止启动脚本并重启后,连接器无法从上次停止的偏移量恢复,抛出如下异常:
DebeziumException: 连接器之前在执行快照时停止,但当前连接器配置为从不允许快照。请重新配置连接器,使其初始或在需要时使用快照。
连接器配置如下:
{ "name": "mssql-dbz-connector", "config": { "connector.class": "io.debezium.connector.sqlserver.SqlServerConnector", "database.hostname": "***", "database.port": "1433", "database.user": "***", "database.password": "***", "database.names": "***", "database.encrypt": "false", "topic.prefix": "***", "schema.history.internal.kafka.bootstrap.servers": "***", "schema.history.internal.kafka.topic": "***", "snapshot.mode": "schema_only" } }
完整异常堆栈信息:
ERROR [mssql-dbz-connector|task-0] WorkerSourceTask{id=mssql-dbz-connector-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:212) io.debezium.DebeziumException: 连接器之前在执行快照时停止,但当前连接器配置为从不允许快照。请重新配置连接器,使其初始或在需要时使用快照。 at io.debezium.connector.common.BaseSourceTask.validateAndLoadSchemaHistory(BaseSourceTask.java:97) at io.debezium.connector.sqlserver.SqlServerConnectorTask.start(SqlServerConnectorTask.java:108) at io.debezium.connector.common.BaseSourceTask.start(BaseSourceTask.java:240) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.initializeAndStart(AbstractWorkerSourceTask.java:280) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:202) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:77) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:236) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) [2024-05-19 23:48:29,610] INFO [mssql-dbz-connector|task-0] Stopping down connector (io.debezium.connector.common.BaseSourceTask:398)
原因分析
当前配置的snapshot.mode为schema_only,该模式仅在连接器首次启动时获取数据库表结构快照。但首次启动时连接器在执行快照过程中被中断,导致Debezium的元数据(偏移量、快照状态)记录为“快照未完成”。重启时,连接器检查到之前的快照未完成,但当前配置不允许重新执行快照,因此抛出异常。
解决方案
方案1:临时调整快照模式完成未结束的快照
修改连接器配置,将snapshot.mode改为initial或when_needed,重启连接器让它完成剩余的快照流程:
{ // ... 其他配置不变 "snapshot.mode": "when_needed" }
待快照完成并正常运行后,可再改回snapshot.mode: schema_only(若不需要全量数据快照)。
方案2:清除连接器的偏移量和元数据记录
如果不需要保留之前的快照进度,可直接删除连接器的偏移量和元数据:
- 删除Kafka Connect存储偏移量的主题(默认是
connect-offsets,若自定义过需对应调整)中该连接器的记录 - 删除Debezium的schema history主题中对应连接器的记录
- 重新部署连接器
方案3:强制重置连接器
使用Kafka Connect的REST API重置连接器任务:
curl -X POST http://<connect-host>:<port>/connectors/mssql-dbz-connector/tasks/0/restart?reset=true
重置后连接器会重新初始化,需确保此时snapshot.mode配置允许执行快照(如when_needed)。
内容的提问来源于stack exchange,提问作者Roobal Jindal
相关产品推荐
相关产品推荐

