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

RabbitMQ中的Publish-Subscribe模式最佳實(shí)踐記錄

 更新時(shí)間:2024年12月19日 14:11:42   作者:AllenBright  
Publish/Subscribe 模式是 RabbitMQ 中一種強(qiáng)大且靈活的消息傳遞模式,適用于需要將消息廣播給多個(gè)訂閱者的場(chǎng)景,這篇文章主要介紹了RabbitMQ中的Publish-Subscribe模式,需要的朋友可以參考下

在現(xiàn)代分布式系統(tǒng)中,消息隊(duì)列(Message Queue)是實(shí)現(xiàn)異步通信和解耦系統(tǒng)的關(guān)鍵組件。RabbitMQ 是一個(gè)功能強(qiáng)大且廣泛使用的開(kāi)源消息代理,支持多種消息傳遞模式。其中,Publish/Subscribe(發(fā)布/訂閱)模式是一種常見(jiàn)且重要的模式,它允許消息發(fā)布者將消息廣播給多個(gè)訂閱者。

本文將深入探討 RabbitMQ 中的 Publish/Subscribe 模式,包括其工作原理、實(shí)現(xiàn)方式、適用場(chǎng)景以及最佳實(shí)踐。

1. Publish/Subscribe 模式簡(jiǎn)介

1.1 什么是 Publish/Subscribe 模式?

Publish/Subscribe(發(fā)布/訂閱)模式是一種消息傳遞模式,它將消息的發(fā)送者(發(fā)布者)和接收者(訂閱者)解耦。發(fā)布者將消息發(fā)布到一個(gè)交換機(jī)(Exchange),而訂閱者通過(guò)綁定到交換機(jī)的**隊(duì)列(Queue)**來(lái)接收消息。

與點(diǎn)對(duì)點(diǎn)模式(如工作隊(duì)列)不同,Publish/Subscribe 模式允許多個(gè)訂閱者接收相同的消息,從而實(shí)現(xiàn)消息的廣播。

1.2 核心概念

在 RabbitMQ 中,Publish/Subscribe 模式依賴(lài)以下核心組件:

  • 發(fā)布者(Publisher):發(fā)送消息的客戶(hù)端。
  • 交換機(jī)(Exchange):接收發(fā)布者發(fā)送的消息,并根據(jù)規(guī)則將消息路由到隊(duì)列。
  • 隊(duì)列(Queue):存儲(chǔ)消息的緩沖區(qū)。
  • 訂閱者(Subscriber):從隊(duì)列中消費(fèi)消息的客戶(hù)端。
  • 綁定(Binding):定義交換機(jī)和隊(duì)列之間的關(guān)系。

2. Publish/Subscribe 模式的工作原理

2.1 交換機(jī)的作用

在 RabbitMQ 中,消息不會(huì)直接發(fā)送到隊(duì)列,而是發(fā)送到交換機(jī)。交換機(jī)根據(jù)綁定規(guī)則將消息路由到相應(yīng)的隊(duì)列。

RabbitMQ 提供了多種類(lèi)型的交換機(jī),其中最常用的是:

  • Fanout 交換機(jī):將消息廣播到所有綁定到它的隊(duì)列,忽略路由鍵(Routing Key)。
  • Direct 交換機(jī):根據(jù)消息的路由鍵將消息路由到匹配的隊(duì)列。
  • Topic 交換機(jī):支持更復(fù)雜的路由規(guī)則,允許使用通配符匹配路由鍵。
  • Headers 交換機(jī):根據(jù)消息的頭部屬性進(jìn)行路由。

在 Publish/Subscribe 模式中,通常使用 Fanout 交換機(jī),因?yàn)樗軌驅(qū)⑾V播到所有綁定的隊(duì)列。

2.2 消息的廣播過(guò)程

  • 發(fā)布者將消息發(fā)送到交換機(jī)。
  • 交換機(jī)接收到消息后,將消息廣播到所有綁定的隊(duì)列。
  • 訂閱者從隊(duì)列中消費(fèi)消息。

3. Java 實(shí)現(xiàn) Publish/Subscribe 模式

以下是使用 Java 和 RabbitMQ Java Client 實(shí)現(xiàn) Publish/Subscribe 模式的完整示例。

