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

SpringBoot+RocketMQ實現(xiàn)延遲消息的示例代碼

 更新時間:2025年10月24日 11:08:29   作者:匆匆忙忙游刃有余  
本文主要介紹了SpringBoot+RocketMQ實現(xiàn)延遲消息案例詳解,包括基于延遲級別和基于具體時間兩種方式的完整實現(xiàn),文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

下面將詳細介紹如何在SpringBoot中使用RocketMQ實現(xiàn)延遲消息,包括基于延遲級別和基于具體時間兩種方式的完整實現(xiàn)。

一、延遲消息概述

RocketMQ提供了兩種類型的延遲消息機制:

  1. 延遲消息:消息發(fā)送后延遲指定的時間長度再被消費
  2. 定時消息:消息在指定的具體時間點被消費

這兩種機制在訂單超時取消、會議提醒、定時任務調(diào)度等場景中有廣泛應用。

二、環(huán)境準備

1. 添加Maven依賴

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.3</version>
</dependency>

2. 配置文件設置

application.yml中配置RocketMQ連接信息:

rocketmq:
  name-server: localhost:9876
  producer:
    group: delay-message-producer-group

三、延遲級別機制實現(xiàn)

1. 默認延遲級別

RocketMQ默認提供18個延遲級別,定義在MessageStoreConfig類中:

messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"

對應關系:

  • level=1: 延遲1秒
  • level=2: 延遲5秒
  • level=3: 延遲10秒
  • level=4: 延遲30秒
  • level=5: 延遲1分鐘
  • level=6: 延遲2分鐘
  • ...以此類推
  • level=18: 延遲2小時

2. 基于延遲級別的生產(chǎn)者實現(xiàn)

import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

@Component
public class DelayLevelProducer {
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 發(fā)送基于延遲級別的消息
     * @param topic 主題
     * @param tag 標簽
     * @param message 消息內(nèi)容
     * @param delayLevel 延遲級別(1-18)
     */
    public void sendMessageByDelayLevel(String topic, String tag, String message, int delayLevel) {
        // 創(chuàng)建消息
        Message<String> springMessage = MessageBuilder.withPayload(message).build();
        
        // 發(fā)送延遲消息
        SendResult sendResult = rocketMQTemplate.syncSend(
            topic + ":" + tag, 
            springMessage, 
            3000, // 超時時間
            delayLevel // 延遲級別
        );
        
        System.out.println("延遲級別消息發(fā)送成功: " + sendResult);
    }
    
    /**
     * 發(fā)送訂單超時取消消息(延遲15分鐘)
     */
    public void sendOrderTimeoutMessage(String orderId) {
        String message = "訂單超時取消: " + orderId;
        // 15分鐘對應level=14(根據(jù)默認配置)
        sendMessageByDelayLevel("OrderTopic", "Timeout", message, 14);
    }
}

四、基于具體時間的延遲消息實現(xiàn)

1. 定時消息生產(chǎn)者

import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

import java.util.Date;

@Component
public class ScheduledMessageProducer {
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 發(fā)送延遲指定毫秒數(shù)的消息
     */
    public void sendMessageWithDelayMs(String topic, String message, long delayMs) {
        // 計算投遞時間
        long deliverTimeMs = System.currentTimeMillis() + delayMs;
        
        // 創(chuàng)建消息并設置投遞時間
        Message<String> springMessage = MessageBuilder.withPayload(message)
            .setHeader(MessageConst.PROPERTY_DELAY_TIME_MS, String.valueOf(delayMs))
            .setHeader(MessageConst.PROPERTY_TIMER_DELIVER_MS, String.valueOf(deliverTimeMs))
            .build();
        
        SendResult sendResult = rocketMQTemplate.syncSend(topic, springMessage);
        System.out.println("延遲毫秒消息發(fā)送成功: " + sendResult);
    }
    
    /**
     * 發(fā)送指定時間點投遞的消息
     */
    public void sendMessageAtTime(String topic, String message, Date deliverTime) {
        long deliverTimeMs = deliverTime.getTime();
        
        // 創(chuàng)建消息并設置投遞時間
        Message<String> springMessage = MessageBuilder.withPayload(message)
            .setHeader(MessageConst.PROPERTY_TIMER_DELIVER_MS, String.valueOf(deliverTimeMs))
            .build();
        
        SendResult sendResult = rocketMQTemplate.syncSend(topic, springMessage);
        System.out.println("定時投遞消息發(fā)送成功: " + sendResult);
    }
    
    /**
     * 發(fā)送10秒后投遞的消息
     */
    public void sendTenSecondsLaterMessage(String topic, String message) {
        sendMessageWithDelayMs(topic, message, 10000L);
    }
}

五、消費者實現(xiàn)

