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

SpringBoot集成Redis消息隊列的實現(xiàn)示例

 更新時間:2025年05月16日 11:08:47   作者:sjsjsbbsbsn  
本文主要介紹了SpringBoot集成Redis消息隊列的實現(xiàn)示例,包括配置和消費邏輯,RedisStream提供了高吞吐量、順序消費和消費組機(jī)制等優(yōu)勢,具有一定的參考價值,感興趣的可以了解一下

一.Redis Stream 消息隊列模版配置類

/**
 * Redis Stream 消息隊列配置
 */
@Configuration
@RequiredArgsConstructor
public class RedisStreamConfiguration {

    private static final Logger log = LoggerFactory.getLogger(RedisStreamConfiguration.class);
    private final RedisConnectionFactory redisConnectionFactory;
    private final Consumer1 Consumer1;
    private final Consumer2 Consumer2;

    // 定義需要自定義的配置常量
    private static final int BATCH_SIZE = 10; // 每次批量拉取的消息數(shù)量
    private static final Duration POLL_TIMEOUT = Duration.ofSeconds(3); // 拉取消息的阻塞超時時間
    private static final String THREAD_NAME_PREFIX = "your-business"; // 線程名稱前綴
    private static final String GROUP_NAME_1 = "group1"; // 第一個消費者組名稱
    private static final String GROUP_NAME_2 = "group2"; // 第二個消費者組名稱
    private static final String CONSUMER_NAME_1 = "consumer1"; // 第一個消費者名稱
    private static final String CONSUMER_NAME_2 = "consumer2"; // 第二個消費者名稱
    private static final String STREAM_TOPIC_KEY = SHORT_LINK_STATS_STREAM_TOPIC_KEY; // Stream的主題鍵

    @Bean
    public ExecutorService asyncStreamConsumer() {
        log.info("Redis Stream 消息隊列配置線程池");
        AtomicInteger index = new AtomicInteger();
        int processors = Runtime.getRuntime().availableProcessors();

        // 創(chuàng)建一個自定義線程池
        return new ThreadPoolExecutor(
                processors,
                processors + (processors >> 1),
                60,
                TimeUnit.SECONDS,
                new LinkedBlockingQueue<>(),
                runnable -> {
                    Thread thread = new Thread(runnable);
                    thread.setName(THREAD_NAME_PREFIX + "_" + index.incrementAndGet());
                    thread.setDaemon(true);
                    return thread;
                }
        );
    }

    @Bean(initMethod = "start", destroyMethod = "stop")
    public StreamMessageListenerContainer<String, MapRecord<String, String, String>> streamMessageListenerContainer(
            ExecutorService asyncStreamConsumer) {

        // 配置 StreamMessageListenerContainer 容器選項
        StreamMessageListenerContainer.StreamMessageListenerContainerOptions<String, MapRecord<String, String, String>> options =
                StreamMessageListenerContainer.StreamMessageListenerContainerOptions.builder()
                        .batchSize(BATCH_SIZE) // 批量拉取消息數(shù)量
                        .executor(asyncStreamConsumer) // 使用配置好的線程池
                        .pollTimeout(POLL_TIMEOUT) // 拉取消息的超時時間
                        .build();

        // 創(chuàng)建 StreamMessageListenerContainer 實例
        StreamMessageListenerContainer<String, MapRecord<String, String, String>> container =
                StreamMessageListenerContainer.create(redisConnectionFactory, options);

        // 配置第一個消息監(jiān)聽器
        container.receiveAutoAck(
                Consumer.from(GROUP_NAME_1, CONSUMER_NAME_1), // 指定第一個消費者組和消費者名稱
                StreamOffset.create(STREAM_TOPIC_KEY, ReadOffset.lastConsumed()), // 指定主題和偏移量
                Consumer1 // 指定第一個消息處理邏輯
        );

        // 配置第二個消息監(jiān)聽器
        container.receiveAutoAck(
                Consumer.from(GROUP_NAME_2, CONSUMER_NAME_2), // 指定第二個消費者組和消費者名稱
                StreamOffset.create(STREAM_TOPIC_KEY, ReadOffset.lastConsumed()), // 指定主題和偏移量
                Consumer2 // 指定第二個消息處理邏輯
        );

        return container;
    }
}

1. 介紹

RedisStreamConfiguration 是一個用于配置 Redis Stream 消息隊列的 Spring 配置類。它通過 Redis Stream 實現(xiàn)消息的異步處理和多消費者消費,適用于需要高吞吐量、低延遲的業(yè)務(wù)場景。

