分布式模式下将Kafka偏移量存储到S3而非主题的方案问询
如何将Debezium Mongo连接器的偏移量存储到S3等外部存储
Debezium Mongo Source Connector捕获MongoDB变更流时,偏移量(核心为resumeToken)的存储逻辑固定:独立模式存在/tmp/connect.offsets文件,分布式模式存在connect-offsetsKafka主题,且Kafka Connect本身不支持通过配置替换默认的偏移量存储实现(默认分布式用KafkaOffsetBackingStore)。无需修改代码重打包,有以下两种可行方案:
先说明偏移量的存储结构(二进制转译后):
Key: ["debezium-mongo-source-in",{"rs":"Cluster0-shard-0","server_id":"xdebezium"}] Value: {"sec":1683566267,"ord":19,"transaction_id":null,"resume_token":"8264592EBB000000132B022C0100296E5A10040586028B20A1484161EAA0F46645F696400646448B6E3D382E3B80004"}
方案1:用S3 Sink Connector实时备份偏移量到S3
直接部署一个S3 Sink Connector,专门监听connect-offsets主题,将偏移量数据同步到S3:
- 配置S3 Sink时,指定监听主题为
connect-offsets,使用org.apache.kafka.connect.json.JsonConverter将二进制偏移量转为可读的JSON格式存储,同时配置S3的访问凭证、存储桶路径等参数。 - 启动该Sink后,所有偏移量的变更都会实时同步到S3。
- 若需恢复偏移量,可将S3中的数据通过
kafka-console-producer等工具导回connect-offsets主题。
这种方案完全基于现有组件配置实现,无需修改任何代码,且能保证偏移量的实时备份,是最易落地的方案。
方案2:通过MirrorMaker2将偏移量主题同步到S3后端存储
如果环境中部署了Kafka MirrorMaker2,可配置它将connect-offsets主题同步到底层绑定S3的存储系统(部分托管Kafka服务支持将主题数据持久化到S3)。这种方式相当于将偏移量主题的底层存储替换为S3,无需额外部署Sink Connector,但依赖MirrorMaker2的配置和对应存储服务的支持。
以上方案均无需修改Kafka Connect或Debezium的核心代码,即可实现将偏移量数据存储到S3的需求。
内容的提问来源于stack exchange,提问作者Ritik
相关产品推荐
相关产品推荐

