SpringBoot WebFlux对接MongoDB集群出现重复数据问题排查
Reactive链中手动调用
subscribe()导致重复触发
在doOnNext中直接调用subscribe()会使状态保存操作脱离原Reactive链独立异步执行,当上游流被多次订阅或触发时,该保存操作会被重复调用。同时手动订阅无法保证与原链的执行顺序,进一步增加重复执行的概率。前置检查存在竞态条件
尝试的"查询-保存"逻辑是两个独立操作,无原子性保障。在高并发或异步场景下,多个请求可能同时通过查询检查(都未找到对应记录),进而执行保存操作,导致重复记录。缺少数据库层面的唯一约束
commandStatus集合未针对commandId和commandStatus创建复合唯一索引,MongoDB无法从底层阻止重复数据插入。readPreference=secondary的潜在影响
查询操作优先从从节点读取,而从节点的数据同步存在延迟,可能导致前置检查时读取到的不是最新数据(主节点已保存但从节点未同步),进而误判为不存在并执行保存。
1. 修正Reactive链调用方式,移除手动subscribe()
将状态保存操作整合到Reactive链中,使用flatMap替代doOnNext+subscribe(),保证操作的顺序性和单次触发:
@Override public Flux<CommandGetResponse> getCommands(CommandGetRequest request) { return deviceInfoService.findByDeviceId(request.getDeviceId()) .flatMap(commandService::findByDevice) // 将状态保存作为链的一部分,避免独立订阅 .flatMap(command -> commandStatusService.createOrGetCommandStatus(command, CommandStatus.SENT_TO_DEVICE) .thenReturn(command) ) .flatMap(command -> Flux.fromIterable(command.getStatuses()) .map(statuses -> CommandGetResponse.from(command, statuses)) ); }
2. 实现原子化的"查询-保存"操作
使用MongoDB的findOneAndUpdate配合upsert参数,实现原子性的创建或查询,彻底避免竞态条件:
@Override public Mono<CommandStatusDocument> createOrGetCommandStatus(CommandDocument commandDocument, CommandStatus status) { // 构建查询条件:commandId + commandStatus唯一 Query query = Query.query(Criteria.where("commandId").is(commandDocument.getId()) .and("commandStatus").is(status)); // 仅当文档不存在时插入字段 Update update = new Update() .setOnInsert("commandId", commandDocument.getId()) .setOnInsert("commandStatus", status) .setOnInsert("timestamp", Instant.now()); // upsert=true:不存在则插入,存在则返回现有文档;returnNew=true:返回操作后的文档 return reactiveMongoTemplate.findAndModify(query, update, new FindAndModifyOptions().upsert(true).returnNew(true), CommandStatusDocument.class); }
3. 添加复合唯一索引
从数据库层面阻止重复数据,在commandStatus集合创建commandId和commandStatus的复合唯一索引:
方式1:MongoDB命令行创建
db.commandStatus.createIndex({commandId: 1, commandStatus: 1}, {unique: true})
方式2:实体类注解创建
@Data @Builder @AllArgsConstructor @NoArgsConstructor @Document(collection = "commandStatus") // 定义复合唯一索引 @CompoundIndex(name = "commandId_status_unique", def = "{'commandId': 1, 'commandStatus': 1}", unique = true) public class CommandStatusDocument { @Id private String id; private String commandId; private Instant timestamp; private CommandStatus commandStatus; }
4. 调整MongoDB读取偏好(可选)
如果业务需要强一致性查询结果,将readPreference改为primary,确保查询读取主节点的最新数据:
spring: data: mongodb: uri: mongodb://${MONGO_USERNAME}:${MONGO_PASSWORD}@${MONGO_CLUSTER}/${DATABASE_NAME}?authSource=admin&readPreference=primary&replicaSet=rs0&minPoolSize=20
内容的提问来源于stack exchange,提问作者Virkom