延遲消息的消費者與普通消息消費者相同,無需特殊配置:

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;

import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;

@Component
@RocketMQMessageListener(
    topic = "OrderTopic",
    consumerGroup = "delay-message-consumer-group",
    selectorExpression = "Timeout"
)
public class OrderTimeoutConsumer implements RocketMQListener<String> {
    
    @Override
    public void onMessage(String message) {
        String now = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
        System.out.println("[" + now + "] 接收到訂單超時消息: " + message);
        
        // 處理訂單取消邏輯
        processOrderCancellation(message);
    }
    
    private void processOrderCancellation(String message) {
        // 提取訂單ID
        String orderId = message.substring(message.indexOf(":") + 2);
        System.out.println("執(zhí)行訂單取消操作,訂單ID: " + orderId);
        // 這里可以調(diào)用訂單服務進行取消操作
    }
}

六、Controller層實現(xiàn)

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.format.annotation.DateTimeFormat;
import org.springframework.web.bind.annotation.*;

import java.util.Date;

@RestController
@RequestMapping("/api/delay")
public class DelayMessageController {
    
    @Autowired
    private DelayLevelProducer delayLevelProducer;
    
    @Autowired
    private ScheduledMessageProducer scheduledMessageProducer;
    
    /**
     * 發(fā)送基于延遲級別的消息
     */
    @PostMapping("/level")
    public String sendByDelayLevel(
            @RequestParam String topic,
            @RequestParam String tag,
            @RequestParam String message,
            @RequestParam(defaultValue = "3") int delayLevel) {
        
        delayLevelProducer.sendMessageByDelayLevel(topic, tag, message, delayLevel);
        return "延遲級別消息發(fā)送成功,延遲級別: " + delayLevel;
    }
    
    /**
     * 發(fā)送訂單超時取消消息
     */
    @PostMapping("/order/timeout")
    public String sendOrderTimeout(@RequestParam String orderId) {
        delayLevelProducer.sendOrderTimeoutMessage(orderId);
        return "訂單超時取消消息已發(fā)送,訂單ID: " + orderId;
    }
    
    /**
     * 發(fā)送延遲指定毫秒的消息
     */
    @PostMapping("/milliseconds")
    public String sendByDelayMs(
            @RequestParam String topic,
            @RequestParam String message,
            @RequestParam long delayMs) {
        
        scheduledMessageProducer.sendMessageWithDelayMs(topic, message, delayMs);
        return "延遲毫秒消息發(fā)送成功,延遲: " + delayMs + "ms";
    }
    
    /**
     * 發(fā)送指定時間點的消息
     */
    @PostMapping("/scheduled")
    public String sendScheduled(
            @RequestParam String topic,
            @RequestParam String message,
            @RequestParam @DateTimeFormat(pattern = "yyyy-MM-dd HH:mm:ss") Date deliverTime) {
        
        scheduledMessageProducer.sendMessageAtTime(topic, message, deliverTime);
        return "定時消息發(fā)送成功,投遞時間: " + deliverTime;
    }
}

七、自定義延遲級別配置

在Broker的配置文件中可以自定義延遲級別:

# 在broker.conf文件中添加
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h 3h 4h 5h

重啟Broker使其生效。注意,修改延遲級別后,所有使用延遲級別的消息都會使用新的配置。

八、兩種實現(xiàn)方式對比

特性基于延遲級別基于具體時間
靈活性較低,只能使用預定義級別高,可以精確到毫秒
適用版本全版本支持RocketMQ 5.x及以上版本完整支持
使用場景固定延遲時間的場景需要精確控制投遞時間的場景
配置復雜度簡單,無需額外配置可能需要在Broker端開啟相關功能

九、使用注意事項

  1. 延遲精度

    • 延遲消息的投遞時間不是完全精確的,有一定誤差
    • 在高并發(fā)場景下,誤差可能會增大
  2. 版本兼容性

    • 基于具體時間的延遲消息在RocketMQ 5.x版本支持更完善
    • 在低版本中可能需要使用延遲級別機制
  3. 性能考慮

    • 大量延遲消息可能會增加Broker的負擔
    • 對于長時間延遲的消息,考慮使用其他方案(如定時任務+消息隊列組合)
  4. 消息可靠性

    • 延遲消息同樣支持持久化,確保Broker重啟后不會丟失
    • 建議開啟消息確認機制確保消息可靠投遞

十、測試示例

  1. 發(fā)送訂單超時取消消息(延遲15分鐘):

    POST /api/delay/order/timeout?orderId=ORDER123456
    
  2. 發(fā)送10秒后投遞的消息:

    POST /api/delay/milliseconds?topic=TestTopic&message=HelloDelay&delayMs=10000
    
  3. 發(fā)送指定時間點的消息:

    POST /api/delay/scheduled?topic=TestTopic&message=HelloScheduled&deliverTime=2024-12-25%2000:00:00
    

