WebFlux/Project Reactor中如何实现锁机制?非阻塞锁方案问询
解决WebFlux/Project Reactor中的分布式锁与并发控制问题
一、基于Redis的非阻塞分布式锁实现
在Reactor生态中,必须避免使用阻塞式Redis客户端,改用Lettuce或Redisson的Reactive API实现分布式锁,完全贴合非阻塞设计原则。
核心逻辑(Lettuce手动实现)
利用Redis的SET key value NX EX timeout命令实现原子性锁获取,通过Reactor操作符串联锁的获取、业务执行、释放流程:
import io.lettuce.core.RedisClient; import io.lettuce.core.api.reactive.RedisStringReactiveCommands; import io.lettuce.core.SetArgs; import reactor.core.publisher.Mono; import java.util.UUID; public class RedisDistributedLock { private final RedisStringReactiveCommands<String, String> redisCommands; private final String lockKeyPrefix = "lock:"; private final long lockTimeoutSeconds = 30; public RedisDistributedLock(RedisClient redisClient) { this.redisCommands = redisClient.connect().reactive(); } public <T> Mono<T> executeWithLock(String entityId, Mono<T> businessLogic) { String lockKey = lockKeyPrefix + entityId; String lockValue = UUID.randomUUID().toString(); // 用唯一值避免误释放其他线程的锁 // 1. 原子性尝试获取锁 return redisCommands.set(lockKey, lockValue, SetArgs.Builder.nx().ex(lockTimeoutSeconds)) .filter("OK"::equals) // 仅当锁获取成功时继续执行 .flatMap(ok -> businessLogic // 2. 执行业务逻辑,无论成功失败都释放锁 .doFinally(signal -> releaseLock(lockKey, lockValue).subscribe()) ) // 3. 锁获取失败时返回明确错误 .switchIfEmpty(Mono.error(new IllegalStateException("实体[" + entityId + "]正在处理中,请稍后重试"))); } private Mono<Boolean> releaseLock(String lockKey, String lockValue) { // Lua脚本保证释放锁的原子性:只有持有锁的线程才能释放 String luaScript = "if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) else return 0 end"; return redisCommands.eval(luaScript, ScriptOutputType.BOOLEAN, new String[]{lockKey}, lockValue); } }
更便捷的方案:Redisson Reactive锁
Redisson提供了开箱即用的Reactive分布式锁实现,无需手动编写Lua脚本或锁逻辑:
import org.redisson.api.RLockReactive; import org.redisson.api.RedissonReactiveClient; import reactor.core.publisher.Mono; import java.util.concurrent.TimeUnit; public class RedissonReactiveLock { private final RedissonReactiveClient redissonClient; public RedissonReactiveLock(RedissonReactiveClient redissonClient) { this.redissonClient = redissonClient; } public <T> Mono<T> executeWithLock(String entityId, Mono<T> businessLogic) { RLockReactive lock = redissonClient.getLock("lock:" + entityId); // 非阻塞获取锁,超时时间30秒,锁自动过期 return lock.lock(30, TimeUnit.SECONDS) .then(businessLogic) .doFinally(signal -> lock.unlock().subscribe()); } }
二、本地并发控制(非分布式场景)
如果仅需本地按条件控制并发,可通过非阻塞式锁检查实现,避免使用synchronized或阻塞式lock():
import java.util.concurrent.locks.ReentrantLock; import reactor.core.publisher.Mono; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; public class LocalNonBlockingLock { // 按实体ID维护锁映射,实现细粒度控制 private final Map<String, ReentrantLock> entityLocks = new ConcurrentHashMap<>(); public Mono<Void> processEntity(String entityId) { ReentrantLock lock = entityLocks.computeIfAbsent(entityId, k -> new ReentrantLock()); return Mono.fromSupplier(lock::tryLock) .filter(locked -> locked) .flatMap(locked -> { // 执行业务逻辑 return doEntityProcessing(entityId) .doFinally(signal -> { if (lock.isHeldByCurrentThread()) { lock.unlock(); } }); }) .switchIfEmpty(Mono.error(new IllegalStateException("实体[" + entityId + "]正在处理"))); } private Mono<Void> doEntityProcessing(String entityId) { // 业务逻辑实现 return Mono.empty(); } }
这里使用tryLock()的非阻塞特性,结合Reactor操作符处理锁获取结果,完全不会阻塞事件循环线程。
三、关键注意事项
- 锁粒度控制:尽量按实体ID等维度加锁,避免全局锁导致的性能瓶颈
- 自动过期机制:必须为锁设置超时时间,避免因服务崩溃导致死锁
- 重试策略:锁获取失败时,可通过
retryWhen操作符实现指数退避重试,提升用户体验:.switchIfEmpty(Mono.error(new LockAcquireException())) .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)) .filter(LockAcquireException.class::isInstance)) - 避免长耗时操作:WebFlux场景下,锁持有时间越短越好,避免阻塞其他请求
内容的提问来源于stack exchange,提问作者Serg
相关产品推荐
相关产品推荐

