最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot集成Redisson實(shí)現(xiàn)消息隊(duì)列的示例代碼

 更新時(shí)間:2024年10月10日 09:12:37   作者:入秋的大橘  
本文介紹了如何在SpringBoot中通過集成Redisson來實(shí)現(xiàn)消息隊(duì)列的功能,包括RedisQueue、RedisQueueInit、RedisQueueListener、RedisQueueService等相關(guān)組件的實(shí)現(xiàn)和測試,感興趣的可以了解一下

包含組件內(nèi)容

  • RedisQueue:消息隊(duì)列監(jiān)聽標(biāo)識
  • RedisQueueInit:Redis隊(duì)列監(jiān)聽器
  • RedisQueueListener:Redis消息隊(duì)列監(jiān)聽實(shí)現(xiàn)
  • RedisQueueService:Redis消息隊(duì)列服務(wù)工具

代碼實(shí)現(xiàn)

RedisQueue

import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;

/**
 * Redis消息隊(duì)列注解
 */
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
public @interface RedisQueue {
    /**
     * 隊(duì)列名
     */
    String value();
}

RedisQueueInit

import jakarta.annotation.Resource;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import lombok.extern.slf4j.Slf4j;
import org.jetbrains.annotations.NotNull;
import org.redisson.RedissonShutdownException;
import org.redisson.api.RBlockingQueue;
import org.redisson.api.RedissonClient;
import org.springframework.beans.BeansException;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationContextAware;
import org.springframework.stereotype.Component;

/**
 * 初始化Redis隊(duì)列監(jiān)聽器
 *
 * @author 十八
 * @createTime 2024-09-09 22:49
 */
@Slf4j
@Component
public class RedisQueueInit implements ApplicationContextAware {

    public static final String REDIS_QUEUE_PREFIX = "redis-queue";
    final AtomicBoolean shutdownRequested = new AtomicBoolean(false);
    @Resource
    private RedissonClient redissonClient;
    private ExecutorService executorService;

    public static String buildQueueName(String queueName) {
        return REDIS_QUEUE_PREFIX + ":" + queueName;
    }

    @Override
    public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
        Map<String, RedisQueueListener> queueListeners = applicationContext.getBeansOfType(RedisQueueListener.class);
        if (!queueListeners.isEmpty()) {
            executorService = createThreadPool();
            for (Map.Entry<String, RedisQueueListener> entry : queueListeners.entrySet()) {
                RedisQueue redisQueue = entry.getValue().getClass().getAnnotation(RedisQueue.class);
                if (redisQueue != null) {
                    String queueName = redisQueue.value();
                    executorService.submit(() -> listenQueue(queueName, entry.getValue()));
                }
            }
        }
    }

    private ExecutorService createThreadPool() {
        return new ThreadPoolExecutor(
                Runtime.getRuntime().availableProcessors() * 2,
                Runtime.getRuntime().availableProcessors() * 4,
                60L, TimeUnit.SECONDS,
                new LinkedBlockingQueue<>(100),
                new NamedThreadFactory(REDIS_QUEUE_PREFIX),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );
    }

    private void listenQueue(String queueName, RedisQueueListener redisQueueListener) {
        queueName = buildQueueName(queueName);
        RBlockingQueue<?> blockingQueue = redissonClient.getBlockingQueue(queueName);
        log.info("Redis隊(duì)列監(jiān)聽開啟: {}", queueName);
        while (!shutdownRequested.get() && !redissonClient.isShutdown()) {
            try {
                Object message = blockingQueue.take();
                executorService.submit(() -> redisQueueListener.consume(message));
            } catch (RedissonShutdownException e) {
                log.info("Redis連接關(guān)閉,停止監(jiān)聽隊(duì)列: {}", queueName);
                break;
            } catch (Exception e) {
                log.error("監(jiān)聽隊(duì)列異常: {}", queueName, e);
            }
        }
    }

    public void shutdown() {
        if (executorService != null) {
            executorService.shutdown();
            try {
                if (!executorService.awaitTermination(60, TimeUnit.SECONDS)) {
                    executorService.shutdownNow();
                }
            } catch (InterruptedException ex) {
                executorService.shutdownNow();
                Thread.currentThread().interrupt();
            }
        }
        shutdownRequested.set(true);
        if (redissonClient != null && !redissonClient.isShuttingDown()) {
            redissonClient.shutdown();
        }
    }

    private static class NamedThreadFactory implements ThreadFactory {
        private final AtomicInteger threadNumber = new AtomicInteger(1);
        private final String namePrefix;

        public NamedThreadFactory(String prefix) {
            this.namePrefix = prefix;
        }

        @Override
        public Thread newThread(@NotNull Runnable r) {
            return new Thread(r, namePrefix + "-" + threadNumber.getAndIncrement());
        }
    }

}