通過以上配置和代碼,您可以在SpringBoot項目中輕松實現(xiàn)基于RocketMQ的延遲消息功能,滿足各種定時任務和延遲處理的業(yè)務需求。

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

相關文章

  • 基于@RequestParam注解之Spring MVC參數(shù)綁定的利器

    基于@RequestParam注解之Spring MVC參數(shù)綁定的利器

    這篇文章主要介紹了基于@RequestParam注解之Spring MVC參數(shù)綁定的利器,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-03-03
  • debug模式遲遲不能啟動問題及解決

    debug模式遲遲不能啟動問題及解決

    在使用Debug模式進行代碼測試時,由于設置了過多的斷點,導致程序加載緩慢甚至無法啟動,解決此問題的方法是取消不必要的斷點,通過IDE的斷點管理功能,檢查并移除問題斷點,從而優(yōu)化調(diào)試效率,分享此經(jīng)驗希望能幫助遇到相同問題的開發(fā)者
    2022-11-11
  • Java中各類日期和時間轉(zhuǎn)換超詳析總結(jié)(Date和LocalDateTime相互轉(zhuǎn)換等)

    Java中各類日期和時間轉(zhuǎn)換超詳析總結(jié)(Date和LocalDateTime相互轉(zhuǎn)換等)

    這篇文章主要介紹了Java中日期和時間處理的幾個階段,包括java.util.Date、java.sql.Date、java.sql.Time、java.sql.Timestamp、java.util.Calendar和java.util.GregorianCalendar等類,文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2025-01-01
  • 論Java Web應用中調(diào)優(yōu)線程池的重要性

    論Java Web應用中調(diào)優(yōu)線程池的重要性

    這篇文章主要論述Java Web應用中調(diào)優(yōu)線程池的重要性,通過了解應用的需求,組合最大線程數(shù)和平均響應時間,得出一個合適的線程池配置
    2016-04-04
  • springboot2中設置@ApiImplicitParam的dataType不起作用的解決

    springboot2中設置@ApiImplicitParam的dataType不起作用的解決

    本文主要介紹了在SpringBoot2中使用Swagger時,@ApiImplicitParam的dataType屬性不起作用的問題及其解決方法,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2026-03-03
  • 解決springboot 連接 mysql 時報錯 using password: NO的方案

    解決springboot 連接 mysql 時報錯 using password: NO的方案

    在本篇文章里小編給大家整理了關于解決springboot 連接 mysql 時報錯 using password: NO的方案,有需要的朋友們可以學習下。
    2020-01-01
  • SpringBoot自定義注解及AOP的開發(fā)和使用詳解

    SpringBoot自定義注解及AOP的開發(fā)和使用詳解

    在公司項目中,如果需要做一些公共的功能,如日志等,最好的方式是使用自定義注解,自定義注解可以實現(xiàn)我們對想要添加日志的方法上添加,這篇文章基于日志功能來講講自定義注解應該如何開發(fā)和使用,需要的朋友可以參考下
    2023-08-08
  • Java中輸入輸出方式詳細講解

    Java中輸入輸出方式詳細講解

    這篇文章主要給大家介紹了關于Java中輸入輸出方式的相關資料,Java輸入輸出是指使用java提供的一些類和方法來實現(xiàn)數(shù)據(jù)的輸入和輸出,文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2023-09-09
  • Java Quartz觸發(fā)器CronTriggerBean配置用法詳解

    Java Quartz觸發(fā)器CronTriggerBean配置用法詳解

    這篇文章主要介紹了Java Quartz觸發(fā)器CronTriggerBean配置用法詳解,本篇文章通過簡要的案例,講解了該項技術(shù)的了解與使用,以下就是詳細內(nèi)容,需要的朋友可以參考下
    2021-08-08
  • Java編程BigDecimal用法實例分享

    Java編程BigDecimal用法實例分享

    這篇文章主要介紹了Java編程BigDecimal用法實例分享,具有一定參考價值,需要的朋友可以了解下。
    2017-11-11

最新評論

新野县| 军事| 府谷县| 沙洋县| 洞口县| 广安市| 泽普县| 昆明市| 博兴县| 武义县| 辽阳县| 鹿泉市| 施秉县| 如皋市| 汶上县| 清河县| 兴仁县| 孙吴县| 达日县| 石棉县| 呼伦贝尔市| 平安县| 左贡县| 讷河市| 防城港市| 石家庄市| 南丹县| 富阳市| 金秀| 灵山县| 周宁县| 临泽县| 梁平县| 体育| 黄冈市| 安龙县| 灵台县| 出国| 盘锦市| 霍林郭勒市| 丰原市|