RabbitMQ集群實(shí)現(xiàn)消息的高可用性和負(fù)載均衡
一、核心概念:集群是什么?為何而生?
RabbitMQ集群的本質(zhì),是將多個(gè)RabbitMQ節(jié)點(diǎn)(Node)連接成一個(gè)邏輯上的消息服務(wù)器。這些節(jié)點(diǎn)共享一部分?jǐn)?shù)據(jù)和狀態(tài)(主要是元數(shù)據(jù)),從而在單個(gè)節(jié)點(diǎn)故障時(shí),其他節(jié)點(diǎn)可以接管其工作,保證服務(wù)不中斷。
它解決的核心問(wèn)題:
- 高可用性:避免單點(diǎn)故障(SPOF)。
- 橫向擴(kuò)展:通過(guò)增加節(jié)點(diǎn),分散連接和信道的壓力,提升整體吞吐量。
- 數(shù)據(jù)冗余:通過(guò)鏡像隊(duì)列,在多節(jié)點(diǎn)間復(fù)制消息,防止數(shù)據(jù)丟失。
需要注意:RabbitMQ集群默認(rèn)是“元數(shù)據(jù)共享,消息不共享”。這意味著交換機(jī)、隊(duì)列的定義和綁定關(guān)系在所有節(jié)點(diǎn)同步,但隊(duì)列中的消息默認(rèn)只存在于聲明它的那個(gè)節(jié)點(diǎn)上。要實(shí)現(xiàn)消息的冗余,必須依賴(lài)于鏡像隊(duì)列(Mirrored Queues)策略。
二、底層原理:Erlang/OTP的分布式基因
RabbitMQ的集群能力并非后天嫁接,而是深深植根于其底層實(shí)現(xiàn)語(yǔ)言——Erlang 的基因之中。Erlang/OTP平臺(tái)生來(lái)就是為了構(gòu)建高并發(fā)、分布式、高可用的電信級(jí)系統(tǒng)。
- 節(jié)點(diǎn)通信:集群節(jié)點(diǎn)間通過(guò) Erlang Cookie 進(jìn)行認(rèn)證。這是一個(gè)相同的字符串,存儲(chǔ)在
$HOME/.erlang.cookie文件中。只有Cookie相同的Erlang節(jié)點(diǎn)才能組成集群。 - 分布式進(jìn)程:在Erlang看來(lái),RabbitMQ的每個(gè)隊(duì)列、信道都是一個(gè)“進(jìn)程”。集群將這些進(jìn)程及其狀態(tài)信息(元數(shù)據(jù))在節(jié)點(diǎn)間通過(guò) Erlang分布式消息傳遞 進(jìn)行同步,效率極高。
- Mnesia數(shù)據(jù)庫(kù):RabbitMQ使用Erlang內(nèi)置的分布式數(shù)據(jù)庫(kù)Mnesia來(lái)存儲(chǔ)集群的元數(shù)據(jù)(交換機(jī)、隊(duì)列、綁定、用戶(hù)、vhost等)。Mnesia確保了元數(shù)據(jù)在集群內(nèi)強(qiáng)一致性。
原理流程圖:元數(shù)據(jù)同步