RedisQueueListener

/**
 * Redis消息隊(duì)列監(jiān)聽實(shí)現(xiàn)
 *
 * @author 十八
 * @createTime 2024-09-09 22:51
 */
public interface RedisQueueListener<T> {

    /**
     * 隊(duì)列消費(fèi)方法
     *
     * @param content 消息內(nèi)容
     */
    void consume(T content);
}

RedisQueueService

import jakarta.annotation.Resource;
import java.util.concurrent.TimeUnit;
import org.redisson.api.RBlockingQueue;
import org.redisson.api.RDelayedQueue;
import org.redisson.api.RedissonClient;
import org.springframework.stereotype.Component;

/**
 * Redis 消息隊(duì)列服務(wù)
 *
 * @author 十八
 * @createTime 2024-09-09 22:52
 */
@Component
public class RedisQueueService {

    @Resource
    private RedissonClient redissonClient;

    /**
     * 添加隊(duì)列
     *
     * @param queueName 隊(duì)列名稱
     * @param content   消息
     * @param <T>       泛型
     */
    public <T> void send(String queueName, T content) {
        RBlockingQueue<T> blockingQueue = redissonClient.getBlockingQueue(RedisQueueInit.buildQueueName(queueName));
        blockingQueue.add(content);
    }

    /**
     * 添加延遲隊(duì)列
     *
     * @param queueName 隊(duì)列名稱
     * @param content   消息類型
     * @param delay     延遲時(shí)間
     * @param timeUnit  單位
     * @param <T>       泛型
     */
    public <T> void sendDelay(String queueName, T content, long delay, TimeUnit timeUnit) {
        RBlockingQueue<T> blockingFairQueue = redissonClient.getBlockingQueue(RedisQueueInit.buildQueueName(queueName));
        RDelayedQueue<T> delayedQueue = redissonClient.getDelayedQueue(blockingFairQueue);
        delayedQueue.offer(content, delay, timeUnit);
    }

    /**
     * 發(fā)送延遲隊(duì)列消息(單位毫秒)
     *
     * @param queueName 隊(duì)列名稱
     * @param content   消息類型
     * @param delay     延遲時(shí)間
     * @param <T>       泛型
     */
    public <T> void sendDelay(String queueName, T content, long delay) {
        RBlockingQueue<T> blockingFairQueue = redissonClient.getBlockingQueue(RedisQueueInit.buildQueueName(queueName));
        RDelayedQueue<T> delayedQueue = redissonClient.getDelayedQueue(blockingFairQueue);
        delayedQueue.offer(content, delay, TimeUnit.MILLISECONDS);
    }
}

測試

創(chuàng)建監(jiān)聽對象

import cn.yiyanc.infrastructure.redis.annotation.RedisQueue;
import cn.yiyanc.infrastructure.redis.queue.RedisQueueListener;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;

/**
 * @author 十八
 * @createTime 2024-09-10 00:09
 */
@Slf4j
@Component
@RedisQueue("test")
public class TestListener implements RedisQueueListener<String> {
    @Override
    public void invoke(String content) {
        log.info("隊(duì)列消息接收 >>> {}", content);
    }
}

測試用例

import jakarta.annotation.Resource;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

/**
 * @author 十八
 * @createTime 2024-09-10 00:11
 */
@RestController
@RequestMapping("queue")
public class QueueController {

    @Resource
    private RedisQueueService redisQueueService;

    @PostMapping("send")
    public void send(String message) {
        redisQueueService.send("test", message);
        redisQueueService.sendDelay("test", "delay messaege -> " + message, 1000);
    }

}

測試結(jié)果

