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

RabbitMQ集群實(shí)現(xiàn)消息的高可用性和負(fù)載均衡

 更新時(shí)間:2026年04月27日 09:36:38   作者:代碼漫談  
文章解釋了RabbitMQ集群的概念及其構(gòu)建目的,詳細(xì)闡述了其底層原理,包括Erlang/OTP的分布式基因、ErlangCookie認(rèn)證和分布式進(jìn)程管理,感興趣的朋友跟隨小編一起看看吧

一、核心概念:集群是什么?為何而生?

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)題

  1. 高可用性:避免單點(diǎn)故障(SPOF)。
  2. 橫向擴(kuò)展:通過(guò)增加節(jié)點(diǎn),分散連接和信道的壓力,提升整體吞吐量。
  3. 數(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)。

  1. 節(jié)點(diǎn)通信:集群節(jié)點(diǎn)間通過(guò) Erlang Cookie 進(jìn)行認(rèn)證。這是一個(gè)相同的字符串,存儲(chǔ)在 $HOME/.erlang.cookie 文件中。只有Cookie相同的Erlang節(jié)點(diǎn)才能組成集群。
  2. 分布式進(jìn)程:在Erlang看來(lái),RabbitMQ的每個(gè)隊(duì)列、信道都是一個(gè)“進(jìn)程”。集群將這些進(jìn)程及其狀態(tài)信息(元數(shù)據(jù))在節(jié)點(diǎn)間通過(guò) Erlang分布式消息傳遞 進(jìn)行同步,效率極高。
  3. 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í)的提升行為。

七、Spring AMQP整合:聲明鏡像隊(duì)列

在Spring Boot中,我們可以通過(guò)RabbitAdminCachingConnectionFactory優(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à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_minorityautoheal 兩種處理策略分別如何應(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í)算法。

延伸閱讀

  1. RabbitMQ官方文檔 - https://www.rabbitmq.com/clustering.html
  2. RabbitMQ官方文檔 - https://www.rabbitmq.com/ha.html
  3. RabbitMQ官方文檔 - https://www.rabbitmq.com/quorum-queues.html
  4. 《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)題

    完美解決在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查看和修改線(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ā)貨接口

    如何用java對(duì)接微信小程序下單后的發(fā)貨接口

    這篇文章主要介紹了在微信小程序后臺(tái)實(shí)現(xiàn)發(fā)貨通知的步驟,包括獲取Access_token、使用RestTemplate調(diào)用發(fā)貨接口、處理AccessToken緩存以及發(fā)貨成功后的提醒,需要的朋友可以參考下
    2025-03-03
  • Java實(shí)現(xiàn)洗牌發(fā)牌的方法

    Java實(shí)現(xiàn)洗牌發(fā)牌的方法

    這篇文章主要介紹了Java實(shí)現(xiàn)洗牌發(fā)牌的方法,涉及java針對(duì)數(shù)組的遍歷與排序操作相關(guān)技巧,具有一定參考借鑒價(jià)值,需要的朋友可以參考下
    2015-07-07
  • SpringBoot結(jié)合mybatis-plus實(shí)現(xiàn)分頁(yè)的項(xiàng)目實(shí)踐

    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é)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 Eureka探究細(xì)枝末節(jié)

    Java Eureka探究細(xì)枝末節(jié)

    Eureka是Netflix開(kāi)發(fā)的服務(wù)發(fā)現(xiàn)框架,本身是一個(gè)基于REST的服務(wù),主要用于定位運(yùn)行在A(yíng)WS域中的中間層服務(wù),以達(dá)到負(fù)載均衡和中間層服務(wù)故障轉(zhuǎn)移的目的
    2022-09-09
  • SpringBoot使用minio及配置代碼

    SpringBoot使用minio及配置代碼

    MinIO是一個(gè)非常輕量的服務(wù),可以很簡(jiǎn)單的和其他應(yīng)用的結(jié)合,類(lèi)似?NodeJS,?Redis?或者?MySQL。本文重點(diǎn)給大家介紹SpringBoot使用minio及配置代碼,感興趣的朋友一起看看吧
    2022-02-02
  • Java concurrency集合之ConcurrentSkipListMap_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理

    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ò)濾器體系代碼分析

    這篇文章主要介紹了Spring Security中的Servlet過(guò)濾器體系,本文通過(guò)圖文并茂的形式給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-07-07

最新評(píng)論

宜丰县| 蓝山县| 榆社县| 自治县| 溆浦县| 义乌市| 贵德县| 高唐县| 南平市| 藁城市| 沂南县| 邮箱| 长宁县| 廊坊市| 从江县| 壶关县| 江源县| 顺昌县| 红安县| 庆城县| 卢龙县| 兴化市| 兴宁市| 宁阳县| 额敏县| 绥棱县| 丰县| 宁陵县| 衡东县| 汝阳县| 泌阳县| 四会市| 萝北县| 宜丰县| 海原县| 岳普湖县| 屏边| 雅江县| 文成县| 淄博市| 芮城县|