2. 關(guān)鍵組件和自定義參數(shù)

此類主要配置了 Redis Stream 消息監(jiān)聽容器 StreamMessageListenerContainer,包括線程池配置、消費批次和超時時間等,方便用戶根據(jù)業(yè)務(wù)需求自定義。

核心參數(shù)

  • BATCH_SIZE:定義每次批量拉取的消息數(shù)量。通過設(shè)定合適的批量大小,可以減少消費請求次數(shù),提升處理效率。
  • POLL_TIMEOUT:設(shè)置從 Redis Stream 拉取消息的超時時間。超時控制允許程序在無消息時保持阻塞,等待消息到達(dá)。
  • THREAD_NAME_PREFIX:設(shè)置線程名稱前綴,幫助識別不同業(yè)務(wù)模塊的線程。
  • GROUP_NAME_1 和 GROUP_NAME_2:定義兩個不同的消費者組,適用于同一 Stream 多個消費者并行處理消息的場景。
  • CONSUMER_NAME_1 和 CONSUMER_NAME_2:為每個消費者組指定獨立的消費者名稱,有助于實現(xiàn)消費任務(wù)的分配和管理。

代碼實現(xiàn)

配置了 StreamMessageListenerContainer 來處理 Stream 消息,并分別為兩個消費者組和消費者注冊不同的監(jiān)聽器。

3. 主要方法說明

ExecutorService(線程池配置)

@Bean
public ExecutorService asyncStreamConsumer() { ... }

用于創(chuàng)建一個自定義線程池,為 Redis Stream 的消息消費提供異步執(zhí)行環(huán)境。processors 設(shè)置了核心線程數(shù)為 CPU 核心數(shù),最大線程數(shù)為 processors + (processors >> 1),即核心數(shù)的 1.5 倍。線程命名使用 THREAD_NAME_PREFIX 前綴,方便日志記錄和排查問題。

StreamMessageListenerContainer(消息監(jiān)聽容器)

@Bean(initMethod = "start", destroyMethod = "stop")
public StreamMessageListenerContainer<String, MapRecord<String, String, String>> streamMessageListenerContainer(...) { ... }

該方法創(chuàng)建并配置了 Redis Stream 的監(jiān)聽容器。關(guān)鍵步驟如下:

  • 構(gòu)建容器選項:包括批次大小、線程池、拉取超時時間等參數(shù)。

  • 容器實例化:通過 StreamMessageListenerContainer.create() 創(chuàng)建容器,初始化時自動啟動。

  • 消息監(jiān)聽器配置

    • 為第一個消費者組 GROUP_NAME_1 和消費者 CONSUMER_NAME_1 配置了消息監(jiān)聽器 Consumer1,實現(xiàn)自動確認(rèn)并消費消息。
    • 為第二個消費者組 GROUP_NAME_2 和消費者 CONSUMER_NAME_2 配置了另一組消息監(jiān)聽器 Consumer2,以便多消費者處理。

4. 應(yīng)用場景

此配置適用于 Redis Stream 在大規(guī)模并發(fā)場景下的消息隊列管理。通過靈活配置多個消費者組和消費者,可以實現(xiàn)負(fù)載均衡的多線程消費邏輯。

二.消費者模版

/**
 * 消息隊列消費者
 */
@RequiredArgsConstructor
@Slf4j
@Component
public class ShortLinkStatsSaveConsumer implements StreamListener<String, MapRecord<String, String, String>> {
   
    private final RedissonClient redissonClient;
    private final StringRedisTemplate stringRedisTemplate;
    private final MessageQueueIdempotentHandler messageQueueIdempotentHandler;


    @Override
    public void onMessage(MapRecord<String, String, String> message) {
        String stream = message.getStream();
        RecordId id = message.getId();
        if (!messageQueueIdempotentHandler.isMessageProcessed(id.toString())) {
            // 判斷當(dāng)前的這個消息流程是否執(zhí)行完成
            if (messageQueueIdempotentHandler.isAccomplish(id.toString())) {
                return;
            }
            throw new ServiceException("消息未完成流程,需要消息隊列重試");
        }
        try {
            Map<String, String> producerMap = message.getValue();
         	//你自己的業(yè)務(wù)邏輯
            }
            // 刪除消息
            stringRedisTemplate.opsForStream().delete(Objects.requireNonNull(stream), id.getValue());
        } catch (Throwable ex) {
            messageQueueIdempotentHandler.delMessageProcessed(id.toString());
            log.error("消費異常", ex);
            throw ex;
        }
        //消費完刪除
        messageQueueIdempotentHandler.setAccomplish(id.toString());
    }

   
}

