From ca2cb70cf963ed2206508c36fcbe85601e9261ec Mon Sep 17 00:00:00 2001 From: ljl <> Date: Wed, 7 Jun 2023 10:24:00 +0800 Subject: [PATCH] =?UTF-8?q?=E5=90=8C=E6=AD=A5=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../demo/init/RedisDelayQueueHandle.java | 11 ++++ .../demo/init/RedisDelayQueueRunner.java | 57 ++++++++++++++++++ .../demo/init/RedisDelayedQueueInit.java | 60 +++++++++++++++++++ ruoyi-file/pom.xml | 2 +- ruoyi-rabbitmq/pom.xml | 34 +++++++++++ ruoyi-rsaencrypt/pom.xml | 2 +- .../ruoyi/system/service/SysLoginService.java | 10 ++-- ruoyi-work/pom.xml | 2 +- 8 files changed, 170 insertions(+), 8 deletions(-) create mode 100644 ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueHandle.java create mode 100644 ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueRunner.java create mode 100644 ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayedQueueInit.java create mode 100644 ruoyi-rabbitmq/pom.xml diff --git a/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueHandle.java b/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueHandle.java new file mode 100644 index 000000000..9dd3268ca --- /dev/null +++ b/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueHandle.java @@ -0,0 +1,11 @@ +package com.ruoyi.demo.init; + +/** + * 延迟队列执行器 + * Created by LPB on 2021/04/20. + */ +public interface RedisDelayQueueHandle { + + void execute(T t); + +} diff --git a/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueRunner.java b/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueRunner.java new file mode 100644 index 000000000..99dc87d56 --- /dev/null +++ b/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayQueueRunner.java @@ -0,0 +1,57 @@ +package com.ruoyi.demo.init;/* +package cn.dbtalents.checktalents.init; +import cn.dbtalents.checktalents.enums.RedisDelayQueueEnum; +import cn.dbtalents.checktalents.handle.RedisDelayQueueHandle; +import cn.dbtalents.checktalents.util.RedisDelayQueueUtil; +import cn.hutool.extra.spring.SpringUtil; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.CommandLineRunner; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import org.springframework.stereotype.Component; +import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +*/ +/** + * 启动延迟队列 + *//* + +@Slf4j +@Component +public class RedisDelayQueueRunner implements CommandLineRunner { + + @Autowired + private RedisDelayQueueUtil redisDelayQueueUtil; + */ +/* @Autowired + private ThreadPoolTaskExecutor threadPool; + ThreadPoolExecutor executorService = new ThreadPoolExecutor(10, 50, 30, TimeUnit.SECONDS, + new LinkedBlockingQueue(1000), Executors.defaultThreadFactory());*//* + + @Override + public void run(String... args) { + */ +/* threadPool.execute(() -> { + while (true){*//* + + try { + RedisDelayQueueEnum[] queueEnums = RedisDelayQueueEnum.values(); + for (RedisDelayQueueEnum queueEnum : queueEnums) { + Object value = redisDelayQueueUtil.getDelayQueue(queueEnum.getCode()); + if (value != null) { + RedisDelayQueueHandle redisDelayQueueHandle = SpringUtil.getBean(queueEnum.getBeanId()); + redisDelayQueueHandle.execute(value); + } + } + } catch (InterruptedException e) { + log.error("(Redis延迟队列异常中断) {}", e.getMessage()); + } +// } +// }); + log.info("(Redis延迟队列启动成功)"); + } +} +*/ diff --git a/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayedQueueInit.java b/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayedQueueInit.java new file mode 100644 index 000000000..51a7242d8 --- /dev/null +++ b/ruoyi-demo/src/main/java/com/ruoyi/demo/init/RedisDelayedQueueInit.java @@ -0,0 +1,60 @@ +package com.ruoyi.demo.init; +import lombok.extern.slf4j.Slf4j; +import org.redisson.api.RBlockingQueue; +import org.redisson.api.RedissonClient; +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.stereotype.Component; +import java.util.Map; + +/** + * redis 延时队列初始化 + */ +@Component +@Slf4j +public class RedisDelayedQueueInit implements ApplicationContextAware { + + @Autowired + private RedissonClient redissonClient; + + /** + * 获取应用上下文并获取相应的接口实现类 + * @param applicationContext + * @throws BeansException + */ + @Override + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + Map map = applicationContext.getBeansOfType(RedisDelayQueueHandle.class); + for (Map.Entry taskEventListenerEntry : map.entrySet()) { + String listenerName = taskEventListenerEntry.getValue().getClass().getName(); + startThread(listenerName, taskEventListenerEntry.getValue()); + } + } + + /** + * 启动线程获取队列 + * @param queueName 队列名称 + * @param redisDelayedQueueListener 任务回调监听 + */ + private void startThread(String queueName, RedisDelayQueueHandle redisDelayedQueueListener) { + RBlockingQueue blockingFairQueue = redissonClient.getBlockingQueue(queueName); + //由于此线程需要常驻,可以新建线程,不用交给线程池管理 + Thread thread = new Thread(() -> { + log.info("启动监听队列线程" + queueName); + while (true) { + try { + T t = blockingFairQueue.take(); + log.info("监听队列线程{},获取到值:{}", queueName, t); + redisDelayedQueueListener.execute(t); + } catch (Exception e) { + log.info("监听队列线程错误,", e); + } + } + }); + thread.setName(queueName); + thread.start(); + log.info("(Redis延迟队列启动成功)"); + } +} diff --git a/ruoyi-file/pom.xml b/ruoyi-file/pom.xml index 117edcebe..807c6da27 100644 --- a/ruoyi-file/pom.xml +++ b/ruoyi-file/pom.xml @@ -5,7 +5,7 @@ ruoyi-vue-plus com.ruoyi - 4.6.0 + 4.7.0 4.0.0 diff --git a/ruoyi-rabbitmq/pom.xml b/ruoyi-rabbitmq/pom.xml new file mode 100644 index 000000000..1107d146a --- /dev/null +++ b/ruoyi-rabbitmq/pom.xml @@ -0,0 +1,34 @@ + + + + ruoyi-vue-plus + com.ruoyi + 4.7.0 + + 4.0.0 + jar + ruoyi-rabbitmq + + + 消息队列 + + + + + + + com.ruoyi + ruoyi-common + + + + org.springframework.boot + spring-boot-starter-amqp + + + + + + diff --git a/ruoyi-rsaencrypt/pom.xml b/ruoyi-rsaencrypt/pom.xml index dc85a7309..7b12a1c33 100644 --- a/ruoyi-rsaencrypt/pom.xml +++ b/ruoyi-rsaencrypt/pom.xml @@ -5,7 +5,7 @@ ruoyi-vue-plus com.ruoyi - 4.6.0 + 4.7.0 4.0.0 diff --git a/ruoyi-system/src/main/java/com/ruoyi/system/service/SysLoginService.java b/ruoyi-system/src/main/java/com/ruoyi/system/service/SysLoginService.java index f35b4660a..dab9376e4 100644 --- a/ruoyi-system/src/main/java/com/ruoyi/system/service/SysLoginService.java +++ b/ruoyi-system/src/main/java/com/ruoyi/system/service/SysLoginService.java @@ -115,8 +115,8 @@ public class SysLoginService { } /** - * 小程序登录 - * @param xcxCode + * 邮箱登录 + * @param * @return */ public String emailLogin(String email, String emailCode) { @@ -125,7 +125,7 @@ public class SysLoginService { checkLogin(LoginType.EMAIL, user.getUserName(), () -> !validateEmailCode(email, emailCode)); // 此处可根据登录用户的数据不同 自行创建 loginUser - LoginUser loginUser = buildLoginUser(user); + LoginUser loginUser = buildLoginSysUser(user); // 生成token LoginHelper.loginByDevice(loginUser, DeviceType.APP); @@ -257,7 +257,7 @@ public class SysLoginService { } private SysUser loadUserByEmail(String email) { - SysUser user = userMapper.selectOne(new LambdaQueryWrapper() + SysUser user = sysUserMapper.selectOne(new LambdaQueryWrapper() .select(SysUser::getPhonenumber, SysUser::getStatus) .eq(SysUser::getEmail, email)); if (ObjectUtil.isNull(user)) { @@ -267,7 +267,7 @@ public class SysLoginService { log.info("登录用户:{} 已被停用.", email); throw new UserException("user.blocked", email); } - return userMapper.selectUserByEmail(email); + return sysUserMapper.selectUserByEmail(email); } private SysUser loadUserByOpenid(String openid) { diff --git a/ruoyi-work/pom.xml b/ruoyi-work/pom.xml index 9b25ed9bb..d6bbfeecf 100644 --- a/ruoyi-work/pom.xml +++ b/ruoyi-work/pom.xml @@ -5,7 +5,7 @@ ruoyi-vue-plus com.ruoyi - 4.6.0 + 4.7.0 4.0.0