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

spring rocketmq集成方案

 更新時間:2026年03月13日 16:51:05   作者:金坷拉  
本文詳細介紹了如何在Spring Boot項目中集成RocketMQ 5.x,包括前置條件、核心依賴、配置、生產(chǎn)者和消費者實現(xiàn)、測試驗證以及關鍵注意事項,感興趣的朋友跟隨小編一起看看吧

在 Spring 項目中集成 RocketMQ 是非常常見的消息隊列應用場景,我會以 Spring Boot + RocketMQ 5.x(當前主流版本)為例,提供完整、可直接運行的集成方案,包括生產(chǎn)者、消費者的核心代碼和配置說明。

一、前置條件

  1. 已安裝并啟動 RocketMQ(NameServer + Broker),默認端口:NameServer 9876
  2. Spring Boot 版本建議:2.7.x3.x(兼容 RocketMQ 官方 starter)
  3. 開發(fā)環(huán)境:JDK 8+

二、核心依賴引入

pom.xml 中添加 RocketMQ 與 Spring Boot 集成的官方 starter:

<!-- Spring Boot 基礎依賴 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter</artifactId>
</dependency>
<!-- RocketMQ Spring Boot Starter(官方推薦) -->
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.3</version> <!-- 適配 RocketMQ 5.x,兼容 Spring Boot 2/3 -->
</dependency>
<!-- 可選:測試依賴 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-test</artifactId>
    <scope>test</scope>
</dependency>

三、核心配置(application.yml)

resources 目錄下配置 RocketMQ 連接信息:

spring:
  application:
    name: rocketmq-demo # 應用名稱
# RocketMQ 核心配置
rocketmq:
  name-server: 127.0.0.1:9876 # NameServer 地址(集群用逗號分隔)
  producer:
    group: demo-producer-group # 生產(chǎn)者組(必填,標識同一類生產(chǎn)者)
    send-message-timeout: 3000 # 發(fā)送超時時間,默認3000ms
    compress-message-body-threshold: 4096 # 消息壓縮閾值,默認4096字節(jié)
    max-message-size: 4194304 # 最大消息大小,默認4MB
    retry-times-when-send-failed: 2 # 同步發(fā)送失敗重試次數(shù)
    retry-times-when-send-async-failed: 2 # 異步發(fā)送失敗重試次數(shù)

四、生產(chǎn)者實現(xiàn)(3 種發(fā)送方式)

1. 基礎同步發(fā)送(最常用)

適用于需要確認發(fā)送結果的場景(如訂單創(chuàng)建、庫存扣減):