本模板實現(xiàn)了一個 Redis Stream 消息隊列消費者的基礎(chǔ)結(jié)構(gòu)。該模板主要圍繞冪等性檢查、消息解析與處理以及消費狀態(tài)管理三個核心功能,確保消息在高并發(fā)環(huán)境下的安全性與一致性。 

具體的冪等校驗看我另一篇文章

三.生產(chǎn)者模版

/**
 * 短鏈接監(jiān)控狀態(tài)保存消息隊列生產(chǎn)者
 */
@Component
@RequiredArgsConstructor
public class Producer implements MessageQueueProducer{

    private final StringRedisTemplate stringRedisTemplate;


    /**
     * 發(fā)送消息
     */
    public void send(Map<String, String> producerMap) {
        stringRedisTemplate.opsForStream().add(YOUR_KEY, producerMap);
    }
}

注意YOUR_KEY 替換成你自己的即可

四.總結(jié)

1. Redis Stream 消息隊列的優(yōu)勢

Redis Stream 是 Redis 提供的一種強(qiáng)大的消息隊列解決方案,適用于高吞吐量、低延遲的業(yè)務(wù)場景。與傳統(tǒng)的消息隊列系統(tǒng)(如 RabbitMQ 或 Kafka)相比,Redis Stream 在集成與配置方面更加簡單,尤其適合基于 Redis 的應(yīng)用程序。Redis Stream 提供了以下優(yōu)勢:

  • 高吞吐量:支持高并發(fā)和快速消息消費,能夠在瞬間處理大量的消息。
  • 順序消費:保證消息的順序消費,適用于需要順序處理的業(yè)務(wù)場景。
  • 消費組機(jī)制:通過消費組管理消息消費,可以通過多個消費者并行消費,提高處理能力。
  • 持久化與備份:可以將消息存儲在 Redis 中,具備一定的持久化能力,防止數(shù)據(jù)丟失。

2. Redis Stream 配置與應(yīng)用

本文介紹了如何在 Spring Boot 中集成 Redis Stream 消息隊列的配置與消費邏輯,主要包括:

  • 消息消費配置:通過 StreamMessageListenerContainer 實現(xiàn)消息的異步消費。配置了批量拉取的數(shù)量、阻塞超時、線程池等自定義參數(shù),幫助提升系統(tǒng)的并發(fā)處理能力。
  • 多消費者并行處理:通過消費者組(Consumer Group)機(jī)制,實現(xiàn)多消費者并行消費同一個 Stream,提高消息處理的吞吐量和效率。
  • 冪等性與消費確認(rèn):通過 MessageQueueIdempotentHandler 來保證消息的冪等性,避免重復(fù)消費的問題。處理邏輯保證每條消息只會被消費一次,且在消費失敗時能夠適當(dāng)回滾,確保系統(tǒng)的可靠性。

3. 消費者與生產(chǎn)者模板

  • 消費者模板:消費者通過實現(xiàn) StreamListener 接口來處理從 Redis Stream 拉取的消息。為了保證冪等性,消費者首先檢查消息是否已經(jīng)處理過,未完成的消息會被標(biāo)記并重試,確保消息處理的安全性。
  • 生產(chǎn)者模板:生產(chǎn)者通過 StringRedisTemplate 將消息發(fā)送到 Redis Stream。當(dāng)業(yè)務(wù)中有新的消息需要處理時,生產(chǎn)者將消息添加到 Redis Stream 進(jìn)行后續(xù)處理。

4. 應(yīng)用場景

Redis Stream 適用于許多場景,特別是需要高并發(fā)、高吞吐量且保證順序消費的業(yè)務(wù)需求。例如,短鏈接生成與訪問統(tǒng)計、訂單處理、日志收集等業(yè)務(wù)場景,都能通過 Redis Stream 實現(xiàn)高效、可靠的消息隊列。

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