圖示:聲明隊(duì)列時(shí),其元數(shù)據(jù)(定義)在集群內(nèi)強(qiáng)一致同步,但承載消息的Master隊(duì)列進(jìn)程只在一個(gè)節(jié)點(diǎn)創(chuàng)建。*
三、集群部署實(shí)戰(zhàn):三步構(gòu)建集群
我們以三臺(tái)機(jī)器(node1, node2, node3)為例,演示如何手動(dòng)搭建集群。
前置條件:所有節(jié)點(diǎn)安裝相同版本的RabbitMQ和Erlang,并確保 ~/.erlang.cookie 文件內(nèi)容一致。
步驟1:?jiǎn)?dòng)各節(jié)點(diǎn)
# 在 node1, node2, node3 上分別執(zhí)行 sudo systemctl start rabbitmq-server # 或 rabbitmq-server -detached
步驟2:將 node2, node3 加入 node1 的集群
# 在 node2 上執(zhí)行 rabbitmqctl stop_app rabbitmqctl reset # 如果是新節(jié)點(diǎn),可不reset。如果是已存數(shù)據(jù)的節(jié)點(diǎn),reset會(huì)清除數(shù)據(jù)! rabbitmqctl join_cluster rabbit@node1 # 注意:rabbit是默認(rèn)的Erlang節(jié)點(diǎn)名前綴 rabbitmqctl start_app # 在 node3 上執(zhí)行相同操作 rabbitmqctl stop_app rabbitmqctl join_cluster rabbit@node1 rabbitmqctl start_app
步驟3:驗(yàn)證集群狀態(tài)
# 在任意節(jié)點(diǎn)執(zhí)行 rabbitmqctl cluster_status
運(yùn)行結(jié)果示例:
Cluster status of node rabbit@node1 ... Basics Cluster name: rabbit@node1 Disk Nodes rabbit@node1 rabbit@node2 rabbit@node3 Running Nodes rabbit@node1 rabbit@node2 rabbit@node3 ...
四、靈魂所在:鏡像隊(duì)列(Mirrored Queues)
默認(rèn)集群無(wú)法解決消息丟失問(wèn)題,鏡像隊(duì)列才是實(shí)現(xiàn)高可用的關(guān)鍵。它會(huì)將隊(duì)列的內(nèi)容(消息)復(fù)制到集群中的其他多個(gè)節(jié)點(diǎn)上。
配置策略:通過(guò)策略(Policy)來(lái)為匹配的隊(duì)列設(shè)置鏡像規(guī)則。
# 設(shè)置一個(gè)策略,將名稱(chēng)以“mirrored.”開(kāi)頭的隊(duì)列鏡像到所有節(jié)點(diǎn)
rabbitmqctl set_policy ha-all "^mirrored\." '{"ha-mode":"all"}'
# 更常見(jiàn)的生產(chǎn)配置:鏡像到多數(shù)節(jié)點(diǎn)(N/2+1),例如3節(jié)點(diǎn)集群中鏡像到2個(gè)
rabbitmqctl set_policy ha-majority "^ha\." '{"ha-mode":"exactly", "ha-params":2, "ha-sync-mode":"automatic"}'- ha-mode:
all,exactly,nodes - ha-sync-mode:
manual(手動(dòng)同步,默認(rèn)) 或automatic(自動(dòng)同步)。生產(chǎn)環(huán)境建議在業(yè)務(wù)低峰期手動(dòng)同步,或新建隊(duì)列時(shí)自動(dòng)同步,避免高峰同步引發(fā)網(wǎng)絡(luò)阻塞。
鏡像隊(duì)列架構(gòu)圖:

