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

SpringBoot實現(xiàn)發(fā)布/訂閱廣播消息的示例代碼

 更新時間:2026年05月28日 08:42:16   作者:希望永不加班  
在 RabbitMQ 的五大工作模式中,發(fā)布/訂閱(Publish/Subscribe)廣播模式是分布式系統(tǒng)中非常核心的通信方式,本文給大家介紹了SpringBoot實現(xiàn)發(fā)布/訂閱廣播消息的示例代碼,需要的朋友可以參考下

引言

在 RabbitMQ 的五大工作模式中,發(fā)布/訂閱(Publish/Subscribe)廣播模式是分布式系統(tǒng)中非常核心的通信方式。

我們?nèi)粘J褂玫钠胀c對點隊列,一條消息只會被一個消費者消費(競爭消費);而廣播模式可以實現(xiàn)一條消息、多服務、多消費者同時接收,完美實現(xiàn)一對多通知。

像 緩存刷新、配置更新、全局通知、多節(jié)點日志同步、服務狀態(tài)廣播 等場景,全部依賴 Fanout 廣播模式。

一、什么是 MQ 廣播(發(fā)布/訂閱)模式?

1. 核心定義

廣播模式基于 FanoutExchange(扇形交換機) 實現(xiàn),核心邏輯:

生產(chǎn)者發(fā)送一條消息到 Fanout 交換機,所有綁定該交換機的隊列,都會完整收到這條消息。

不管路由鍵是什么、不管隊列名稱,只要完成綁定,就會無條件廣播投遞。

2. 核心特性

  • 無視routingKey,路由鍵傳空、傳任意值都不生效
  • 純廣播、全量投遞、一對多分發(fā)
  • 每條消息獨立進入每一個綁定隊列
  • 天然支持多服務、多節(jié)點同步通知
  • 無匹配規(guī)則,綁定即接收

3. 適用業(yè)務場景

  • 分布式緩存全局刷新(多節(jié)點統(tǒng)一清空緩存)
  • 系統(tǒng)配置動態(tài)推送、熱更新
  • 全站公告、全局消息推送
  • 微服務多節(jié)點日志采集、鏈路追蹤
  • 服務上下線、狀態(tài)同步廣播
  • 多端消息同步(PC/APP/小程序)

二、四大交換機模式核心對比

交換機類型匹配規(guī)則消費模式核心場景
Direct(直連)完全匹配 routingKey點對點競爭消費訂單、支付、任務處理
Topic(主題)通配符模糊匹配選擇性多消費日志分級、消息訂閱
Fanout(廣播)無視路由鍵,全部投遞全員訂閱消費緩存刷新、全局通知
Headers匹配消息頭參數(shù)自定義匹配極少使用

三、關鍵認知誤區(qū)

1:同一個隊列多消費者可以實現(xiàn)廣播

絕對錯誤!

同一個隊列下的多個消費者,默認是競爭消費,一條消息只會被一個消費者消費。

廣播必備條件:每個消費者對應一個獨立隊列,全部綁定同一個 Fanout 交換機。

2:Fanout 交換機需要配置路由鍵

Fanout 交換機底層邏輯直接忽略 routingKey,無論發(fā)送時傳什么值,都不會影響廣播效果。

3:廣播消息天然可靠、不會丟失

默認非持久化、自動ACK 場景下,廣播消息極易丟失,生產(chǎn)必須做持久化+手動ACK。

四、SpringBoot 完整實現(xiàn)

1. 基礎依賴

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

2. 生產(chǎn)級配置文件

spring:
  rabbitmq:
    host: 127.0.0.1
    port: 5672
    username: guest
    password: guest
    virtual-host: /
    listener:
      simple:
        # 手動ACK 保證廣播消息不丟
        acknowledge-mode: manual
        # 限制預取數(shù),防止單節(jié)點消息堆積
        prefetch: 5
        # 開啟消費重試
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000

3. 廣播交換機、隊列、綁定配置類