3.1 添加依賴(lài)

在 Maven 項(xiàng)目中,添加 RabbitMQ Java Client 依賴(lài):

<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.20.0</version>
</dependency>

3.2 創(chuàng)建發(fā)布者(Publisher)

發(fā)布者負(fù)責(zé)將消息發(fā)送到交換機(jī)。以下是發(fā)布者的代碼:

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
public class Publisher {
    private static final String EXCHANGE_NAME = "publisher_subscriber";
    public static void main(String[] argv) throws Exception {
        // 創(chuàng)建連接工廠
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("192.168.200.138");
        factory.setPort(5672);
        factory.setVirtualHost("/test");
        factory.setUsername("test");
        factory.setPassword("test");
        try (Connection connection = factory.newConnection();
             Channel channel = connection.createChannel()) {
            // 聲明一個(gè) Fanout 交換機(jī)
            channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
            // 發(fā)布消息
            String message = "Hello, Subscribers!";
            channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes(StandardCharsets.UTF_8));
            System.out.println(" [x] Sent '" + message + "'");
        }
    }
}

3.3 創(chuàng)建訂閱者(Subscriber)

訂閱者負(fù)責(zé)從隊(duì)列中消費(fèi)消息。以下是訂閱者的代碼:

import com.rabbitmq.client.*;
import java.nio.charset.StandardCharsets;
public class Subscriber {
    private static final String EXCHANGE_NAME = "publisher_subscriber";
    public static void main(String[] argv) throws Exception {
        // 創(chuàng)建連接工廠
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("192.168.200.138");
        factory.setPort(5672);
        factory.setVirtualHost("/test");
        factory.setUsername("test");
        factory.setPassword("test");
        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();
        // 聲明一個(gè) Fanout 交換機(jī)
        channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
        // 創(chuàng)建一個(gè)臨時(shí)隊(duì)列,并綁定到交換機(jī)
        String queueName = channel.queueDeclare().getQueue();
        channel.queueBind(queueName, EXCHANGE_NAME, "");
        System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
        // 定義消息處理函數(shù)
        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
            System.out.println(" [x] Received '" + message + "'");
        };
        // 開(kāi)始消費(fèi)消息
        channel.basicConsume(queueName, true, deliverCallback, consumerTag -> {
        });
    }
}

3.4 運(yùn)行示例

啟動(dòng)多個(gè)訂閱者,在不同的終端窗口中運(yùn)行多個(gè)訂閱者實(shí)例

啟動(dòng)多個(gè)訂閱者后,能在RabbitMQ終端頁(yè)面,能看到多個(gè)臨時(shí)的隊(duì)列,但交換機(jī)只有一個(gè)publisher_subscriber。

啟動(dòng)發(fā)布者,在另一個(gè)終端窗口中運(yùn)行發(fā)布者 3.4.1 觀察輸出

所有訂閱者都會(huì)收到發(fā)布者發(fā)送的消息。例如:

發(fā)布者輸出:

 [x] Sent 'Hello, Subscribers!'

訂閱者輸出:

 [*] Waiting for messages. To exit press CTRL+C
 [x] Received 'Hello, Subscribers!'

4. 代碼解析

4.1 發(fā)布者代碼解析

  • 連接工廠ConnectionFactory 用于創(chuàng)建到 RabbitMQ 服務(wù)器的連接。
  • 交換機(jī)聲明channel.exchangeDeclare(EXCHANGE_NAME, "fanout") 聲明一個(gè) Fanout 交換機(jī)。
  • 消息發(fā)布channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes(StandardCharsets.UTF_8)) 將消息發(fā)送到交換機(jī)。

4.2 訂閱者代碼解析

  • 臨時(shí)隊(duì)列channel.queueDeclare().getQueue() 創(chuàng)建一個(gè)非持久化的、獨(dú)占的臨時(shí)隊(duì)列。
  • 隊(duì)列綁定channel.queueBind(queueName, EXCHANGE_NAME, "") 將隊(duì)列綁定到交換機(jī)。
  • 消息處理DeliverCallback 定義了如何處理接收到的消息。
  • 消費(fèi)消息channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { }) 開(kāi)始消費(fèi)消息。

