QueueUtils.java 7.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222
  1. package com.ruoyi.common.utils.redis;
  2. import com.ruoyi.common.utils.spring.SpringUtils;
  3. import lombok.AccessLevel;
  4. import lombok.NoArgsConstructor;
  5. import org.redisson.api.*;
  6. import java.util.Comparator;
  7. import java.util.concurrent.TimeUnit;
  8. import java.util.function.Consumer;
  9. /**
  10. * 分布式队列工具
  11. * 轻量级队列 重量级数据量 请使用 MQ
  12. * 要求 redis 5.X 以上
  13. *
  14. * @author Lion Li
  15. * @version 3.6.0 新增
  16. */
  17. @NoArgsConstructor(access = AccessLevel.PRIVATE)
  18. public class QueueUtils {
  19. private static final RedissonClient CLIENT = SpringUtils.getBean(RedissonClient.class);
  20. /**
  21. * 获取客户端实例
  22. */
  23. public static RedissonClient getClient() {
  24. return CLIENT;
  25. }
  26. /**
  27. * 添加延迟队列数据 默认毫秒
  28. *
  29. * @param queueName 队列名
  30. * @param data 数据
  31. * @param time 延迟时间
  32. */
  33. public static <T> void addDelayedQueueObject(String queueName, T data, long time) {
  34. addDelayedQueueObject(queueName, data, time, TimeUnit.MILLISECONDS);
  35. }
  36. /**
  37. * 添加延迟队列数据
  38. *
  39. * @param queueName 队列名
  40. * @param data 数据
  41. * @param time 延迟时间
  42. * @param timeUnit 单位
  43. */
  44. public static <T> void addDelayedQueueObject(String queueName, T data, long time, TimeUnit timeUnit) {
  45. RBlockingQueue<T> queue = CLIENT.getBlockingQueue(queueName);
  46. RDelayedQueue<T> delayedQueue = CLIENT.getDelayedQueue(queue);
  47. delayedQueue.offer(data, time, timeUnit);
  48. }
  49. /**
  50. * 获取一个延迟队列数据 没有数据返回 null
  51. *
  52. * @param queueName 队列名
  53. */
  54. public static <T> T getDelayedQueueObject(String queueName) {
  55. RBlockingQueue<T> queue = CLIENT.getBlockingQueue(queueName);
  56. RDelayedQueue<T> delayedQueue = CLIENT.getDelayedQueue(queue);
  57. return delayedQueue.poll();
  58. }
  59. /**
  60. * 删除延迟队列数据
  61. */
  62. public static <T> boolean removeDelayedQueueObject(String queueName, T data) {
  63. RBlockingQueue<T> queue = CLIENT.getBlockingQueue(queueName);
  64. RDelayedQueue<T> delayedQueue = CLIENT.getDelayedQueue(queue);
  65. return delayedQueue.remove(data);
  66. }
  67. /**
  68. * 销毁延迟队列 所有阻塞监听 报错
  69. */
  70. public static <T> void destroyDelayedQueue(String queueName) {
  71. RBlockingQueue<T> queue = CLIENT.getBlockingQueue(queueName);
  72. RDelayedQueue<T> delayedQueue = CLIENT.getDelayedQueue(queue);
  73. delayedQueue.destroy();
  74. }
  75. /**
  76. * 尝试设置 优先队列比较器 用于排序优先级
  77. *
  78. * @param queueName 队列名
  79. * @param comparator 比较器
  80. */
  81. public static <T> boolean trySetPriorityQueueComparator(String queueName, Comparator<T> comparator) {
  82. RPriorityBlockingQueue<T> priorityBlockingQueue = CLIENT.getPriorityBlockingQueue(queueName);
  83. return priorityBlockingQueue.trySetComparator(comparator);
  84. }
  85. /**
  86. * 尝试设置 优先队列比较器 用于排序优先级
  87. *
  88. * @param queueName 队列名
  89. * @param comparator 比较器
  90. * @param destroy 已存在是否销毁
  91. */
  92. public static <T> boolean trySetPriorityQueueComparator(String queueName, Comparator<T> comparator, boolean destroy) {
  93. RPriorityBlockingQueue<T> priorityBlockingQueue = CLIENT.getPriorityBlockingQueue(queueName);
  94. if (priorityBlockingQueue.isExists() && destroy) {
  95. destroyPriorityQueueObject(queueName);
  96. }
  97. return priorityBlockingQueue.trySetComparator(comparator);
  98. }
  99. /**
  100. * 添加优先队列数据
  101. *
  102. * @param queueName 队列名
  103. * @param data 数据
  104. */
  105. public static <T> boolean addPriorityQueueObject(String queueName, T data) {
  106. RPriorityBlockingQueue<T> priorityBlockingQueue = CLIENT.getPriorityBlockingQueue(queueName);
  107. return priorityBlockingQueue.offer(data);
  108. }
  109. /**
  110. * 获取一个优先队列数据 没有数据返回 null
  111. *
  112. * @param queueName 队列名
  113. */
  114. public static <T> T getPriorityQueueObject(String queueName) {
  115. RPriorityBlockingQueue<T> priorityBlockingQueue = CLIENT.getPriorityBlockingQueue(queueName);
  116. return priorityBlockingQueue.poll();
  117. }
  118. /**
  119. * 删除优先队列数据
  120. */
  121. public static <T> boolean removePriorityQueueObject(String queueName, T data) {
  122. RPriorityBlockingQueue<T> priorityBlockingQueue = CLIENT.getPriorityBlockingQueue(queueName);
  123. return priorityBlockingQueue.remove(data);
  124. }
  125. /**
  126. * 销毁优先队列
  127. */
  128. public static boolean destroyPriorityQueueObject(String queueName) {
  129. RPriorityBlockingQueue<?> priorityBlockingQueue = CLIENT.getPriorityBlockingQueue(queueName);
  130. return priorityBlockingQueue.delete();
  131. }
  132. /**
  133. * 尝试设置 有界队列 容量 用于限制数量
  134. *
  135. * @param queueName 队列名
  136. * @param capacity 容量
  137. */
  138. public static <T> boolean trySetBoundedQueueCapacity(String queueName, int capacity) {
  139. RBoundedBlockingQueue<T> boundedBlockingQueue = CLIENT.getBoundedBlockingQueue(queueName);
  140. return boundedBlockingQueue.trySetCapacity(capacity);
  141. }
  142. /**
  143. * 尝试设置 有界队列 容量 用于限制数量
  144. *
  145. * @param queueName 队列名
  146. * @param capacity 容量
  147. * @param destroy 已存在是否销毁
  148. */
  149. public static <T> boolean trySetBoundedQueueCapacity(String queueName, int capacity, boolean destroy) {
  150. RBoundedBlockingQueue<T> boundedBlockingQueue = CLIENT.getBoundedBlockingQueue(queueName);
  151. if (boundedBlockingQueue.isExists() && destroy) {
  152. destroyBoundedQueueObject(queueName);
  153. }
  154. return boundedBlockingQueue.trySetCapacity(capacity);
  155. }
  156. /**
  157. * 添加有界队列数据
  158. *
  159. * @param queueName 队列名
  160. * @param data 数据
  161. * @return 添加成功 true 已达到界限 false
  162. */
  163. public static <T> boolean addBoundedQueueObject(String queueName, T data) {
  164. RBoundedBlockingQueue<T> boundedBlockingQueue = CLIENT.getBoundedBlockingQueue(queueName);
  165. return boundedBlockingQueue.offer(data);
  166. }
  167. /**
  168. * 获取一个有界队列数据 没有数据返回 null
  169. *
  170. * @param queueName 队列名
  171. */
  172. public static <T> T getBoundedQueueObject(String queueName) {
  173. RBoundedBlockingQueue<T> boundedBlockingQueue = CLIENT.getBoundedBlockingQueue(queueName);
  174. return boundedBlockingQueue.poll();
  175. }
  176. /**
  177. * 删除有界队列数据
  178. */
  179. public static <T> boolean removeBoundedQueueObject(String queueName, T data) {
  180. RBoundedBlockingQueue<T> boundedBlockingQueue = CLIENT.getBoundedBlockingQueue(queueName);
  181. return boundedBlockingQueue.remove(data);
  182. }
  183. /**
  184. * 销毁有界队列
  185. */
  186. public static boolean destroyBoundedQueueObject(String queueName) {
  187. RBoundedBlockingQueue<?> boundedBlockingQueue = CLIENT.getBoundedBlockingQueue(queueName);
  188. return boundedBlockingQueue.delete();
  189. }
  190. /**
  191. * 订阅阻塞队列(可订阅所有实现类 例如: 延迟 优先 有界 等)
  192. */
  193. public static <T> void subscribeBlockingQueue(String queueName, Consumer<T> consumer) {
  194. RBlockingQueue<T> queue = CLIENT.getBlockingQueue(queueName);
  195. queue.subscribeOnElements(consumer);
  196. }
  197. }