相關(guān)文章

  • Java實現(xiàn)窗體程序顯示日歷表

    Java實現(xiàn)窗體程序顯示日歷表

    這篇文章主要為大家詳細(xì)介紹了Java實現(xiàn)窗體程序顯示日歷表,文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-06-06
  • java(Java?25?LTS)的下載、安裝、配置圖文教程?(IDEA?2025?為例)

    java(Java?25?LTS)的下載、安裝、配置圖文教程?(IDEA?2025?為例)

    在Java開發(fā)中選擇合適的JDK版本至關(guān),重要這篇文章主要介紹了java(Java?25?LTS)的下載、安裝、配置的相關(guān)資料,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2025-11-11
  • PowerJob的DispatchStrategy方法工作流程源碼解讀

    PowerJob的DispatchStrategy方法工作流程源碼解讀

    這篇文章主要為大家介紹了PowerJob的DispatchStrategy方法工作流程源碼解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2024-01-01
  • 在java上使用亞馬遜云儲存方法

    在java上使用亞馬遜云儲存方法

    這篇文章主要介紹了在java上使用亞馬遜云儲存方法,首先寫一個配置類,寫一個controller接口調(diào)用方法存儲文件,本文結(jié)合示例代碼給大家介紹的非常詳細(xì),需要的朋友參考下吧
    2024-01-01
  • 深入理解SpringMVC的參數(shù)綁定與數(shù)據(jù)響應(yīng)機(jī)制

    深入理解SpringMVC的參數(shù)綁定與數(shù)據(jù)響應(yīng)機(jī)制

    本文將深入探討SpringMVC的參數(shù)綁定方式,包括基本類型、對象、集合等類型的綁定方式,以及如何處理參數(shù)校驗和異常。同時,本文還將介紹SpringMVC的數(shù)據(jù)響應(yīng)機(jī)制,包括如何返回JSON、XML等格式的數(shù)據(jù),以及如何處理文件上傳和下載。
    2023-06-06
  • hadoop?切片機(jī)制分析與應(yīng)用

    hadoop?切片機(jī)制分析與應(yīng)用

    切片這個詞對于做過python開發(fā)的同學(xué)一定不陌生,但是與hadoop中的切片有所區(qū)別,hadoop中的切片是為了優(yōu)化hadoop的job在處理過程中MapTask階段的性能達(dá)到最優(yōu)而言
    2022-02-02
  • springboot使用RedisRepository操作數(shù)據(jù)的實現(xiàn)

    springboot使用RedisRepository操作數(shù)據(jù)的實現(xiàn)

    本文主要介紹了springboot使用RedisRepository操作數(shù)據(jù)的實現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2022-05-05
  • SpringBoot實現(xiàn)聯(lián)表查詢的代碼詳解

    SpringBoot實現(xiàn)聯(lián)表查詢的代碼詳解

    這篇文章主要介紹了SpringBoot中如何實現(xiàn)聯(lián)表查詢,文中通過代碼示例和圖文結(jié)合的方式講解的非常詳細(xì),對大家的學(xué)習(xí)或工作有一定的幫助,需要的朋友可以參考下
    2024-05-05
  • 在Spring Boot中實現(xiàn)多環(huán)境配置的方法

    在Spring Boot中實現(xiàn)多環(huán)境配置的方法

    在SpringBoot中,實現(xiàn)多環(huán)境配置是一項重要且常用的功能,它允許開發(fā)者為不同的運行環(huán)境,這種方式簡化了環(huán)境切換的復(fù)雜度,提高了項目的可維護(hù)性和靈活性,本文給大家介紹在Spring Boot中實現(xiàn)多環(huán)境配置的方法,感興趣的朋友跟隨小編一起看看吧
    2024-09-09
  • Java垃圾回收之標(biāo)記壓縮算法詳解

    Java垃圾回收之標(biāo)記壓縮算法詳解

    今天小編就為大家分享一篇關(guān)于Java垃圾回收之標(biāo)記壓縮算法詳解,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2018-10-10

最新評論

柞水县| 唐山市| 晋江市| 嘉荫县| 亚东县| 临洮县| 锡林郭勒盟| 修武县| 丰台区| 平定县| 瑞昌市| 牟定县| 黑龙江省| 汝州市| 鹤峰县| 桂东县| 荥阳市| 凭祥市| 衡山县| 庆元县| 苏州市| 遵化市| 寿宁县| 元朗区| 大埔区| 马鞍山市| 尉犁县| 平乐县| 南宫市| 安义县| 察雅县| 丽水市| 阳新县| 兰州市| 庄河市| 胶州市| 台北县| 博野县| 洪泽县| 福建省| 南乐县|