5. Publish/Subscribe 模式的適用場(chǎng)景

5.1 日志記錄

在分布式系統(tǒng)中,日志記錄是一個(gè)常見(jiàn)的需求。使用 Publish/Subscribe 模式,可以將日志消息廣播給多個(gè)日志處理器,分別將日志寫(xiě)入文件、數(shù)據(jù)庫(kù)或發(fā)送到監(jiān)控系統(tǒng)。

5.2 實(shí)時(shí)通知

在社交網(wǎng)絡(luò)或即時(shí)通訊應(yīng)用中,可以使用 Publish/Subscribe 模式向多個(gè)用戶(hù)發(fā)送實(shí)時(shí)通知。例如,當(dāng)用戶(hù)發(fā)布新動(dòng)態(tài)時(shí),通知所有關(guān)注者。

5.3 分布式緩存更新

在分布式緩存系統(tǒng)中,當(dāng)緩存數(shù)據(jù)更新時(shí),可以使用 Publish/Subscribe 模式通知所有緩存節(jié)點(diǎn)同步更新。

5.4 事件驅(qū)動(dòng)架構(gòu)

在事件驅(qū)動(dòng)架構(gòu)中,Publish/Subscribe 模式用于實(shí)現(xiàn)事件的廣播。例如,當(dāng)用戶(hù)注冊(cè)成功時(shí),發(fā)布一個(gè)事件,通知多個(gè)服務(wù)(如郵件服務(wù)、積分服務(wù))執(zhí)行相應(yīng)的操作。

6. 最佳實(shí)踐

6.1 使用持久化

為了確保消息不會(huì)丟失,建議將交換機(jī)和隊(duì)列設(shè)置為持久化。例如:

channel.exchangeDeclare(EXCHANGE_NAME, "fanout", true);
channel.queueDeclare("my_queue", true, false, false, null);

6.2 處理消息確認(rèn)

在生產(chǎn)環(huán)境中,建議啟用消息確認(rèn)機(jī)制,確保消息被成功消費(fèi)。例如:

channel.basicConsume(queueName, false, deliverCallback, consumerTag -> { });

6.3 避免消息積壓

在高并發(fā)場(chǎng)景下,可能會(huì)出現(xiàn)消息積壓的情況。可以通過(guò)設(shè)置隊(duì)列的最大長(zhǎng)度或使用**死信隊(duì)列(DLX)**來(lái)處理積壓的消息。

6.4 監(jiān)控和報(bào)警

使用 RabbitMQ 的管理界面或監(jiān)控工具(如 Prometheus + Grafana)監(jiān)控消息隊(duì)列的狀態(tài),并設(shè)置報(bào)警規(guī)則,及時(shí)發(fā)現(xiàn)和解決問(wèn)題。

7. 總結(jié)

Publish/Subscribe 模式是 RabbitMQ 中一種強(qiáng)大且靈活的消息傳遞模式,適用于需要將消息廣播給多個(gè)訂閱者的場(chǎng)景。通過(guò)使用 Fanout 交換機(jī),可以輕松實(shí)現(xiàn)消息的廣播,同時(shí)結(jié)合持久化、消息確認(rèn)和監(jiān)控機(jī)制,可以構(gòu)建高可靠性的分布式系統(tǒng)。

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

