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

spring rocketmq集成方案

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

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

一、前置條件

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

二、核心依賴引入

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

<!-- Spring Boot 基礎(chǔ)依賴 -->
<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>
<!-- 可選:測(cè)試依賴 -->
<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 # 應(yīng)用名稱
# RocketMQ 核心配置
rocketmq:
  name-server: 127.0.0.1:9876 # NameServer 地址(集群用逗號(hào)分隔)
  producer:
    group: demo-producer-group # 生產(chǎn)者組(必填,標(biāo)識(shí)同一類生產(chǎn)者)
    send-message-timeout: 3000 # 發(fā)送超時(shí)時(shí)間,默認(rèn)3000ms
    compress-message-body-threshold: 4096 # 消息壓縮閾值,默認(rèn)4096字節(jié)
    max-message-size: 4194304 # 最大消息大小,默認(rèn)4MB
    retry-times-when-send-failed: 2 # 同步發(fā)送失敗重試次數(shù)
    retry-times-when-send-async-failed: 2 # 異步發(fā)送失敗重試次數(shù)

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

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

適用于需要確認(rèn)發(fā)送結(jié)果的場(chǎng)景(如訂單創(chuàng)建、庫(kù)存扣減):

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ā)送消息(阻塞等待結(jié)果)
     * @param topic 消息主題(必填,需提前創(chuàng)建)
     * @param message 消息內(nèi)容
     * @return 發(fā)送結(jié)果
     */
    public SendResult sendSyncMessage(String topic, String message) {
        try {
            // 發(fā)送格式:"topic:tag"(tag可選,用于消息過(guò)濾)
            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è)務(wù)異常處理(如重試、記錄日志、告警)
            throw new RuntimeException("消息發(fā)送失敗", e);
        }
    }
    /**
     * 異步發(fā)送消息(非阻塞,回調(diào)通知結(jié)果)
     * @param topic 主題
     * @param message 消息內(nèi)容
     */
    public void sendAsyncMessage(String topic, String message) {
        rocketMQTemplate.asyncSend(
            topic + ":demoTag",
            message,
            // 發(fā)送成功回調(diào)
            sendResult -> System.out.println("異步發(fā)送成功,消息ID:" + sendResult.getMsgId()),
            // 發(fā)送失敗回調(diào)
            throwable -> System.err.println("異步發(fā)送失?。? + throwable.getMessage())
        );
    }
    /**
     * 單向發(fā)送消息(無(wú)回調(diào),適用于日志、埋點(diǎn)等不關(guān)心結(jié)果的場(chǎng)景)
     * @param topic 主題
     * @param message 消息內(nèi)容
     */
    public void sendOneWayMessage(String topic, String message) {
        rocketMQTemplate.sendOneWay(topic + ":demoTag", message);
        System.out.println("單向消息發(fā)送請(qǐng)求已提交");
    }
}

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

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

// 自定義消息實(shí)體(實(shí)現(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);
}

五、消費(fèi)者實(shí)現(xiàn)(2 種消費(fèi)模式)

1. 普通消費(fèi)(默認(rèn)集群模式)

適用于多實(shí)例負(fù)載均衡消費(fèi)(同一組消費(fèi)者分?jǐn)傁ⅲ?/p>

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

2. 消費(fèi)自定義對(duì)象消息

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

@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è)務(wù)邏輯
    }
}

六、測(cè)試驗(yàn)證

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

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秒,確保消費(fèi)者能接收到消息
        try {
            Thread.sleep(3000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

七、關(guān)鍵注意事項(xiàng)

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

總結(jié)

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

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

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

相關(guān)文章

最新評(píng)論

广饶县| 新宁县| 故城县| 灵璧县| 始兴县| 昌都县| 梨树县| 揭阳市| 孙吴县| 穆棱市| 铅山县| 东宁县| 定襄县| 佛教| 奎屯市| 康乐县| 丹阳市| 阿尔山市| 滕州市| 南昌县| 定州市| 仪征市| 台南市| 奈曼旗| 清涧县| 沂源县| 武清区| 晴隆县| 浦江县| 长子县| 油尖旺区| 探索| 西华县| 郯城县| 永嘉县| 仁怀市| 乡宁县| 临颍县| 乾安县| 潍坊市| 双流县|