如何在WebFlux中实现REST方法同步解决并发余额更新异常?
解决Reactive环境下用户余额并发更新丢失的问题
你遇到的是典型并发更新丢失问题:多个请求同时查询到同一个用户的余额(日志里三个请求都拿到770.0),各自在内存中修改后保存,最终只有最后一次更新生效,前面的操作都被覆盖了。以下是几种针对Reactive场景的解决方案:
方案一:本地信号量串行化(单实例场景适用)
Reactive环境下不能用普通synchronized(会阻塞线程违背非阻塞原则),可以用ConcurrentHashMap维护每个用户ID对应的信号量,保证同一用户的更新请求串行执行:
@Slf4j @RestController @RequestMapping("/users") @RequiredArgsConstructor public class UserController { private final UserRepository userRepository; // 维护每个userId的串行信号量,隔离不同用户的更新请求 private final ConcurrentHashMap<Long, MonoProcessor<Void>> userLocks = new ConcurrentHashMap<>(); @PutMapping("/{userId}:updateBalance") public Mono<User> updateUser(@PathVariable("userId") Long userId, @RequestParam("diff") Double diff) { // 获取或创建当前用户的锁,等待前一个请求完成后再执行更新 return Mono.defer(() -> userLocks.computeIfAbsent(userId, k -> MonoProcessor.create())) .then(doUpdate(userId, diff)) .doFinally(signalType -> userLocks.remove(userId)); } // 抽离核心更新逻辑 private Mono<User> doUpdate(Long userId, Double diff) { return userRepository.findById(userId) .flatMap(user -> { log.info(String.format("current balance - %s, will be - %s", user.getBalance(), user.getBalance() + diff)); user.setBalance(user.getBalance() + diff); return userRepository.save(user); }); } }
原理:每个用户ID对应一个MonoProcessor,新请求会等待前一个请求的信号量完成后再执行,执行完毕后移除锁,既保证同一用户的更新串行,又不影响不同用户的并发操作。
方案二:数据库乐观锁(多实例场景适用)
给用户实体添加版本号字段,利用数据库的乐观锁机制控制并发更新,失败的请求通过重试保证最终一致性:
- 修改User实体:
@Entity public class User { // 原有字段 private Double balance; @Version private Integer version; // 乐观锁版本号 // getter、setter方法 }
- 修改Controller更新逻辑,添加重试:
@PutMapping("/{userId}:updateBalance") public Mono<User> updateUser(@PathVariable("userId") Long userId, @RequestParam("diff") Double diff) { return Mono.defer(() -> userRepository.findById(userId) .flatMap(user -> { log.info(String.format("current balance - %s, will be - %s", user.getBalance(), user.getBalance() + diff)); user.setBalance(user.getBalance() + diff); return userRepository.save(user); })) // 乐观锁失败后重试3次,采用指数退避策略 .retryWhen(Retry.backoff(3, Duration.ofMillis(100)) .filter(throwable -> throwable instanceof OptimisticLockingFailureException)); }
原理:并发更新时只有第一个请求的版本号能匹配数据库记录,其他请求会抛出OptimisticLockingFailureException,通过重试机制重新查询最新数据再更新,最终保证数据正确。
方案三:数据库原子更新(最可靠方案)
直接让数据库执行余额加减的原子操作,彻底避免先查后改的步骤,从根源解决并发问题:
- 在UserRepository中添加原子更新方法:
public interface UserRepository extends ReactiveCrudRepository<User, Long> { @Modifying @Query("UPDATE User u SET u.balance = u.balance + :diff WHERE u.id = :userId") Mono<Integer> updateBalance(@Param("userId") Long userId, @Param("diff") Double diff); // 更新后查询最新用户信息返回 Mono<User> findById(Long userId); }
- 修改Controller逻辑:
@PutMapping("/{userId}:updateBalance") public Mono<User> updateUser(@PathVariable("userId") Long userId, @RequestParam("diff") Double diff) { return userRepository.updateBalance(userId, diff) .then(userRepository.findById(userId)); }
原理:数据库的UPDATE语句是原子执行的,多个并发请求的更新会被数据库依次处理,不会出现丢失更新的情况,是多实例部署场景下的最优选择。
内容的提问来源于stack exchange,提问作者Dmirii Berrington
相关产品推荐
相关产品推荐