相關(guān)文章

  • APT?注解處理器實(shí)現(xiàn)?Lombok?常用注解功能詳解

    APT?注解處理器實(shí)現(xiàn)?Lombok?常用注解功能詳解

    這篇文章主要為大家介紹了使用APT?注解處理器實(shí)現(xiàn)?Lombok?常用注解功能詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-09-09
  • Spring?boot整合jsp和tiles模板示例

    Spring?boot整合jsp和tiles模板示例

    這篇文章主要介紹了Spring?boot整合jsp模板和tiles模板的示例演示過(guò)程,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-03-03
  • Java Hibernate中的查詢(xún)策略和抓取策略

    Java Hibernate中的查詢(xún)策略和抓取策略

    Hibernate是一種Java對(duì)象關(guān)系映射框架,提供了多種查詢(xún)和抓取策略,用于優(yōu)化數(shù)據(jù)庫(kù)訪問(wèn)性能。查詢(xún)策略包括延遲加載、立即加載、查詢(xún)緩存等,抓取策略包括join抓取、子查詢(xún)抓取、批量抓取等。這些策略可以根據(jù)實(shí)際應(yīng)用場(chǎng)景進(jìn)行選擇和配置,提高數(shù)據(jù)訪問(wèn)的效率和穩(wěn)定性
    2023-04-04
  • java仿windows記事本小程序

    java仿windows記事本小程序

    這篇文章主要為大家詳細(xì)介紹了java仿windows記事本小程序,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2019-03-03
  • Maven高頻配置錯(cuò)誤總結(jié)與解決方案

    Maven高頻配置錯(cuò)誤總結(jié)與解決方案

    在日常 Java 開(kāi)發(fā)中,Maven 幾乎是標(biāo)配構(gòu)建工具,但很多人在寫(xiě) pom.xml、搭建多模塊項(xiàng)目、處理依賴(lài)時(shí),總會(huì)遇到各種莫名其妙的報(bào)錯(cuò),本文把最常見(jiàn)、最容易踩坑的 Maven 配置問(wèn)題整理成一篇完整文章,需要的朋友可以參考下
    2026-03-03
  • Java集合的組內(nèi)平均值的計(jì)算方法總結(jié)

    Java集合的組內(nèi)平均值的計(jì)算方法總結(jié)

    在Java中,經(jīng)常需要對(duì)集合進(jìn)行各種操作,其中之一就是計(jì)算集合的組內(nèi)平均值,本文將介紹如何使用Java集合來(lái)計(jì)算組內(nèi)平均值,并提供一些示例代碼和實(shí)用技巧
    2024-08-08
  • JAVA WEB中Servlet和Servlet容器的區(qū)別

    JAVA WEB中Servlet和Servlet容器的區(qū)別

    這篇文章主要介紹了JAVA WEB中Servlet和Servlet容器的區(qū)別,文中示例代碼非常詳細(xì),供大家參考和學(xué)習(xí),感興趣的朋友可以了解下
    2020-06-06
  • SpringBoot整合EasyExcel實(shí)現(xiàn)批量導(dǎo)入導(dǎo)出

    SpringBoot整合EasyExcel實(shí)現(xiàn)批量導(dǎo)入導(dǎo)出

    這篇文章主要為大家詳細(xì)介紹了SpringBoot整合EasyExcel實(shí)現(xiàn)批量導(dǎo)入導(dǎo)出功能的相關(guān)知識(shí),文中的示例代碼講解詳細(xì),需要的小伙伴可以參考下
    2024-03-03
  • Java中字符串拼接的一些細(xì)節(jié)分析

    Java中字符串拼接的一些細(xì)節(jié)分析

    這篇文章主要介紹了Java中字符串拼接的一些細(xì)節(jié)分析,本文著重剖析了字符串拼接的一些性能問(wèn)題、技巧等內(nèi)容,需要的朋友可以參考下
    2015-01-01
  • java 中設(shè)計(jì)模式(裝飾設(shè)計(jì)模式)的實(shí)例詳解

    java 中設(shè)計(jì)模式(裝飾設(shè)計(jì)模式)的實(shí)例詳解

    這篇文章主要介紹了java 中設(shè)計(jì)模式(裝飾設(shè)計(jì)模式)的實(shí)例詳解的相關(guān)資料,希望通過(guò)本文能幫助到大家,需要的朋友可以參考下
    2017-09-09

最新評(píng)論

民丰县| 英超| 浠水县| 酒泉市| 溆浦县| 修武县| 资中县| 宣武区| 巴中市| 三明市| 平顶山市| 邢台县| 时尚| 郑州市| 勐海县| 赣榆县| 望江县| 左贡县| 于田县| 汪清县| 开化县| 南宫市| 即墨市| 敦煌市| 增城市| 阜新市| 奎屯市| 陵水| 神池县| 紫阳县| 锡林郭勒盟| 平乡县| 华池县| 旺苍县| 宁河县| 博乐市| 荆州市| 通许县| 乌兰察布市| 乳山市| 景东|