rblockingqueue是redisson提供的分布式阻塞队列,完整实现java blockingqueue接口,基于redis实现跨jvm任务分发,支持put/take阻塞语义、动态容量调整及原子性操作,适用于异步任务处理、流量控制等场景。

在线程池内部利用 Redisson 实现分布式任务的“异步领取”,核心是把线程池作为本地执行载体,而任务的分发、抢占与状态同步交由 Redisson 的分布式能力完成。这不是让线程池去“拉取”Redis里的任务,而是用 Redisson 的分布式队列(如 RBlockingQueue)或分布式信号量(RPermitExpirableSemaphore)来协调多个 JVM 实例中的线程池,确保同一任务只被一个节点的一个线程“领走”并执行。
用 RBlockingQueue 实现任务的异步领取
这是最贴近“领取”语义的方式:任务提前入队,各节点线程池中的线程竞争式地从队列中 poll 或 take 任务,Redisson 保证操作原子性。
- 任务发布方(比如 API 接口)调用
queue.offer(task)入队,支持序列化(需实现Serializable) - 每个服务实例启动一个后台线程(或固定线程池中的工作线程),循环执行:
Task task = queue.poll(1, TimeUnit.SECONDS); // 非阻塞,带超时
若拿到任务,则交由业务逻辑处理;没拿到就继续下一轮 - 推荐搭配
RExecutorService提交真正耗时逻辑,避免阻塞领取线程
用 RPermitExpirableSemaphore 控制并发领取数
适用于需要限制“同时最多几个线程能尝试领取新任务”的场景(例如防雪崩、控资源),配合本地线程池做节流。
- 初始化一个可过期许可的信号量:
RPermitExpirableSemaphore semaphore = redisson.getPermitExpirableSemaphore("task:acquire:limit"); - 线程在准备领取前先异步申请许可:
semaphore.tryAcquireAsync(1, 30, TimeUnit.SECONDS),返回CompletableFuture<string></string>(许可 ID) - 获取成功后,再去查 Redis 中待领取的任务列表(如
RSet<string> pendingTasks</string>),用 Lua 脚本原子性地SPOP一个任务并标记为“已领取” - 执行完任务后,记得
semaphore.tryReleaseAsync(permitId)归还许可
结合 RExecutorService 实现“提交即领取”的轻量调度
如果你不希望自己维护领取循环,Redisson 自带的分布式执行器本身就是一种“自动领取”机制——任务提交到命名执行器后,任意注册了该执行器的节点都会自动争抢执行。
- 定义一个
Callable<string></string>或Runnable任务类(必须Serializable) - 通过
redisson.getExecutorService("myTaskPool").submit(task)提交 - 在每个服务节点上启动
RedissonNode并注册同名执行器(setExecutorServiceWorkers(...)),它会监听 Redis 中的任务队列,有任务就用本地线程池执行 - 整个过程对调用方透明,天然支持失败重试、结果回调(
RFuture.onComplete())
关键注意事项
避免常见陷阱才能让“异步领取”真正可靠:
-
任务必须可序列化:所有字段不能含线程局部变量、Spring 上下文引用、不可序列化资源(如
Connection) - 领取和执行要分离:领取线程只负责“拿任务 ID”,执行交给独立线程池,防止领取卡住导致后续任务积压
-
失败任务要有兜底:比如领取后设置 Redis 过期时间(
EXPIRE task:{id} 300),配合定时扫描补偿未完成任务 - 锁粒度要合理:不要用全局锁控制领取,优先用队列或信号量,减少竞争
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











