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

多实例部署下Spring Data MongoDB ChangeStream重复写审计记录求解

解决方案

方案1:分布式锁主备模式(适合低负载场景)

该模式下同一时间只有1个实例运行变更流监听,实例故障时自动切换到其他存活实例,实现逻辑如下:

  • 引入Spring Integration MongoDB分布式锁组件,也可以自己基于MongoDB唯一索引实现轻量锁:
    // build.gradle新增依赖
    implementation 'org.springframework.integration:spring-integration-mongodb'
    
  • 改造监听器容器的启动逻辑,只有拿到锁的实例才启动监听:
    @Configuration
    @Slf4j
    public class MongoChangeStreamListenerConfig {
      // 锁过期时间根据业务容忍的故障切换时长设置,比如30秒
      private static final long LOCK_TTL = 30000L;
      private static final String LOCK_KEY = "change-stream-lock:my-collection";
    
      @Bean
      public MongoLockRegistry mongoLockRegistry(MongoTemplate mongoTemplate) {
        return new MongoLockRegistry(mongoTemplate, "change-stream-locks");
      }
    
      @Bean
      public SmartLifecycle changeStreamStarter(MessageListenerContainer container, MongoLockRegistry lockRegistry) {
        return new SmartLifecycle() {
          private Lock lock;
          private volatile boolean running = false;
    
          @Override
          public void start() {
            lock = lockRegistry.obtain(LOCK_KEY);
            try {
              // 尝试加锁,加锁成功才启动监听
              if (lock.tryLock(LOCK_TTL, TimeUnit.MILLISECONDS)) {
                container.start();
                running = true;
                log.info("当前实例获取到变更流锁,启动监听");
                // 开启后台线程定时续期锁,避免锁过期导致监听中断
                new Thread(() -> {
                  while (running) {
                    try {
                      Thread.sleep(LOCK_TTL / 2);
                      // Spring Integration MongoLockRegistry会自动处理锁续期,无需额外操作
                    } catch (InterruptedException e) {
                      Thread.currentThread().interrupt();
                    }
                  }
                }).start();
              }
            } catch (InterruptedException e) {
              Thread.currentThread().interrupt();
            }
          }
    
          @Override
          public void stop() {
            running = false;
            container.stop();
            if (lock != null) {
              lock.unlock();
            }
          }
    
          @Override
          public boolean isRunning() {
            return running;
          }
    
          @Override public int getPhase() { return 0; }
          @Override public boolean isAutoStartup() { return true; }
          @Override public void stop(Runnable callback) { stop(); callback.run(); }
        };
      }
    
      // 原来的MessageListenerContainer Bean保留,设置autoStartup为false即可
      @Bean
      MessageListenerContainer changeStreamListenerContainer(
          MongoTemplate template,
          PartyConsentAuditListener consentAuditListener,
          ErrorHandler errorHandler) {
        MessageListenerContainer messageListenerContainer =
            new MongoStreamListenerContainer(template, errorHandler);
        // 改为手动控制启动
        messageListenerContainer.setAutoStartup(false);
        ChangeStreamRequest<PartyConsentEntity> request =
            ChangeStreamRequest.builder(consentAuditListener)
                .collection("my-collection")
                .filter(newAggregation(match(where("operationType").in("insert", "update", "replace"))))
                .fullDocumentLookup(FullDocument.UPDATE_LOOKUP)
                .build();
        messageListenerContainer.register(request, MyEntity.class, errorHandler);
        log.info("mongo stream listener is registered");
        return messageListenerContainer;
      }
      // 其他Bean保持不变
    }
    

方案2:变更流分片负载模式(适合高负载场景,同时满足不重复+负载分担需求)

该模式利用MongoDB变更流的过滤能力,将事件按规则分片到不同实例处理,每个实例只处理自己分片的事件,完全避免重复同时实现流量分担:

  • 给每个实例分配唯一的分片编号(范围0~n-1,n为实例总数),可以通过配置参数、K8s StatefulSet序号、注册中心实例列表计算等方式获取
  • 改造ChangeStreamRequest的过滤规则,新增分片过滤条件:
    // 假设当前实例分片编号为instanceIndex,总实例数为instanceTotal
    ChangeStreamRequest<PartyConsentEntity> request =
        ChangeStreamRequest.builder(consentAuditListener)
            .collection("my-collection")
            .filter(newAggregation(
                match(where("operationType").in("insert", "update", "replace")),
                // 按文档_id哈希取模分片,保证同一个文档的所有变更都落到同一个实例处理
                match(AggregationExpression.from("{$eq:[{$mod:[{$hash:'$fullDocument._id'}, " + instanceTotal + "]}, " + instanceIndex + "]}"))
            ))
            .fullDocumentLookup(FullDocument.UPDATE_LOOKUP)
            .build();
    
  • 实例扩缩容时需要同步调整所有实例的分片编号和总实例数,配合变更流的resume token可以保证事件不丢失不重复。

兜底方案:审计集合幂等校验

无论采用哪种方案,都建议增加幂等校验避免极端场景下的重复数据:

  • 给审计集合新增唯一索引,索引字段为原文档_id + 操作类型 + 操作时间戳
  • 写入审计记录时如果出现唯一键冲突,直接忽略即可:
    @Override
    public void onMessage(Message<ChangeStreamDocument<Document>, MyEntity > message) {
      var update = message.getBody();
      if (update == null) {
        return;
      }
      try {
        AuditEntity audit = buildAuditEntity(update);
        mongoTemplate.insert(audit);
      } catch (DuplicateKeyException e) {
        log.info("重复审计事件,忽略,文档id:{}", update.getId());
      }
    }
    

内容的提问来源于stack exchange,提问作者dmca

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 06:24:03