圖示:生產(chǎn)者/消費(fèi)者只與Master副本交互。Master故障后,最老的副本會(huì)被提升為新的Master。*
五、客戶(hù)端連接:負(fù)載均衡與故障轉(zhuǎn)移
客戶(hù)端連接集群時(shí),不應(yīng)寫(xiě)死一個(gè)節(jié)點(diǎn)地址,而應(yīng)提供節(jié)點(diǎn)列表,客戶(hù)端會(huì)按順序嘗試連接。
Java客戶(hù)端連接示例:
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
public class ClusterConnectionExample {
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setUsername("guest");
factory.setPassword("guest");
factory.setVirtualHost("/");
// 關(guān)鍵:設(shè)置集群節(jié)點(diǎn)地址數(shù)組
// 格式:amqp:// 可省略,默認(rèn)端口5672
Address[] addresses = new Address[] {
new Address("node1", 5672),
new Address("node2", 5672),
new Address("node3", 5672)
};
// 自動(dòng)重試連接。也可以使用外部負(fù)載均衡器(如HAProxy、F5)
Connection connection = factory.newConnection(addresses, "MyApp-Client");
System.out.println("成功連接到集群: " + connection.getAddress().getHostAddress());
// ... 使用 connection 創(chuàng)建 Channel 進(jìn)行后續(xù)操作
connection.close();
}
}運(yùn)行結(jié)果:
成功連接到集群: node1 (假設(shè)node1是第一個(gè)可用的節(jié)點(diǎn))
六、企業(yè)級(jí)最佳實(shí)踐與配置詳解
- 節(jié)點(diǎn)類(lèi)型:
- 磁盤(pán)節(jié)點(diǎn)(Disc Node):將元數(shù)據(jù)和消息存儲(chǔ)到磁盤(pán)。集群中至少需要一個(gè)磁盤(pán)節(jié)點(diǎn)來(lái)保證元數(shù)據(jù)可靠性。 生產(chǎn)環(huán)境通常所有節(jié)點(diǎn)都是磁盤(pán)節(jié)點(diǎn)。
- 內(nèi)存節(jié)點(diǎn)(RAM Node):將元數(shù)據(jù)僅存儲(chǔ)在內(nèi)存,性能極高,但重啟后數(shù)據(jù)丟失(依賴(lài)磁盤(pán)節(jié)點(diǎn)同步)。適用于臨時(shí)、高性能的中間節(jié)點(diǎn),但增加了運(yùn)維復(fù)雜度,不推薦新手使用。
- 集群規(guī)模:建議3或5個(gè)節(jié)點(diǎn)(奇數(shù)個(gè)),便于維持仲裁多數(shù)(Quorum)。超過(guò)5個(gè)節(jié)點(diǎn),同步開(kāi)銷(xiāo)會(huì)顯著增大,收益遞減。
- 網(wǎng)絡(luò)要求:集群節(jié)點(diǎn)間需要穩(wěn)定的低延遲網(wǎng)絡(luò)。避免跨地域部署集群,應(yīng)使用同地域多可用區(qū)。
- 鏡像隊(duì)列策略:
- ha-sync-mode: 新鏡像節(jié)點(diǎn)加入時(shí),
automatic會(huì)同步所有歷史消息,可能阻塞隊(duì)列。建議在業(yè)務(wù)低峰期通過(guò)rabbitmqctl sync_queue {queue_name}手動(dòng)同步。 - ha-promote-on-shutdown: 控制主節(jié)點(diǎn)優(yōu)雅關(guān)閉時(shí)是否提升鏡像。默認(rèn)為
when-synced,只有已同步的鏡像才會(huì)被提升,更安全。 - ha-promote-on-failure:類(lèi)似,控制故障時(shí)的提升行為。
- ha-sync-mode: 新鏡像節(jié)點(diǎn)加入時(shí),
七、Spring AMQP整合:聲明鏡像隊(duì)列
在Spring Boot中,我們可以通過(guò)RabbitAdmin和CachingConnectionFactory優(yōu)雅地集成集群。
import org.springframework.amqp.core.*;
import org.springframework.amqp.rabbit.connection.CachingConnectionFactory;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class RabbitMQClusterConfig {
@Bean
public CachingConnectionFactory connectionFactory() {
CachingConnectionFactory factory = new CachingConnectionFactory();
// 設(shè)置集群節(jié)點(diǎn)地址
factory.setAddresses("node1:5672,node2:5672,node3:5672");
factory.setUsername("guest");
factory.setPassword("guest");
factory.setVirtualHost("/");
// 開(kāi)啟Publisher Confirms,保證消息可靠投遞
factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED);
// 開(kāi)啟Publisher Returns,監(jiān)聽(tīng)不可路由消息
factory.setPublisherReturns(true);
return factory;
}
@Bean
public RabbitAdmin rabbitAdmin(CachingConnectionFactory connectionFactory) {
return new RabbitAdmin(connectionFactory);
}
@Bean
public Queue mirroredQueue() {
Map<String, Object> args = new HashMap<>();
// 在代碼中也可以聲明隊(duì)列參數(shù),但通常更推薦用Policy在服務(wù)端統(tǒng)一定義
// args.put("x-ha-policy", "all"); // 舊參數(shù),已棄用
// 推薦在RabbitMQ管理后臺(tái)或通過(guò)rabbitmqctl set_policy設(shè)置
return new Queue("order.queue",
true, // durable 持久化
false, // exclusive 非獨(dú)占
false, // autoDelete 不自動(dòng)刪除
args);
}
@Bean
public DirectExchange orderExchange() {
return new DirectExchange("order.exchange", true, false);
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(mirroredQueue())
.to(orderExchange())
.with("order.created");
}
}八、生產(chǎn)環(huán)境高可用架構(gòu):集群+負(fù)載均衡
單一集群仍可能因機(jī)房故障而癱瘓。真正的企業(yè)級(jí)方案是 “多集群鏡像” 或 “集群+負(fù)載均衡器”。
推薦架構(gòu):
[生產(chǎn)者/消費(fèi)者]
|
[負(fù)載均衡器: HAProxy/Nginx] (Virtual IP)
|
[RabbitMQ集群]
/ | \
Node1 Node2 Node3 (磁盤(pán)節(jié)點(diǎn),同機(jī)房)- 負(fù)載均衡器:對(duì)外提供單一入口,實(shí)現(xiàn)節(jié)點(diǎn)故障自動(dòng)切換。
- 客戶(hù)端:只需連接VIP。當(dāng)Node1宕機(jī),連接自動(dòng)轉(zhuǎn)移到Node2。
HAProxy配置示例片段:
listen rabbitmq_cluster
bind 0.0.0.0:5670
mode tcp
balance roundrobin
timeout client 3h
timeout server 3h
option clitcpka
server node1 node1:5672 check inter 5s rise 2 fall 3
server node2 node2:5672 check inter 5s rise 2 fall 3
server node3 node3:5672 check inter 5s rise 2 fall 3客戶(hù)端連接地址改為 haproxy-host:5670。
九、監(jiān)控、運(yùn)維與故障排查
- 監(jiān)控指標(biāo):
- 節(jié)點(diǎn)狀態(tài):
rabbitmqctl cluster_status,rabbitmqctl node_health_check - 隊(duì)列狀態(tài):消息數(shù)、消費(fèi)者數(shù)、內(nèi)存占用、同步狀態(tài)(對(duì)于鏡像隊(duì)列)。
- Erlang運(yùn)行時(shí):內(nèi)存、進(jìn)程數(shù)、ETS表數(shù)量。
- 管理界面:訪(fǎng)問(wèn)
http://node-ip:15672,在 Admin -> Cluster 查看節(jié)點(diǎn)信息,在 Queues 頁(yè)查看隊(duì)列的“Features”是否包含HA(表示是鏡像隊(duì)列)及同步狀態(tài)(+1表示有一個(gè)鏡像,+2表示兩個(gè))。
- 節(jié)點(diǎn)狀態(tài):
- 常見(jiàn)故障排查:
- 節(jié)點(diǎn)無(wú)法加入集群:檢查Cookie、主機(jī)名解析、防火墻、Erlang版本。
- 鏡像隊(duì)列不同步:檢查網(wǎng)絡(luò),查看
rabbitmqctl list_queues name slave_pids synchronised_slave_pids,手動(dòng)觸發(fā)同步。 - 腦裂(Network Partition):這是最嚴(yán)重的問(wèn)題。RabbitMQ會(huì)檢測(cè)到網(wǎng)絡(luò)分區(qū),并顯示在管理界面。處理需謹(jǐn)慎,可設(shè)置
cluster_partition_handling參數(shù)為pause_minority(少數(shù)派節(jié)點(diǎn)自動(dòng)暫停)或autoheal(自動(dòng)恢復(fù)),但都有風(fēng)險(xiǎn),最好從網(wǎng)絡(luò)層面避免。
十、易錯(cuò)點(diǎn)與常見(jiàn)誤區(qū)
誤區(qū)1:搭建了集群,消息就不會(huì)丟失了。
- 正確理解:默認(rèn)集群只同步元數(shù)據(jù)。必須配合使用【持久化消息(Delivery mode=2)】、【持久化隊(duì)列】和【鏡像隊(duì)列策略】, 才能實(shí)現(xiàn)節(jié)點(diǎn)故障時(shí)的消息不丟失。同時(shí),生產(chǎn)者需要使用Publisher Confirm,消費(fèi)者需要手動(dòng)ACK。
誤區(qū)2:鏡像隊(duì)列越多,可靠性越高,性能越好。
- 正確理解:鏡像隊(duì)列帶來(lái)可靠性的同時(shí),會(huì)顯著增加網(wǎng)絡(luò)開(kāi)銷(xiāo)和磁盤(pán)IO。每一條消息都要同步到所有鏡像節(jié)點(diǎn),性能會(huì)下降。通常鏡像到2-3個(gè)節(jié)點(diǎn)(包含Master)是性?xún)r(jià)比較高的選擇。
誤區(qū)3:消費(fèi)者可以連接任意節(jié)點(diǎn)消費(fèi)任意隊(duì)列。
- 正確理解:消費(fèi)者可以連接集群任意節(jié)點(diǎn)。但如果該節(jié)點(diǎn)不是隊(duì)列Master所在節(jié)點(diǎn),RabbitMQ會(huì)在內(nèi)部建立跨節(jié)點(diǎn)信道,將消息從Master節(jié)點(diǎn)路由到消費(fèi)者連接的節(jié)點(diǎn),這會(huì)增加內(nèi)部網(wǎng)絡(luò)流量。最佳實(shí)踐是讓消費(fèi)者盡量連接到隊(duì)列Master所在的節(jié)點(diǎn)(但這通常由負(fù)載均衡器決定,難以精確控制)。
易混淆概念:持久化 vs 鏡像
- 持久化:將消息和隊(duì)列定義寫(xiě)入磁盤(pán),防止RabbitMQ服務(wù)重啟導(dǎo)致數(shù)據(jù)丟失。
- 鏡像:將隊(duì)列的內(nèi)容(消息)復(fù)制到其他節(jié)點(diǎn),防止單個(gè)節(jié)點(diǎn)物理故障導(dǎo)致數(shù)據(jù)丟失。
- 關(guān)系:兩者需結(jié)合使用。一個(gè)隊(duì)列可以持久化但不鏡像(單點(diǎn)重啟不丟,宕機(jī)丟);也可以鏡像但不持久化(節(jié)點(diǎn)宕機(jī)切換不丟,重啟丟)。生產(chǎn)環(huán)境必須同時(shí)開(kāi)啟。
十一、高級(jí)主題:仲裁隊(duì)列(Quorum Queues) vs 鏡像隊(duì)列
自RabbitMQ 3.8版本引入了仲裁隊(duì)列,它是為集群環(huán)境設(shè)計(jì)的現(xiàn)代隊(duì)列類(lèi)型,旨在解決經(jīng)典鏡像隊(duì)列的一些復(fù)雜性問(wèn)題。
對(duì)比表:
| 特性 | 經(jīng)典鏡像隊(duì)列 (Classic Mirrored) | 仲裁隊(duì)列 (Quorum Queue) |
|---|---|---|
| 設(shè)計(jì)目標(biāo) | 在傳統(tǒng)主從復(fù)制上增加高可用 | 基于Raft共識(shí)算法,強(qiáng)一致性、數(shù)據(jù)安全 |
| 數(shù)據(jù)一致性 | 最終一致性(異步復(fù)制) | 強(qiáng)一致性(多數(shù)節(jié)點(diǎn)確認(rèn)) |
| 節(jié)點(diǎn)故障處理 | 自動(dòng)故障轉(zhuǎn)移,可能丟消息(未同步部分) | 自動(dòng)故障轉(zhuǎn)移,保證已確認(rèn)消息不丟 |
| 配置復(fù)雜度 | 高(需理解Policy, ha參數(shù)) | 低(聲明時(shí)指定類(lèi)型即可) |
| 性能 | 較高(異步復(fù)制) | 略低(強(qiáng)一致需多數(shù)節(jié)點(diǎn)確認(rèn)) |
| 推薦場(chǎng)景 | 高吞吐,允許極小概率消息丟失的常規(guī)業(yè)務(wù) | 金融交易、訂單狀態(tài)等對(duì)一致性要求極高的核心業(yè)務(wù) |
聲明仲裁隊(duì)列:
import org.springframework.amqp.core.QueueBuilder;
@Bean
public Queue quorumOrderQueue() {
return QueueBuilder.durable("quorum.order.queue")
.quorum() // 設(shè)置為仲裁隊(duì)列
.build();
}十二、Java生產(chǎn)-消費(fèi)示例
場(chǎng)景:模擬一個(gè)訂單創(chuàng)建后,發(fā)送到高可用集群的鏡像隊(duì)列。
1. 生產(chǎn)者 (OrderProducer.java)
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import java.nio.charset.StandardCharsets;
public class OrderProducer {
private final static String EXCHANGE_NAME = "order.exchange.ha";
private final static String ROUTING_KEY = "order.created";
private final static String QUEUE_NAME = "order.queue.ha";
public static void main(String[] argv) throws Exception {
// 1. 連接工廠(chǎng)
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost"); // 生產(chǎn)環(huán)境應(yīng)為負(fù)載均衡器地址
factory.setPort(5672);
factory.setUsername("guest");
factory.setPassword("guest");
factory.setVirtualHost("/");
// 2. 創(chuàng)建連接和信道
// 使用try-with-resources確保資源關(guān)閉
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 3. 聲明持久化的直連交換機(jī)和持久化的隊(duì)列
// 注意:這些定義會(huì)在集群節(jié)點(diǎn)間同步(元數(shù)據(jù))
channel.exchangeDeclare(EXCHANGE_NAME, "direct", true);
// 隊(duì)列聲明,假設(shè)服務(wù)端已通過(guò)Policy將此隊(duì)列設(shè)為鏡像隊(duì)列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
// 4. 開(kāi)啟Publisher Confirms (異步確認(rèn))
channel.confirmSelect();
// 5. 發(fā)送持久化消息
String message = "訂單消息: ORDER-20240426-001";
channel.basicPublish(EXCHANGE_NAME,
ROUTING_KEY,
MessageProperties.PERSISTENT_TEXT_PLAIN, // 關(guān)鍵:消息持久化
message.getBytes(StandardCharsets.UTF_8));
System.out.println(" [x] 發(fā)送 '" + message + "'");
// 6. 等待Broker確認(rèn)
if (channel.waitForConfirms(5000)) {
System.out.println(" [?] 消息已得到Broker確認(rèn),投遞成功。");
} else {
System.out.println(" [x] 消息未得到Broker確認(rèn),可能投遞失敗。");
// 實(shí)際生產(chǎn)環(huán)境應(yīng)實(shí)現(xiàn)重試或落庫(kù)補(bǔ)償邏輯
}
} // try-with-resources會(huì)自動(dòng)關(guān)閉channel和connection
}
}運(yùn)行結(jié)果:
[x] 發(fā)送 '訂單消息: ORDER-20240426-001'
[?] 消息已得到Broker確認(rèn),投遞成功。
2. 消費(fèi)者 (OrderConsumer.java)
import com.rabbitmq.client.*;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeoutException;
public class OrderConsumer {
private final static String QUEUE_NAME = "order.queue.ha";
public static void main(String[] argv) throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("guest");
factory.setPassword("guest");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
// 為了保證公平分發(fā),在消費(fèi)者端設(shè)置Qos
int prefetchCount = 1; // 每次只預(yù)取1條消息
channel.basicQos(prefetchCount);
System.out.println(" [*] 等待消息。按 CTRL+C 退出");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println(" [x] 收到 '" + message + "'");
try {
// 模擬業(yè)務(wù)處理耗時(shí)
Thread.sleep(1000);
System.out.println(" [?] 業(yè)務(wù)處理完成: " + message);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
// 關(guān)鍵:手動(dòng)確認(rèn)消息,保證可靠性
// deliveryTag, multiple
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
System.out.println(" [?] 消息已手動(dòng)確認(rèn)。");
}
};
// 消費(fèi)消息,關(guān)閉自動(dòng)確認(rèn)(autoAck=false),采用手動(dòng)確認(rèn)
boolean autoAck = false;
channel.basicConsume(QUEUE_NAME, autoAck, deliverCallback, consumerTag -> {
System.out.println(" [x] 消費(fèi)者被取消: " + consumerTag);
});
}
}運(yùn)行結(jié)果:
[*] 等待消息。按 CTRL+C 退出
[x] 收到 '訂單消息: ORDER-20240426-001'
[?] 業(yè)務(wù)處理完成: 訂單消息: ORDER-20240426-001
[?] 消息已手動(dòng)確認(rèn)。
十三、性能、可靠性與監(jiān)控
性能考量:
- 網(wǎng)絡(luò)延遲是最大瓶頸:鏡像和仲裁隊(duì)列的性能?chē)?yán)重依賴(lài)節(jié)點(diǎn)間網(wǎng)絡(luò)延遲。務(wù)必保證集群節(jié)點(diǎn)在低延遲的網(wǎng)絡(luò)環(huán)境中(通常要求<1ms)。
- 磁盤(pán)速度:使用SSD硬盤(pán),特別是對(duì)于持久化隊(duì)列和仲裁隊(duì)列。
- 內(nèi)存:確保有足夠的內(nèi)存,Erlang VM會(huì)利用內(nèi)存進(jìn)行緩存。通過(guò)
rabbitmqctl status監(jiān)控memory部分。
可靠性 checklist:
- 集群節(jié)點(diǎn)數(shù) >= 3(且為奇數(shù))
- 所有節(jié)點(diǎn)為磁盤(pán)節(jié)點(diǎn)
- 配置了鏡像隊(duì)列或使用仲裁隊(duì)列
- 生產(chǎn)者使用 Publisher Confirm
- 消息設(shè)置為持久化 (
PERSISTENT) - 消費(fèi)者使用手動(dòng) ACK
- 有完善的監(jiān)控和告警(隊(duì)列長(zhǎng)度、節(jié)點(diǎn)狀態(tài)、內(nèi)存/磁盤(pán)使用率)
十四、延伸思考
思考一:在腦裂(網(wǎng)絡(luò)分區(qū))發(fā)生后,RabbitMQ 的兩個(gè)分區(qū)各自接管了部分隊(duì)列的Master,網(wǎng)絡(luò)恢復(fù)后,數(shù)據(jù)會(huì)出現(xiàn)什么沖突?RabbitMQ 的 pause_minority 和 autoheal 兩種處理策略分別如何應(yīng)對(duì)?各自的優(yōu)缺點(diǎn)是什么?
提示:從CAP理論的角度思考,RabbitMQ在腦裂時(shí)優(yōu)先保證了可用性(A),犧牲了部分一致性(C)。
引申閱讀:Raft共識(shí)算法、分布式系統(tǒng)腦裂處理。
思考二:如果我想實(shí)現(xiàn)跨地域(如北京-上海)的高可用,使用一個(gè)RabbitMQ集群是否合適?如果不合適,應(yīng)該采用什么架構(gòu)?
提示:考慮網(wǎng)絡(luò)延遲(通常>30ms)對(duì)鏡像同步性能和可靠性的致命影響。
引申閱讀:RabbitMQ Federation / Shovel 插件,實(shí)現(xiàn)集群間的消息轉(zhuǎn)發(fā)。
思考三:仲裁隊(duì)列(Quorum Queue)基于Raft協(xié)議,為什么它能保證強(qiáng)一致性和數(shù)據(jù)安全,但在寫(xiě)入延遲上會(huì)比經(jīng)典鏡像隊(duì)列高?
提示:思考Raft協(xié)議中“Leader選舉”、“日志復(fù)制”和“多數(shù)派提交”的過(guò)程。
引申閱讀:Raft論文,分布式共識(shí)算法。
延伸閱讀:
- RabbitMQ官方文檔 - https://www.rabbitmq.com/clustering.html
- RabbitMQ官方文檔 - https://www.rabbitmq.com/ha.html
- RabbitMQ官方文檔 - https://www.rabbitmq.com/quorum-queues.html
- 《RabbitMQ in Depth》 - 第7章 Clustering and High Availability
到此這篇關(guān)于RabbitMQ集群實(shí)現(xiàn)消息的高可用性和負(fù)載均衡的文章就介紹到這了,更多相關(guān)RabbitMQ集群高可用性和負(fù)載均衡內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
完美解決在Servlet中出現(xiàn)一個(gè)輸出中文亂碼的問(wèn)題
下面小編就為大家?guī)?lái)一篇完美解決在Servlet中出現(xiàn)一個(gè)輸出中文亂碼的問(wèn)題。小編覺(jué)得挺不錯(cuò)的現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2017-01-01
Java查看和修改線(xiàn)程優(yōu)先級(jí)操作詳解
JAVA中每個(gè)線(xiàn)程都有優(yōu)化級(jí)屬性,默認(rèn)情況下,新建的線(xiàn)程和創(chuàng)建該線(xiàn)程的線(xiàn)程優(yōu)先級(jí)是一樣的。本文將為大家詳解Java查看和修改線(xiàn)程優(yōu)先級(jí)操作的方法,需要的可以參考一下2022-08-08
如何用java對(duì)接微信小程序下單后的發(fā)貨接口
這篇文章主要介紹了在微信小程序后臺(tái)實(shí)現(xiàn)發(fā)貨通知的步驟,包括獲取Access_token、使用RestTemplate調(diào)用發(fā)貨接口、處理AccessToken緩存以及發(fā)貨成功后的提醒,需要的朋友可以參考下2025-03-03
SpringBoot結(jié)合mybatis-plus實(shí)現(xiàn)分頁(yè)的項(xiàng)目實(shí)踐
本文主要介紹了SpringBoot結(jié)合mybatis-plus實(shí)現(xiàn)分頁(yè)的項(xiàng)目實(shí)踐,主要基于MyBatis-Plus 自帶的分頁(yè)插件 PaginationInterceptor,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2023-06-06
手把手教學(xué)Win10同時(shí)安裝兩個(gè)版本的JDK并隨時(shí)切換(JDK8和JDK11)
最近在學(xué)習(xí)JDK11的一些新特性,但是日常使用基本上都是基于JDK8,因此,需要在win環(huán)境下安裝多個(gè)版本的JDK,下面這篇文章主要給大家介紹了手把手教學(xué)Win10同時(shí)安裝兩個(gè)版本的JDK(JDK8和JDK11)并隨時(shí)切換的相關(guān)資料,需要的朋友可以參考下2023-03-03
Java concurrency集合之ConcurrentSkipListMap_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理
這篇文章主要為大家詳細(xì)介紹了Java concurrency集合之ConcurrentSkipListMap的相關(guān)資料,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2017-06-06
Spring Security中的Servlet過(guò)濾器體系代碼分析
這篇文章主要介紹了Spring Security中的Servlet過(guò)濾器體系,本文通過(guò)圖文并茂的形式給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2020-07-07

