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

如何在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,新请求会等待前一个请求的信号量完成后再执行,执行完毕后移除锁,既保证同一用户的更新串行,又不影响不同用户的并发操作。

方案二:数据库乐观锁(多实例场景适用)

给用户实体添加版本号字段,利用数据库的乐观锁机制控制并发更新,失败的请求通过重试保证最终一致性:

  1. 修改User实体:
@Entity
public class User {
    // 原有字段
    private Double balance;
    @Version
    private Integer version; // 乐观锁版本号

    // getter、setter方法
}
  1. 修改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,通过重试机制重新查询最新数据再更新,最终保证数据正确。

方案三:数据库原子更新(最可靠方案)

直接让数据库执行余额加减的原子操作,彻底避免先查后改的步骤,从根源解决并发问题:

  1. 在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);
}
  1. 修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 21:30:07