import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
@Component
public class RocketMQProducer {
    // 注入官方封裝的 RocketMQ 模板(類似 RabbitTemplate)
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    /**
     * 同步發(fā)送消息(阻塞等待結果)
     * @param topic 消息主題(必填,需提前創(chuàng)建)
     * @param message 消息內容
     * @return 發(fā)送結果
     */
    public SendResult sendSyncMessage(String topic, String message) {
        try {
            // 發(fā)送格式:"topic:tag"(tag可選,用于消息過濾)
            SendResult sendResult = rocketMQTemplate.syncSend(topic + ":demoTag", message);
            System.out.println("同步發(fā)送成功,消息ID:" + sendResult.getMsgId());
            return sendResult;
        } catch (Exception e) {
            System.err.println("同步發(fā)送失?。? + e.getMessage());
            // 業(yè)務異常處理(如重試、記錄日志、告警)
            throw new RuntimeException("消息發(fā)送失敗", e);
        }
    }
    /**
     * 異步發(fā)送消息(非阻塞,回調通知結果)
     * @param topic 主題
     * @param message 消息內容
     */
    public void sendAsyncMessage(String topic, String message) {
        rocketMQTemplate.asyncSend(
            topic + ":demoTag",
            message,
            // 發(fā)送成功回調
            sendResult -> System.out.println("異步發(fā)送成功,消息ID:" + sendResult.getMsgId()),
            // 發(fā)送失敗回調
            throwable -> System.err.println("異步發(fā)送失?。? + throwable.getMessage())
        );
    }
    /**
     * 單向發(fā)送消息(無回調,適用于日志、埋點等不關心結果的場景)
     * @param topic 主題
     * @param message 消息內容
     */
    public void sendOneWayMessage(String topic, String message) {
        rocketMQTemplate.sendOneWay(topic + ":demoTag", message);
        System.out.println("單向消息發(fā)送請求已提交");
    }
}

2. 發(fā)送自定義對象消息

如果需要發(fā)送 Java 對象(而非字符串),只需保證對象可序列化:

// 自定義消息實體(實現(xiàn) Serializable)
public class OrderMessage implements Serializable {
    private Long orderId;
    private String orderNo;
    private BigDecimal amount;
    // 省略 getter/setter/toString
}
// 生產(chǎn)者中新增方法
public SendResult sendObjectMessage(String topic, OrderMessage orderMessage) {
    return rocketMQTemplate.syncSend(topic + ":orderTag", orderMessage);
}

五、消費者實現(xiàn)(2 種消費模式)

1. 普通消費(默認集群模式)

適用于多實例負載均衡消費(同一組消費者分攤消息):

import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
/**
 * RocketMQ 消費者
 * - topic:訂閱的主題(需與生產(chǎn)者一致)
 * - consumerGroup:消費者組(必填,同一組消費同一主題)
 * - messageModel:消費模式(CLUSTERING 集群模式,BROADCASTING 廣播模式)
 * - consumeMode:消費方式(CONCURRENTLY 并發(fā)消費,ORDERLY 順序消費)
 */
@Component
@RocketMQMessageListener(
    topic = "demo_topic", // 訂閱主題
    consumerGroup = "demo-consumer-group", // 消費者組
    messageModel = MessageModel.CLUSTERING, // 集群模式(默認)
    consumeMode = ConsumeMode.CONCURRENTLY // 并發(fā)消費(默認)
)
public class RocketMQConsumer implements RocketMQListener<String> {
    /**
     * 消息消費邏輯(接收到消息時觸發(fā))
     * @param message 消息內容(與生產(chǎn)者發(fā)送類型一致)
     */
    @Override
    public void onMessage(String message) {
        try {
            // 核心業(yè)務邏輯:如解析消息、處理訂單、更新庫存等
            System.out.println("接收到消息:" + message);
            // 消費成功無需返回,拋出異常則會觸發(fā)重試
        } catch (Exception e) {
            System.err.println("消息消費失?。? + e.getMessage());
            // 異常拋出后,RocketMQ 會自動重試(默認最多16次)
            throw new RuntimeException("消費失敗", e);
        }
    }
}

2. 消費自定義對象消息

如果生產(chǎn)者發(fā)送的是自定義對象,消費者需指定泛型為對應類型:

@Component
@RocketMQMessageListener(
    topic = "demo_topic",
    consumerGroup = "order-consumer-group"
)
public class OrderMessageConsumer implements RocketMQListener<OrderMessage> {
    @Override
    public void onMessage(OrderMessage orderMessage) {
        System.out.println("接收到訂單消息:" + orderMessage);
        // 處理訂單業(yè)務邏輯
    }
}

六、測試驗證

編寫測試類,驗證生產(chǎn)者發(fā)送、消費者接收是否正常:

import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
@SpringBootTest
public class RocketMQDemoTest {
    @Autowired
    private RocketMQProducer rocketMQProducer;
    @Test
    public void testSendSyncMessage() {
        // 發(fā)送消息到 demo_topic 主題
        rocketMQProducer.sendSyncMessage("demo_topic", "Hello RocketMQ + Spring Boot!");
        // 暫停3秒,確保消費者能接收到消息
        try {
            Thread.sleep(3000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

七、關鍵注意事項

  1. 主題 / 組命名規(guī)范:避免特殊字符,建議用 業(yè)務_模塊_主題 格式(如 order_pay_topic)。
  2. 重試機制:消費失敗默認重試 16 次,可通過 maxReconsumeTimes 配置重試次數(shù)。
  3. 消息持久化:RocketMQ 默認持久化消息,即使消費者宕機,重啟后仍能消費未處理的消息。
  4. 順序消費:如需保證消息順序,需將 consumeMode 設為 ORDERLY,且生產(chǎn)者發(fā)送時指定同一消息隊列。
  5. 異常處理:生產(chǎn)環(huán)境建議對接告警(如釘釘、短信),避免消費失敗無感知。

總結

  1. Spring Boot 集成 RocketMQ 的核心是引入官方 rocketmq-spring-boot-starter,配置 NameServer 地址和生產(chǎn) / 消費組。
  2. 生產(chǎn)者通過 RocketMQTemplate 實現(xiàn)同步 / 異步 / 單向發(fā)送,支持字符串和自定義對象消息。
  3. 消費者通過 @RocketMQMessageListener 注解聲明訂閱關系,實現(xiàn) RocketMQListener 接口處理消息邏輯,默認集群模式并發(fā)消費。

核心關鍵點:主題與消費組必須配置正確,消費失敗拋出異常會觸發(fā)自動重試,生產(chǎn)環(huán)境需做好異常監(jiān)控和重試次數(shù)限制。

到此這篇關于spring rocketmq集成的文章就介紹到這了,更多相關spring rocketmq集成內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • maven 環(huán)境變量的配置詳解

    maven 環(huán)境變量的配置詳解

    這篇文章主要介紹了maven 環(huán)境變量的配置詳解,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2020-09-09
  • Spring Boot 菜單刪除實現(xiàn)代碼與事務管理

    Spring Boot 菜單刪除實現(xiàn)代碼與事務管理

    本文將詳細介紹Spring Boot環(huán)境下菜單刪除功能的實現(xiàn)邏輯,包括關聯(lián)數(shù)據(jù)處理、事務控制和異常處理等關鍵環(huán)節(jié),強調需處理多級嵌套、角色關聯(lián)及數(shù)據(jù)一致性,感興趣的朋友跟隨小編一起看看吧
    2025-08-08
  • 快速排序和分治排序介紹

    快速排序和分治排序介紹

    這篇文章主要介紹了快速排序和分治排序,需要的朋友可以參考下
    2015-04-04
  • 關于Hibernate的一些學習心得總結

    關于Hibernate的一些學習心得總結

    Hibernate是一個優(yōu)秀的Java 持久化層解決方案,是當今主流的對象—關系映射(ORM)工具
    2013-07-07
  • 詳解Java是如何通過接口來創(chuàng)建代理并進行http請求

    詳解Java是如何通過接口來創(chuàng)建代理并進行http請求

    今天給大家?guī)淼闹R是關于Java的,文章圍繞Java是如何通過接口來創(chuàng)建代理并進行http請求展開,文中有非常詳細的介紹及代碼示例,需要的朋友可以參考下
    2021-06-06
  • Spring Cloud實現(xiàn)5分鐘級區(qū)域切換的操作方法

    Spring Cloud實現(xiàn)5分鐘級區(qū)域切換的操作方法

    Spring Cloud 2023.x通過智能路由預熱、多活數(shù)據(jù)同步和自動化流量切換,實現(xiàn)5分鐘內完成跨區(qū)域故障轉移,本文以某電商平臺從AWS亞太切換至阿里云華東的實戰(zhàn)為例,詳解關鍵技術路徑,需要的朋友可以參考下
    2025-04-04
  • springboot使用注解實現(xiàn)鑒權功能

    springboot使用注解實現(xiàn)鑒權功能

    這篇文章主要介紹了springboot使用注解實現(xiàn)鑒權功能,本文通過實例代碼給大家介紹的非常詳細,感興趣的朋友跟隨小編一起看看吧
    2024-12-12
  • springboot+thymeleaf國際化之LocaleResolver接口的示例

    springboot+thymeleaf國際化之LocaleResolver接口的示例

    本篇文章主要介紹了springboot+thymeleaf國際化之LocaleResolver的示例 ,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-11-11
  • Java實現(xiàn)map轉換成json的方法詳解

    Java實現(xiàn)map轉換成json的方法詳解

    這篇文章主要為大家詳細介紹了Java語言實現(xiàn)map轉換成json的幾種方法,文中的示例代碼講解詳細,對我們學習Java有一定幫助,需要的可以參考一下
    2022-05-05
  • JAVA學習筆記:注釋、變量的聲明和定義操作實例分析

    JAVA學習筆記:注釋、變量的聲明和定義操作實例分析

    這篇文章主要介紹了JAVA學習筆記:注釋、變量的聲明和定義操作,結合實例形式分析了Java注釋、變量的聲明和定義相關原理、實現(xiàn)方法及操作注意事項,需要的朋友可以參考下
    2020-04-04

最新評論

丹阳市| 海伦市| 准格尔旗| 大邑县| 大方县| 永丰县| 乳源| 九寨沟县| 蕉岭县| 武鸣县| 徐州市| 奎屯市| 桓台县| 遂宁市| 穆棱市| 婺源县| 阳江市| 西乌珠穆沁旗| 余姚市| 湘潭市| 成武县| 隆回县| 卓尼县| 海阳市| 诏安县| 闸北区| 中西区| 辽宁省| 大同市| 彭州市| 惠来县| 宝应县| 太保市| 翼城县| 祥云县| 陆丰市| 搜索| 江永县| 武陟县| 盱眙县| 柳江县|