import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.FanoutExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class FanoutBroadcastConfig {
    // 廣播交換機名稱
    public static final String FANOUT_EXCHANGE = "system.fanout.broadcast.exchange";
    // 三個獨立消費者隊列
    public static final String QUEUE_CACHE_REFRESH = "queue.cache.refresh";
    public static final String QUEUE_NOTICE = "queue.system.notice";
    public static final String QUEUE_LOG = "queue.log.collect";
    // 聲明 Fanout 廣播交換機:持久化、不自動刪除
    @Bean
    public FanoutExchange fanoutExchange() {
        return new FanoutExchange(FANOUT_EXCHANGE, true, false);
    }
    // 隊列1:緩存刷新隊列
    @Bean
    public Queue cacheRefreshQueue() {
        return new Queue(QUEUE_CACHE_REFRESH, true);
    }
    // 隊列2:系統(tǒng)通知隊列
    @Bean
    public Queue noticeQueue() {
        return new Queue(QUEUE_NOTICE, true);
    }
    // 隊列3:日志采集隊列
    @Bean
    public Queue logQueue() {
        return new Queue(QUEUE_LOG, true);
    }
    // 全部綁定到廣播交換機
    @Bean
    public Binding bindingCacheRefresh(Queue cacheRefreshQueue, FanoutExchange fanoutExchange) {
        return BindingBuilder.bind(cacheRefreshQueue).to(fanoutExchange);
    }
    @Bean
    public Binding bindingNotice(Queue noticeQueue, FanoutExchange fanoutExchange) {
        return BindingBuilder.bind(noticeQueue).to(fanoutExchange);
    }
    @Bean
    public Binding bindingLog(Queue logQueue, FanoutExchange fanoutExchange) {
        return BindingBuilder.bind(logQueue).to(fanoutExchange);
    }
}

4. 廣播消息生產(chǎn)者

import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class BroadcastProducer {
    @Autowired
    private RabbitTemplate rabbitTemplate;
    @GetMapping("/send/broadcast")
    public String sendBroadcastMsg(@RequestParam String content) {
        // Fanout廣播:路由鍵傳空字符串
        rabbitTemplate.convertAndSend(
                FanoutBroadcastConfig.FANOUT_EXCHANGE,
                "",
                content
        );
        return "? 廣播消息發(fā)送成功:" + content;
    }
}

5. 多消費者實現(xiàn)(全員接收)

消費者1:緩存刷新消費者

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import java.io.IOException;
@Component
public class CacheRefreshConsumer {
    @RabbitListener(queues = FanoutBroadcastConfig.QUEUE_CACHE_REFRESH)
    public void consume(String msg, Message message, Channel channel) throws IOException {
        try {
            System.out.println("【緩存服務】接收廣播消息:" + msg);
            // 執(zhí)行緩存刷新業(yè)務邏輯
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            // 消費失敗,重回隊列重試
            channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
        }
    }
}

消費者2:系統(tǒng)通知消費者

@Component
public class SystemNoticeConsumer {
    @RabbitListener(queues = FanoutBroadcastConfig.QUEUE_NOTICE)
    public void consume(String msg, Message message, Channel channel) throws IOException {
        try {
            System.out.println("【通知服務】接收廣播消息:" + msg);
            // 執(zhí)行消息推送業(yè)務
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
        }
    }
}

消費者3:日志采集消費者

@Component
public class LogCollectConsumer {
    @RabbitListener(queues = FanoutBroadcastConfig.QUEUE_LOG)
    public void consume(String msg, Message message, Channel channel) throws IOException {
        try {
            System.out.println("【日志服務】接收廣播消息:" + msg);
            // 執(zhí)行日志采集業(yè)務
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
        }
    }
}

五、測試效果

訪問接口:

http://localhost:8080/send/broadcast?content=全局緩存刷新通知

控制臺輸出:

【緩存服務】接收廣播消息:全局緩存刷新通知
【通知服務】接收廣播消息:全局緩存刷新通知
【日志服務】接收廣播消息:全局緩存刷新通知

一條消息,多服務同時消費,廣播生效!

六、總結

1. 必須開啟持久化

交換機、隊列全部設置持久化,防止重啟丟失廣播配置。

2. 強制手動ACK

廣播場景多為重要通知、緩存同步,自動ACK會導致業(yè)務未執(zhí)行完成消息丟失。

3. 每個服務獨立隊列

不同微服務必須使用獨立隊列,避免競爭消費,保證廣播全覆蓋。

4. 廣播消息建議做冪等

MQ 重試、網(wǎng)絡抖動會導致廣播消息重復推送,核心業(yè)務必須基于消息ID做冪等防重。

5. 禁止設置復雜路由鍵

