MongoDB变更流副本集恢复问题:Spring实现仅主节点正常时可用
我之前在做Spring整合MongoDB变更流的时候,也碰到过副本集故障后订阅断连、丢失事件的问题,折腾了一阵后整理出几个实用的恢复方案,分享给你:
MongoDB副本集故障后变更流恢复方案
1. 基于resumeToken的精准断点恢复
这是MongoDB变更流原生支持的核心恢复能力,也是最可靠的方案。你需要在消费每个变更事件时,持久化保存最新的resumeToken(推荐存在独立的Mongo集合、Redis这类持久化存储里,别存在内存里,避免服务重启丢失)。
具体实现步骤:
- 处理变更事件时,每次完成业务逻辑后,把
ChangeStreamDocument.getResumeToken()的值存到持久化存储 - 当检测到变更流连接中断(比如捕获
MongoException、SubscriptionException),重新订阅时用resumeAfter()传入保存的token:
// 伪代码示例 private void restartChangeStream() { // 从持久化存储获取最后保存的token BsonDocument lastResumeToken = resumeTokenRepository.getLatestToken(); MongoDatabase db = mongoClient.getDatabase("experiment"); MongoCollection<Document> coll = db.getCollection("your-target-collection"); try { coll.watch() .resumeAfter(lastResumeToken) .forEach(this::handleChangeEvent); } catch (Exception e) { logger.error("变更流订阅中断,将尝试重启", e); restartChangeStream(); // 或者结合重试机制 } } private void handleChangeEvent(ChangeStreamDocument<Document> event) { // 执行你的业务逻辑 // ... // 更新持久化的resumeToken resumeTokenRepository.saveToken(event.getResumeToken()); }
注意:如果故障时间过长,resumeToken对应的oplog被清理了,可以改用startAtOperationTime()指定一个时间点恢复,不过这种方式可能会出现少量重复消费,需要业务逻辑配合处理。
2. Spring环境下的自动重连与重试配置
结合Spring的重试机制,让变更流订阅在故障后自动重启,不用人工介入:
- 用
@Retryable注解包裹订阅方法,指定需要重试的异常类型,配置指数退避避免频繁重试 - 搭配
@Recover方法处理重试多次失败后的告警逻辑:
@Service public class ChangeEventService { private static final Logger logger = LoggerFactory.getLogger(ChangeEventService.class); private final MongoClient mongoClient; private final ResumeTokenRepository resumeTokenRepository; public ChangeEventService(MongoClient mongoClient, ResumeTokenRepository resumeTokenRepository) { this.mongoClient = mongoClient; this.resumeTokenRepository = resumeTokenRepository; } @PostConstruct @Retryable(value = {MongoConnectionException.class, MongoException.class}, maxAttempts = -1, // 无限重试(可根据业务调整次数) backoff = @Backoff(delay = 5000, multiplier = 2)) // 5秒后重试,每次间隔翻倍 public void subscribeToChangeStream() { MongoDatabase db = mongoClient.getDatabase("experiment"); MongoCollection<Document> coll = db.getCollection("your-target-collection"); coll.watch() .resumeAfter(resumeTokenRepository.getLatestToken()) .forEach(this::handleChangeEvent); } @Recover public void handleRetryFailure(Throwable e) { logger.error("变更流重试多次仍无法恢复,需人工介入", e); // 这里可以加告警逻辑:发邮件、钉钉通知等 } private void handleChangeEvent(ChangeStreamDocument<Document> event) { // 业务逻辑处理 resumeTokenRepository.saveToken(event.getResumeToken()); } }
同时要确保MongoClient配置了自动发现副本集节点的参数:
@Bean public MongoClient mongoClient() { return MongoClients.create(MongoClientSettings.builder() .applyConnectionString(new ConnectionString("mongodb://node1:27017,node2:27017,node3:27017/?replicaSet=rs0")) .applyToSocketSettings(builder -> builder.connectTimeout(Duration.ofSeconds(10))) .applyToConnectionPoolSettings(builder -> builder.maxWaitTime(Duration.ofSeconds(10))) .build()); }
3. 副本集主节点切换的适配处理
当副本集发生主节点切换时,MongoDB驱动会自动发现新主,但变更流可能短暂中断。可以做这些优化:
- 连接字符串必须包含所有副本集节点,不能只写主节点地址,让驱动能自动感知节点变化
- 注册
ClusterListener监听副本集节点变更事件,在主节点切换后可以主动触发变更流重启:
mongoClient.subscribe(new ClusterListener() { @Override public void clusterDescriptionChanged(ClusterDescriptionChangedEvent event) { ClusterDescription newDesc = event.getNewDescription(); if (newDesc.hasPrimary()) { logger.info("副本集主节点切换完成,新主节点:{}", newDesc.getPrimary()); // 可选:主动重启变更流订阅 // subscribeToChangeStream(); } } // 其他接口方法空实现即可 @Override public void clusterOpening(ClusterOpeningEvent event) {} @Override public void clusterClosed(ClusterClosedEvent event) {} });
4. 必须做的幂等性保障
不管用哪种恢复方式,都可能出现重复消费的情况(比如故障时已经处理的事件,恢复后又被消费一次)。所以一定要保证业务处理逻辑是幂等的:
- 用变更事件的
_id(每个变更事件都有唯一的_id)作为去重标识,处理前先检查是否已经处理过 - 业务操作本身支持重复执行:比如更新用
$set而非直接覆盖文档,插入操作使用upsert模式
内容的提问来源于stack exchange,提问作者Sydney
相关产品推荐
相关产品推荐

