You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 10:30:09