Fanout 無視路由鍵,統(tǒng)一傳空字符串,保持代碼規(guī)范。

寫在最后

廣播發(fā)布訂閱模式是微服務分布式通信的重要基石,區(qū)別于傳統(tǒng)的點對點任務消費,它主打全局通知、多節(jié)點同步、狀態(tài)廣播,是緩存刷新、配置熱更新、系統(tǒng)公告等場景的最優(yōu)解。

很多開發(fā)者一直混淆“競爭消費”和“廣播消費”的本質(zhì),導致線上通知不全、同步失效等隱性問題。吃透 Fanout 交換機的底層原理與落地規(guī)范,能幫你徹底解決分布式多節(jié)點同步難題。

以上就是SpringBoot實現(xiàn)發(fā)布/訂閱廣播消息的示例代碼的詳細內(nèi)容,更多關于SpringBoot發(fā)布/訂閱廣播消息的資料請關注腳本之家其它相關文章!

相關文章

  • Apache?Commons?Imaging處理圖像實例深究

    Apache?Commons?Imaging處理圖像實例深究

    這篇文章主要為大家介紹了Apache?Commons?Imaging處理圖像的實例探索深究,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-12-12
  • Java中的Callable實現(xiàn)多線程詳解

    Java中的Callable實現(xiàn)多線程詳解

    這篇文章主要介紹了Java中的Callable實現(xiàn)多線程詳解,接口Callable中有一個call方法,其返回值類型為V,這是一個泛型,值得關注的是這個call方法有返回值,這意味著線程執(zhí)行完畢后可以將處理結果返回,需要的朋友可以參考下
    2023-08-08
  • java中新生代和老生代的關系說明

    java中新生代和老生代的關系說明

    這篇文章主要介紹了java中新生代和老生代的關系說明,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-07-07
  • java8 對象轉Map時重復 key Duplicate key xxxx的解決

    java8 對象轉Map時重復 key Duplicate key xxxx的解決

    這篇文章主要介紹了java8 對象轉Map時重復 key Duplicate key xxxx的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-09-09
  • 淺談java實現(xiàn)背包算法(0-1背包問題)

    淺談java實現(xiàn)背包算法(0-1背包問題)

    本篇文章主要介紹了淺談java實現(xiàn)背包算法(0-1背包問題) ,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-08-08
  • Java實現(xiàn)字符串切割的方法詳解

    Java實現(xiàn)字符串切割的方法詳解

    這篇文章主要為大家介紹了一些Java中切割字符串的小技巧,可以把性能提升5~10倍。文中的示例代碼講解詳細,快跟隨小編一起學習一下
    2022-03-03
  • java解析xml常用的幾種方式總結

    java解析xml常用的幾種方式總結

    這篇文章主要介紹了java解析xml常用的幾種方式總結,有需要的朋友可以參考一下
    2013-11-11
  • java中json-diff簡單使用及對象是否一致詳解

    java中json-diff簡單使用及對象是否一致詳解

    這篇文章主要為大家介紹了java中json-diff簡單使用及對象是否一致對比詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2023-03-03
  • SpringBoot中的PropertySource原理詳解

    SpringBoot中的PropertySource原理詳解

    這篇文章主要介紹了SpringBoot中的PropertySource原理詳解,PropertySource?是一個非常重要的概念,它允許您在應用程序中定義屬性,并將這些屬性注入到?Spring?環(huán)境中,需要的朋友可以參考下
    2023-07-07
  • java代理 jdk動態(tài)代理應用案列

    java代理 jdk動態(tài)代理應用案列

    java代理有jdk動態(tài)代理、cglib代理,這里只說下jdk動態(tài)代理,jdk動態(tài)代理主要使用的是java反射機制,需要了解的朋友可以參考下
    2012-11-11

最新評論

柳林县| 叶城县| 蕉岭县| 子长县| 克什克腾旗| 山东| 阳山县| 嘉祥县| 镇雄县| 安岳县| 大邑县| 天台县| 永康市| 夏河县| 新建县| 大化| 常熟市| 岳池县| 芷江| 运城市| 集贤县| 柘荣县| 郧西县| 西林县| 宣汉县| 镇远县| 东山县| 通化市| 花莲县| 定日县| 绥江县| 沛县| 饶河县| 延吉市| 北票市| 穆棱市| 西和县| 广丰县| 锡林郭勒盟| 龙山县| 五寨县|