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: 10003. 廣播交換機、隊列、綁定配置類
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處理圖像的實例探索深究,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪2023-12-12
java8 對象轉Map時重復 key Duplicate key xxxx的解決
這篇文章主要介紹了java8 對象轉Map時重復 key Duplicate key xxxx的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-09-09
SpringBoot中的PropertySource原理詳解
這篇文章主要介紹了SpringBoot中的PropertySource原理詳解,PropertySource?是一個非常重要的概念,它允許您在應用程序中定義屬性,并將這些屬性注入到?Spring?環(huán)境中,需要的朋友可以參考下2023-07-07

