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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 19:18:19