update ruoyi-demo/src/main/java/com/ruoyi/demo/controller/queue/DelayedQueueController.java.

Signed-off-by: 月夜 <768242801@qq.com>
This commit is contained in:
月夜 2023-07-04 11:17:30 +00:00 committed by Gitee
parent a491534489
commit 6cfc161462
No known key found for this signature in database
GPG Key ID: 173E9B9CA92EEF8F

View File

@ -28,6 +28,9 @@ import java.util.concurrent.TimeUnit;
@RequestMapping("/demo/queue/delayed") @RequestMapping("/demo/queue/delayed")
public class DelayedQueueController { public class DelayedQueueController {
//注入线程池-需要在application.yml中将thread-pool.enabled设置为true
private final ThreadPoolTaskExecutor threadPoolTaskExecutor;
/** /**
* 订阅队列 * 订阅队列
* *
@ -38,9 +41,15 @@ public class DelayedQueueController {
log.info("通道: {} 监听中......", queueName); log.info("通道: {} 监听中......", queueName);
// 项目初始化设置一次即可 // 项目初始化设置一次即可
QueueUtils.subscribeBlockingQueue(queueName, (String orderNum) -> { QueueUtils.subscribeBlockingQueue(queueName, (String orderNum) -> {
//如业务代码部分使用到了redis相关操作需要将业务逻辑异步否则将会报错Sync methods can't be invoked from async/rx/reactive listeners
// 观察接收时间 // 观察接收时间
log.info("通道: {}, 收到数据: {}", queueName, orderNum); log.info("通道: {}, 收到数据: {}", queueName, orderNum);
//业务代码
//如业务代码部分使用到了redis相关操作需要将业务逻辑异步否则将会报错Sync methods can't be invoked from async/rx/reactive listeners
//示例
threadPoolTaskExecutor.execute(() -> {
String str = RedisUtils.getCacheObject("test");
});
}); });
return R.ok("操作成功"); return R.ok("操作成功");
} }