到此這篇關(guān)于SpringBoot集成Redisson實(shí)現(xiàn)消息隊(duì)列的示例代碼的文章就介紹到這了,更多相關(guān)SpringBoot Redisson消息隊(duì)列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java8中新特性O(shè)ptional、接口中默認(rèn)方法和靜態(tài)方法詳解

    Java8中新特性O(shè)ptional、接口中默認(rèn)方法和靜態(tài)方法詳解

    Java 8 已經(jīng)發(fā)布很久了,很多報(bào)道表明Java 8 是一次重大的版本升級。下面這篇文章主要給大家介紹了關(guān)于Java8中新特性O(shè)ptional、接口中默認(rèn)方法和靜態(tài)方法的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),需要的朋友可以參考下。
    2017-12-12
  • SpringBoot中Tomcat配置的示例代碼

    SpringBoot中Tomcat配置的示例代碼

    本文分享了在SpringBoot項(xiàng)目中配置Tomcat的一些心得和經(jīng)驗(yàn),包括Tomcat版本選擇、調(diào)整配置參數(shù)、自定義連接器、監(jiān)控和日志管理等方面,通過這些配置,可以有效提升應(yīng)用的性能、響應(yīng)速度和并發(fā)處理能力
    2024-11-11
  • java解析json數(shù)組方式

    java解析json數(shù)組方式

    這篇文章主要介紹了java解析json數(shù)組方式,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-06-06
  • maven package 打包報(bào)錯(cuò) Failed to execute goal的解決

    maven package 打包報(bào)錯(cuò) Failed to execute goal的解決

    這篇文章主要介紹了maven package 打包報(bào)錯(cuò) Failed to execute goal的解決,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-11-11
  • IDEA使用Tomcat運(yùn)行web項(xiàng)目教程分享

    IDEA使用Tomcat運(yùn)行web項(xiàng)目教程分享

    在非Spring Boot項(xiàng)目中運(yùn)行Nacos示例,需要手動配置Tomcat容器,本文介紹了如何在IDEA中配置Tomcat,并詳細(xì)解決了配置過程中可能遇到的異常情況,步驟包括修改IDEA項(xiàng)目結(jié)構(gòu)、添加Web模塊、配置Artifacts和Tomcat Server
    2024-10-10
  • spring學(xué)習(xí)之參數(shù)傳遞與檢驗(yàn)詳解

    spring學(xué)習(xí)之參數(shù)傳遞與檢驗(yàn)詳解

    這篇文章主要給大家介紹了關(guān)于spring參數(shù)傳遞與檢驗(yàn)的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作能帶來一定的幫助,需要的朋友們下面跟著小編來一起學(xué)習(xí)學(xué)習(xí)吧。
    2017-07-07
  • 詳解用Kotlin寫一個(gè)基于Spring Boot的RESTful服務(wù)

    詳解用Kotlin寫一個(gè)基于Spring Boot的RESTful服務(wù)

    這篇文章主要介紹了詳解用Kotlin寫一個(gè)基于Spring Boot的RESTful服務(wù) ,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-05-05
  • java實(shí)現(xiàn)對對碰小游戲

    java實(shí)現(xiàn)對對碰小游戲

    這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)對對碰小游戲,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2019-12-12
  • 簡要分析Java的Hibernate框架中的自定義類型

    簡要分析Java的Hibernate框架中的自定義類型

    這篇文章主要介紹了Java的Hibernate框架中的自定義類型,Hibernate是Java的SSH三大web開發(fā)框架之一,需要的朋友可以參考下
    2016-01-01
  • MyBatis-Plus 復(fù)雜查詢Lambda+Wrapper 多條件功能實(shí)現(xiàn)

    MyBatis-Plus 復(fù)雜查詢Lambda+Wrapper 多條件功能實(shí)現(xiàn)

    通過本文的介紹,我們深入了解了MyBatis-Plus中Lambda+Wrapper的強(qiáng)大功能,它不僅極大地簡化了SQL的編寫,還提高了代碼的安全性和可維護(hù)性,感興趣的朋友跟隨小編一起看看吧
    2026-01-01

最新評論

永吉县| 广安市| 墨江| 三明市| 孟津县| 高密市| 铁岭市| 彭水| 岑巩县| 丹江口市| 深圳市| 忻州市| 青冈县| 荃湾区| 翁牛特旗| 布拖县| 运城市| 鹤峰县| 罗甸县| 尼木县| 霍州市| 呼玛县| 呼和浩特市| 齐齐哈尔市| 巴塘县| 渝中区| 临夏县| 宜黄县| 德州市| 同心县| 农安县| 道孚县| 登封市| 东平县| 五河县| 铅山县| 福安市| 大城县| 芦山县| 抚顺市| 新乡市|