如何着手修复Debezium Engine的偏移量文件损坏问题
问题背景
使用Debezium Engine同步MySQL数据,采用org.apache.kafka.connect.storage.FileOffsetBackingStore记录偏移量,因意外断电导致偏移量文件损坏,启动时触发如下报错:
ERROR io.debezium.embedded.EmbeddedEngine - Unable to configure and start the 'org.apache.kafka.connect.storage.FileOffsetBackingStore' offset backing store
org.apache.kafka.connect.errors.ConnectException: java.io.StreamCorruptedException: invalid stream header: 00000000
at org.apache.kafka.connect.storage.FileOffsetBackingStore.load(FileOffsetBackingStore.java:86)
at org.apache.kafka.connect.storage.FileOffsetBackingStore.start(FileOffsetBackingStore.java:59)
at io.debezium.embedded.EmbeddedEngine.run(EmbeddedEngine.java:691)
at io.debezium.embedded.ConvertingEngineBuilder$2.run(ConvertingEngineBuilder.java:192)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: java.io.StreamCorruptedException: invalid stream header: 00000000
at java.base/java.io.ObjectInputStream.readStreamHeader(ObjectInputStream.java:987)
at java.base/java.io.ObjectInputStream.<init>(ObjectInputStream.java:414)
at org.apache.kafka.connect.util.SafeObjectInputStream.<init>(SafeObjectInputStream.java:48)
at org.apache.kafka.connect.storage.FileOffsetBackingStore.load(FileOffsetBackingStore.java:71)
... 8 common frames omitted
需要修复偏移量文件,确定合适的目标偏移量(故障发生时或更早的临近值),并完成文件修改。
解决方案步骤
1. 尝试从损坏文件中提取有效偏移量
偏移量文件是Java序列化的HashMap,头部损坏不代表全部内容失效:
- 用十六进制编辑器打开损坏文件,查找序列化数据的特征起始标识
AC ED(正常Java序列化流的开头) - 如果找到后续的有效片段,将该片段复制到新文件中,尝试用Debezium加载,或通过Java反序列化工具读取其中的偏移量键值对
2. 通过MySQL Binlog定位时间对应的偏移量
若文件无法修复,直接从MySQL端获取故障前的有效偏移量:
- 执行MySQL命令查看所有binlog文件:
SHOW BINARY LOGS; - 找到故障时间点前后的目标binlog文件,用
mysqlbinlog解析并筛选时间范围内的事件:mysqlbinlog --start-datetime="YYYY-MM-DD HH:MM:SS" --stop-datetime="YYYY-MM-DD HH:MM:SS" mysql-bin.0000XX > binlog_events.txt - 从解析结果中找到故障发生前最后一个事务的
position值,同时记录对应的binlog文件名,这两个值就是Debezium偏移量中需要的file和pos字段
3. 从Kafka主题历史数据中提取偏移量
如果Debezium同步的Kafka主题保留了历史消息:
- 消费目标主题的历史消息,查看每条消息的
source字段,其中包含file(binlog文件名)、pos(binlog位置)和ts_ms(时间戳) - 按
ts_ms排序,找到故障发生前的最后一条有效消息对应的偏移量,作为目标值
4. 修改偏移量文件并验证
确定目标偏移量后,用HashmapEditor工具完成修改:
- 生成一个正常的空白偏移量文件(可临时启动Debezium同步测试库生成,或用工具创建)
- 打开文件,替换MySQL连接器对应的偏移量字段:通常包含
server_id、file、pos、gtid(如果启用GTID)等,填入你找到的有效值 - 将修改后的文件替换原损坏的偏移量文件,启动Debezium Engine验证同步是否正常
内容的提问来源于stack exchange,提问作者omatase

