多实例部署下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
相关产品推荐
相关产品推荐

