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

.NET6+中使用RabbitMQ詳細指南

 更新時間:2026年05月04日 08:45:15   作者:回憶2012初秋  
本文詳細介紹如何在.NET6中使用RabbitMQ實現(xiàn)異步通信,首先創(chuàng)建.NET6應(yīng)用并通過NuGet安裝RabbitMQ客戶端庫,接著配置RabbitMQ選項,創(chuàng)建連接工廠類,然后實現(xiàn)RabbitMQ生產(chǎn)者和消費者,本文介紹的非常詳細,感興趣的朋友一起看看吧

解釋:
    RabbitMQ 是一個流行的開源消息隊列系統(tǒng),廣泛用于實現(xiàn)異步通信、解耦組件、負載均衡等場景。在本篇博客中,我們將詳細介紹如何在 .NET 6 中使用 RabbitMQ,包括生產(chǎn)者和消費者的實現(xiàn),以及如何通過依賴注入來管理它們。

一、創(chuàng)建 .NET 6 應(yīng)用

    首先,確保你已經(jīng)安裝了 .NET 6 SDK??梢允褂妹钚泄ぞ邉?chuàng)建一個新的 .NET 控制臺應(yīng)用:

dotnet new console -n RabbitMQDemo
cd RabbitMQDemo

   接著,你需要安裝 RabbitMQ 的客戶端庫,通過 NuGet 包管理器來安裝 RabbitMQ.Client:

dotnet add package RabbitMQ.Client;       //這里我們選擇 6.4.0 版本
dotnet add package Masuit.Tools.Core;     //這個為一個

二、配置 RabbitMQ

    接下來,我們需要定義連接 RabbitMQ 所需的配置選項。通常我們會將這些配置選項存儲在一個類中。以下是配置類(RabbitMQServiceOptions)的實現(xiàn):

/// <summary>
/// RabbitMQ服務(wù)配置
/// </summary>
public class RabbitMQServiceOptions
{
        /// <summary>
        /// 服務(wù)地址
        /// </summary>
        public string Host { get; set; }
        /// <summary>
        /// 端口
        /// </summary>
        public int Port { get; set; }
        /// <summary>
        /// 用戶名
        /// </summary>
        public string UserName { get; set; }
        /// <summary>
        /// 密碼
        /// </summary>
        public string Password { get; set; }
}

三、實現(xiàn) RabbitMQ 連接工廠

    為確保我們能夠高效且安全地建立 RabbitMQ 連接,我們將創(chuàng)建一個名為 RabbitMQContext 的連接工廠類:

    /// <summary>
    /// RabbitMQ連接工廠
    /// </summary>
    public class RabbitMQContext
    {
        private static ConnectionFactory? factory;
        private static readonly object lockObj = new();
        /// <summary>
        /// 獲取單個RabbitMQ連接
        /// </summary>
        /// <returns></returns>
        public static IConnection GetConnection(string hostName, int port, string userName, string password)
        {
            if (factory == null)
            {
                lock (lockObj)
                {
                    factory ??= new ConnectionFactory
                    {
                        HostName = hostName,
                        Port = port,
                        UserName = userName,
                        Password = password
                    };
                }
            }
            return factory.CreateConnection();
        }
    }

四、實現(xiàn) RabbitMQ 生產(chǎn)者

    接下來我們來實現(xiàn)一個 RabbitMQ 生產(chǎn)者,用于發(fā)送消息到隊列。創(chuàng)建一個名為 RabbitMQProducer 的類:

/// <summary>
 /// RabbitMQ 客戶端,用于發(fā)送消息到 RabbitMQ 隊列或交換機。
 /// </summary>
 public class RabbitMQProducer
 {
     private readonly ILogger<RabbitMQProducer> _logger;
     private readonly RabbitMQServiceOptions _options;
     /// <summary>
     /// 初始化 RabbitMQ 客戶端。
     /// </summary>
     /// <param name="logger">日志記錄器。</param>
     /// <param name="options">RabbitMQ 連接配置。</param>
     public RabbitMQProducer(ILogger<RabbitMQProducer> logger, RabbitMQServiceOptions options)
     {
         _logger = logger;
         _options = options;
     }
     /****
      * RabbitMQ 交換機類型說明:
      * 1. Direct Exchange – 處理路由鍵。需要將一個隊列綁定到交換機上,要求該消息與一個特定的路由鍵完全匹配。
      *    例如,如果一個隊列綁定到該交換機上要求路由鍵 “dog”,則只有被標記為“dog”的消息才被轉(zhuǎn)發(fā)。
      * 2. Fanout Exchange – 不處理路由鍵。只需將隊列綁定到交換機上,發(fā)送到交換機的消息會被轉(zhuǎn)發(fā)到所有綁定的隊列。
      *    類似于廣播,所有綁定的隊列都會收到消息。
      * 3. Topic Exchange – 將路由鍵和某模式進行匹配。隊列需要綁定到一個模式上。
      *    符號“#”匹配一個或多個詞,符號“*”匹配不多不少一個詞。
      *    例如,“audit.#”能夠匹配到“audit.irs.corporate”,但“audit.*”只會匹配到“audit.irs”。
      ****/
     /// <summary>
     /// 發(fā)布消息(工作隊列模式,適用于多消費者負載均衡)。
     /// </summary>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="message">消息內(nèi)容。</param>
     public void WorkQueueSendMessage(string queueName, string message)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到指定隊列
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchange: string.Empty, routingKey: queueName, basicProperties: properties, body: body);
     }
     /// <summary>
     /// 推送消息(簡單模式,適用于單生產(chǎn)者和單消費者)。
     /// </summary>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="message">消息內(nèi)容。</param>
     public void SimpleSendMessage(string queueName, string message)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到指定隊列
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchange: string.Empty, routingKey: queueName, mandatory: false, basicProperties: properties, body: body);
     }
     /// <summary>
     /// 發(fā)布消息(Fanout Exchange 模式,適用于廣播消息到所有綁定隊列)。
     /// </summary>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="message">消息內(nèi)容。</param>
     public void FanoutSendMessage(string queueName, string message)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 創(chuàng)建 Fanout 類型的交換機
         var exchangeName = $"{queueName}_fanout_exchange";
         channel.ExchangeDeclare(exchangeName, ExchangeType.Fanout);
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
         // 將隊列綁定到交換機,routingKey 無需指定
         channel.QueueBind(queueName, exchangeName, routingKey: string.Empty);
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到交換機,所有綁定隊列都會收到消息
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchangeName, routingKey: string.Empty, properties, body: body);
     }
     /// <summary>
     /// 發(fā)布消息(Fanout Exchange 模式,適用于廣播消息到所有綁定隊列)。
     /// </summary>
     /// <param name="exchangeName">交換機名稱。</param>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="message">消息內(nèi)容。</param>
     public void FanoutSendMessage(string exchangeName, string queueName, string message)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 創(chuàng)建 Fanout 類型的交換機
         channel.ExchangeDeclare(exchangeName, ExchangeType.Fanout);
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
         // 將隊列綁定到交換機,routingKey 無需指定
         channel.QueueBind(queueName, exchangeName, routingKey: string.Empty);
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到交換機,所有綁定隊列都會收到消息
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchangeName, routingKey: string.Empty, properties, body: body);
     }
     /// <summary>
     /// 發(fā)布消息(Direct Exchange 模式,適用于路由鍵完全匹配的消息分發(fā))。
     /// </summary>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="message">消息內(nèi)容。</param>
     public void DirectSendMessage(string queueName, string message)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 創(chuàng)建 Direct 類型的交換機
         var exchangeName = $"{queueName}_direct_exchange";
         channel.ExchangeDeclare(exchangeName, ExchangeType.Direct);
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
         // 將隊列綁定到交換機,并指定路由鍵
         var routingKey = $"{queueName}";
         channel.QueueBind(queueName, exchangeName, routingKey);
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到交換機,只有路由鍵完全匹配的隊列才會收到消息
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchangeName, routingKey, properties, body: body);
     }
     /// <summary>
     /// 發(fā)布消息(Direct Exchange 模式,適用于路由鍵完全匹配的消息分發(fā))。
     /// </summary>
     /// <param name="exchangeName">交換機名稱。</param>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="routingKey">路由鍵。</param>
     /// <param name="message">消息內(nèi)容。</param>
     public void DirectSendMessage(string exchangeName, string queueName, string routingKey, string message)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 創(chuàng)建 Direct 類型的交換機
         channel.ExchangeDeclare(exchangeName, ExchangeType.Direct);
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
         // 將隊列綁定到交換機,并指定路由鍵
         channel.QueueBind(queueName, exchangeName, routingKey);
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到交換機,只有路由鍵完全匹配的隊列才會收到消息
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchangeName, routingKey, properties, body: body);
     }
     /// <summary>
     /// 發(fā)布消息(Topic Exchange 模式,適用于路由鍵模式匹配的消息分發(fā))。
     /// </summary>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="routingKey">路由鍵。</param>
     /// <param name="message">消息內(nèi)容。</param>
     /// <param name="bindingKeys">綁定規(guī)則(可選)。</param>
     public void TopicSendMessage(string queueName, string routingKey, string message, IEnumerable<string>? bindingKeys = null)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 創(chuàng)建 Topic 類型的交換機
         var exchangeName = $"{queueName}_topic_exchange";
         channel.ExchangeDeclare(exchangeName, ExchangeType.Topic);
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
         // 將隊列綁定到交換機,并指定綁定規(guī)則
         if (!bindingKeys.IsNullOrEmpty())
         {
             foreach (string bindingKey in bindingKeys)
             {
                 channel.QueueBind(queueName, exchangeName, bindingKey);
             }
         }
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到交換機,只有路由鍵與綁定規(guī)則匹配的隊列才會收到消息
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchangeName, routingKey, properties, body: body);
     }
     /// <summary>
     /// 發(fā)布消息(Topic Exchange 模式,適用于路由鍵模式匹配的消息分發(fā))。
     /// </summary>
     /// <param name="exchangeName">交換機名稱。</param>
     /// <param name="queueName">隊列名稱。</param>
     /// <param name="routingKey">路由鍵。</param>
     /// <param name="message">消息內(nèi)容。</param>
     /// <param name="bindingKeys">綁定規(guī)則(可選)。</param>
     public void TopicSendMessage(string exchangeName, string queueName, string routingKey, string message, IEnumerable<string>? bindingKeys = null)
     {
         using var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
         using var channel = connection.CreateModel();
         // 創(chuàng)建 Topic 類型的交換機
         channel.ExchangeDeclare(exchangeName, ExchangeType.Topic);
         // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
         channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
         // 將隊列綁定到交換機,并指定綁定規(guī)則
         if (!bindingKeys.IsNullOrEmpty())
         {
             foreach (string bindingKey in bindingKeys)
             {
                 channel.QueueBind(queueName, exchangeName, bindingKey);
             }
         }
         // 設(shè)置消息持久化,確保消息在 RabbitMQ 重啟后不會丟失
         var properties = channel.CreateBasicProperties();
         properties.Persistent = true;
         // 發(fā)送消息到交換機,只有路由鍵與綁定規(guī)則匹配的隊列才會收到消息
         var body = Encoding.UTF8.GetBytes(message);
         channel.BasicPublish(exchangeName, routingKey, properties, body: body);
     }
 }

五、實現(xiàn) RabbitMQ 消費者

    接下來我們來實現(xiàn)一個 RabbitMQ 消費者,用于發(fā)送消息到隊列。創(chuàng)建一個名為 RabbitMQConsumer 的類:

   /// <summary>
    /// RabbitMQ 消費者,用于從 RabbitMQ 隊列或交換機中消費消息。
    /// </summary>
    public class RabbitMQConsumer
    {
        private readonly ILogger<RabbitMQConsumer> _logger;
        private readonly RabbitMQServiceOptions _options;
        /// <summary>
        /// 初始化 RabbitMQ 消費者。
        /// </summary>
        /// <param name="logger">日志記錄器。</param>
        /// <param name="options">RabbitMQ 連接配置。</param>
        public RabbitMQConsumer(ILogger<RabbitMQConsumer> logger, RabbitMQServiceOptions options)
        {
            _logger = logger;
            _options = options;
        }
        /// <summary>
        /// 簡單消費者(簡單模式,適用于單生產(chǎn)者和單消費者)。
        /// </summary>
        /// <param name="queueName">隊列名稱。</param>
        /// <param name="handler">消息處理邏輯。</param>
        public void SimpleConsumer(string queueName, Action<string> handler)
        {
            var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
            var channel = connection.CreateModel();
            // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
            channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
            // 創(chuàng)建消費者
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (model, ea) =>
            {
                // 處理消息
                var message = Encoding.UTF8.GetString(ea.Body.ToArray());
                handler(message);
            };
            // 開始消費,自動確認消息
            channel.BasicConsume(queue: queueName, autoAck: true, consumer: consumer);
        }
        /// <summary>
        /// 消費者(工作隊列模式,適用于多消費者負載均衡)。
        /// </summary>
        /// <param name="queueName">隊列名稱。</param>
        /// <param name="handler">消息處理邏輯。</param>
        public void WorkConsumer(string queueName, Func<string, bool> handler)
        {
            var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
            var channel = connection.CreateModel();
            // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
            channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
            // 設(shè)置限流,避免消費者一次性接收過多消息
            channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);
            // 創(chuàng)建消費者
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (model, ea) =>
            {
                // 處理消息
                var message = Encoding.UTF8.GetString(ea.Body.ToArray());
                var result = handler(message);
                // 如果消息處理成功,手動確認消息
                if (result)
                {
                    channel.BasicAck(ea.DeliveryTag, false);
                }
            };
            // 開始消費,手動確認消息
            channel.BasicConsume(queueName, autoAck: false, consumer: consumer);
        }
        /// <summary>
        /// 消費者(發(fā)布/訂閱模式,適用于廣播消息到所有綁定隊列)。
        /// </summary>
        /// <param name="exchangeName">交換機名稱。</param>
        /// <param name="queueName">隊列名稱。</param>
        /// <param name="handler">消息處理邏輯。</param>
        public void PubSubConsumer(string exchangeName, string queueName, Action<string> handler)
        {
            var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
            var channel = connection.CreateModel();
            // 聲明 Fanout 類型的交換機
            channel.ExchangeDeclare(exchange: exchangeName, type: "fanout");
            // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
            var queueNameResult = channel.QueueDeclare(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
            // 將隊列綁定到交換機
            channel.QueueBind(queue: queueName, exchange: exchangeName, routingKey: "");
            // 創(chuàng)建消費者
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (model, ea) =>
            {
                // 處理消息
                var message = Encoding.UTF8.GetString(ea.Body.ToArray());
                handler(message);
                // 手動確認消息
                channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
            };
            // 開始消費,自動確認消息
            channel.BasicConsume(queue: queueName, autoAck: true, consumer: consumer);
        }
        /// <summary>
        /// 消費者(路由模式,適用于路由鍵完全匹配的消息分發(fā))。
        /// </summary>
        /// <param name="exchangeName">交換機名稱。</param>
        /// <param name="queueName">隊列名稱。</param>
        /// <param name="routingKey">路由鍵。</param>
        /// <param name="handler">消息處理邏輯。</param>
        public void RoutingConsumer(string exchangeName, string queueName, string routingKey, Action<string> handler)
        {
            var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
            var channel = connection.CreateModel();
            // 聲明 Direct 類型的交換機
            channel.ExchangeDeclare(exchange: exchangeName, type: "direct");
            // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
            var queueNameResult = channel.QueueDeclare(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
            // 將隊列綁定到交換機,并指定路由鍵
            channel.QueueBind(queue: queueName, exchange: exchangeName, routingKey: routingKey);
            // 創(chuàng)建消費者
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (model, ea) =>
            {
                // 處理消息
                var message = Encoding.UTF8.GetString(ea.Body.ToArray());
                handler(message);
                // 手動確認消息
                channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
            };
            // 開始消費,自動確認消息
            channel.BasicConsume(queue: queueName, autoAck: true, consumer: consumer);
        }
        /// <summary>
        /// 消費者(主題模式,適用于路由鍵模式匹配的消息分發(fā))。
        /// </summary>
        /// <param name="exchangeName">交換機名稱。</param>
        /// <param name="queueName">隊列名稱。</param>
        /// <param name="routingKey">路由鍵。</param>
        /// <param name="handler">消息處理邏輯。</param>
        public void TopicConsumer(string exchangeName, string queueName, string routingKey, Action<string> handler)
        {
            var connection = RabbitMQContext.GetConnection(_options.Host, _options.Port, _options.UserName, _options.Password);
            var channel = connection.CreateModel();
            // 聲明 Topic 類型的交換機
            channel.ExchangeDeclare(exchange: exchangeName, type: ExchangeType.Topic);
            // 聲明隊列(如果不存在則創(chuàng)建),并設(shè)置為持久化
            var queueNameResult = channel.QueueDeclare(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
            // 將隊列綁定到交換機,并指定路由鍵
            channel.QueueBind(queue: queueName, exchange: exchangeName, routingKey: routingKey);
            // 創(chuàng)建消費者
            var consumer = new EventingBasicConsumer(channel);
            consumer.Received += (model, ea) =>
            {
                // 處理消息
                var message = Encoding.UTF8.GetString(ea.Body.ToArray());
                handler(message);
                // 手動確認消息
                channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
            };
            // 開始消費,手動確認消息
            channel.BasicConsume(queue: queueName, autoAck: false, consumer: consumer);
        }
    }

六、整合到.Net6的注入依賴:

    在 .NET 應(yīng)用中,通常我們會使用依賴注入來管理服務(wù)的生命周期。以下是一個擴展類,用于將 RabbitMQ 的生產(chǎn)者和消費者注冊到服務(wù)集合中:

   /// <summary>
    /// RabbitMQ 服務(wù)集合擴展類,用于將 RabbitMQ 客戶端和監(jiān)聽器添加到依賴注入容器。
    /// </summary>
    public static class RabbitMQServiceCollectionExtensions
    {
        /// <summary>
        /// 添加 RabbitMQ 服務(wù)到服務(wù)集合。
        /// </summary>
        /// <param name="services">服務(wù)集合。</param>
        /// <returns>服務(wù)集合。</returns>
        /// <exception cref="ArgumentNullException">當(dāng)配置選項無效時拋出。</exception>
        public static IServiceCollection AddRabbmitMQ(this IServiceCollection services)
        {
            // 從容器中獲取配置的 RabbitMQServiceOptions
            var serviceProvider = services.BuildServiceProvider();
            var serviceOptions = serviceProvider.GetRequiredService<IOptions<RabbitMQServiceOptions>>().Value;
            // 驗證服務(wù)選項是否有效
            if (serviceOptions == null || serviceOptions.Host.IsNullOrEmpty() || serviceOptions.Port < 0)
            {
                throw new ArgumentNullException(nameof(serviceOptions), "RabbitMQ service options must be provided with valid settings.");
            }
            // 注冊 RabbitMQ 客戶端作為單例服務(wù)
            services.AddSingleton(sp => new RabbitMQProducer(
                sp.GetRequiredService<ILogger<RabbitMQProducer>>(),
                serviceOptions));
            // 注冊 RabbitMQ消費端作為單例服務(wù)
            services.AddSingleton(sp => new RabbitMQConsumer(
                sp.GetRequiredService<ILogger<RabbitMQConsumer>>(),
                serviceOptions));
            return services;
        }
    }

七、注入 RabbitMQ到程序(Program):

var builder = WebApplication.CreateBuilder(args);
builder.Services.AddControllers();
// 加載 RabbitMQ 配置
builder.Services.Configure<RabbitMQServiceOptions>(builder.Configuration.GetSection("RabbitMQ"));
// 添加 RabbitMQ 客戶端和監(jiān)聽器
builder.Services.AddRabbmitMQ();
var app = builder.Build();
// Configure the HTTP request pipeline.
app.UseAuthorization();
app.MapControllers();
app.Run();

八、生產(chǎn)者的Demo:

using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using System;
using System.Collections.Generic;
using System.Text;
class Program
{
    static void Main(string[] args)
    {
        // 配置 RabbitMQ 連接選項
        var options = new RabbitMQServiceOptions
        {
            Host = "localhost",
            Port = 5672,
            UserName = "guest",
            Password = "guest"
        };
        // 創(chuàng)建日志記錄器(這里使用控制臺日志)
        using var loggerFactory = LoggerFactory.Create(builder => builder.AddConsole());
        var logger = loggerFactory.CreateLogger<RabbitMQProducer>();
        // 創(chuàng)建 RabbitMQ 生產(chǎn)者
        var producer = new RabbitMQProducer(logger, options);
        // 發(fā)送工作隊列消息
        producer.WorkQueueSendMessage("work_queue", "Hello Work Queue");
        // 發(fā)送簡單隊列消息
        producer.SimpleSendMessage("simple_queue", "Hello Simple Queue");
        // 發(fā)送 Fanout 交換機消息(自動生成交換機名稱)
        producer.FanoutSendMessage("fanout_queue", "Hello Fanout Queue");
        // 發(fā)送 Fanout 交換機消息(指定交換機名稱)
        producer.FanoutSendMessage("my_fanout_exchange", "fanout_queue", "Hello Fanout Queue");
        // 發(fā)送 Direct 交換機消息(自動生成交換機名稱)
        producer.DirectSendMessage("direct_queue", "Hello Direct Queue");
        // 發(fā)送 Direct 交換機消息(指定交換機名稱和路由鍵)
        producer.DirectSendMessage("my_direct_exchange", "direct_queue", "direct_key", "Hello Direct Queue");
        // 發(fā)送 Topic 交換機消息(自動生成交換機名稱)
        producer.TopicSendMessage("topic_queue", "topic.key", "Hello Topic Queue", new List<string> { "topic.*" });
        // 發(fā)送 Topic 交換機消息(指定交換機名稱和路由鍵)
        producer.TopicSendMessage("my_topic_exchange", "topic_queue", "topic.key", "Hello Topic Queue", new List<string> { "topic.*" });
        Console.WriteLine("消息發(fā)送完成!");
    }
}
  • 通過以上示例代碼,你可以輕松地在 .NET 6 中使用 RabbitMQ 實現(xiàn)多種消息發(fā)送模式。每種模式都有其特定的應(yīng)用場景,例如:
    :適用于多消費者負載均衡。
  • 工作隊列模式
  • 簡單模式:適用于單生產(chǎn)者和單消費者。
  • Fanout 交換機模式:適用于廣播消息到所有綁定隊列。
  • Direct 交換機模式:適用于路由鍵完全匹配的消息分發(fā)。
  • Topic 交換機模式:適用于路由鍵模式匹配的消息分發(fā)。

九、消費者Demo:

using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using System;
using System.Text;
class Program
{
    static void Main(string[] args)
    {
        // 配置 RabbitMQ 連接選項
        var options = new RabbitMQServiceOptions
        {
            Host = "localhost",
            Port = 5672,
            UserName = "guest",
            Password = "guest"
        };
        // 創(chuàng)建日志記錄器(這里使用控制臺日志)
        using var loggerFactory = LoggerFactory.Create(builder => builder.AddConsole());
        var logger = loggerFactory.CreateLogger<RabbitMQConsumer>();
        // 創(chuàng)建 RabbitMQ 消費者
        var consumer = new RabbitMQConsumer(logger, options);
        // 簡單消費者
        consumer.SimpleConsumer("simple_queue", message =>
        {
            Console.WriteLine($"接收到簡單隊列消息: {message}");
        });
        // 工作隊列消費者
        consumer.WorkConsumer("work_queue", message =>
        {
            Console.WriteLine($"接收到工作隊列消息: {message}");
            return true; // 處理成功,手動確認消息
        });
        // 發(fā)布/訂閱消費者
        consumer.PubSubConsumer("my_fanout_exchange", "fanout_queue", message =>
        {
            Console.WriteLine($"接收到發(fā)布/訂閱消息: {message}");
        });
        // 路由消費者
        consumer.RoutingConsumer("my_direct_exchange", "direct_queue", "direct_key", message =>
        {
            Console.WriteLine($"接收到路由消息: {message}");
        });
        // 主題消費者
        consumer.TopicConsumer("my_topic_exchange", "topic_queue", "topic.key", message =>
        {
            Console.WriteLine($"接收到主題消息: {message}");
        });
        Console.WriteLine("消費者已啟動,等待接收消息...");
        Console.ReadLine(); // 保持程序運行
    }
}
  • 通過以上示例代碼,你可以輕松地在 .NET 6 中使用 RabbitMQ 實現(xiàn)多種消息消費模式。每種模式都有其特定的應(yīng)用場景,例如:
    :適用于單生產(chǎn)者和單消費者。
  • 簡單消費者
  • 工作隊列消費者:適用于多消費者負載均衡。
  • 發(fā)布/訂閱消費者:適用于廣播消息到所有綁定隊列。
  • 路由消費者:適用于路由鍵完全匹配的消息分發(fā)。
  • 主題消費者:適用于路由鍵模式匹配的消息分發(fā)。

到此這篇關(guān)于.NET6+中使用RabbitMQ詳細指南的文章就介紹到這了,更多相關(guān).net使用rabbitmq內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Visual Studio 2017 (VS 2017)離線安裝包制作方法

    Visual Studio 2017 (VS 2017)離線安裝包制作方法

    這篇文章主要為大家詳細介紹了Visual Studio 2017離線安裝包的制作方法,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-03-03
  • 詳解.NET中使用Redis數(shù)據(jù)庫

    詳解.NET中使用Redis數(shù)據(jù)庫

    Redis是一個用的比較廣泛的Key/Value的內(nèi)存數(shù)據(jù)庫,這篇文章主要介紹了詳解.NET中使用Redis數(shù)據(jù)庫,有興趣的可以了解一下。
    2016-12-12
  • Devexpress中Gridcontrol查找分組

    Devexpress中Gridcontrol查找分組

    本文通過實例代碼給大家介紹了Devexpress中Gridcontrol查找分組的方法,非常不錯,具有一定的參考價誒接價值,需要的朋友一起看看吧
    2018-08-08
  • 總結(jié)十條.NET異常處理建議

    總結(jié)十條.NET異常處理建議

    .NET中從始至終要緊記異常處理的策略:拋出具體的一個異常,而不是只拋出Exception類型的異常,這樣能方便我們捕獲對應(yīng)類型的異常。我們在編寫代碼時要注意考慮到應(yīng)用程序最差的情況;顯示有好的信息,并提供適當(dāng)?shù)墓芾韱T聯(lián)系信息
    2015-11-11
  • 在Code First模式中自動創(chuàng)建Entity模型

    在Code First模式中自動創(chuàng)建Entity模型

    這篇文章介紹了在Code First模式中自動創(chuàng)建Entity模型的方法,文中通過示例代碼介紹的非常詳細。對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-06-06
  • asp.net 通過httpModule計算頁面的執(zhí)行時間

    asp.net 通過httpModule計算頁面的執(zhí)行時間

    有時候我們想檢測一下網(wǎng)頁的執(zhí)行效率。記錄下開始請求時的時間和頁面執(zhí)行完畢后的時間點,這段時間差就是頁面的執(zhí)行時間了。要實現(xiàn)這個功能,通過HttpModule來實現(xiàn)是最方便而且準確的。
    2011-02-02
  • ASP.NET動態(tài)加載用戶控件的實現(xiàn)方法

    ASP.NET動態(tài)加載用戶控件的實現(xiàn)方法

    動態(tài)加載用戶控件的方法,用asp.net的朋友推薦
    2008-10-10
  • 一個簡單MVC5 + EF6示例分享

    一個簡單MVC5 + EF6示例分享

    本文小編跟大家分享了一個簡單MVC5 + EF6示例,感興趣的小伙伴們可以參考一下
    2015-09-09
  • ASP.NET MVC自定義異常過濾器

    ASP.NET MVC自定義異常過濾器

    這篇文章介紹了ASP.NET MVC自定義異常過濾器的方法,文中通過示例代碼介紹的非常詳細。對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-03-03
  • asp.net 預(yù)防SQL注入攻擊之我見

    asp.net 預(yù)防SQL注入攻擊之我見

    說起防止SQL注入攻擊,感覺很郁悶,這么多年了大家一直在討論,也一直在爭論,可是到了現(xiàn)在似乎還是沒有定論。當(dāng)不知道注入原理的時候會覺得很神奇,怎么就被注入了呢?會覺得很難預(yù)防。但是當(dāng)知道了注入原理之后預(yù)防不就是很簡單的事情了嗎?
    2009-11-11

最新評論

嘉荫县| 张家港市| 玉门市| 克什克腾旗| 禄丰县| 金秀| 大名县| 清流县| 新宁县| 噶尔县| 大荔县| 林周县| 甘洛县| 清镇市| 天柱县| 清远市| 武城县| 辽阳市| 桦甸市| 台安县| 望都县| 南涧| 区。| 疏附县| 灵川县| 黄山市| 马公市| 汝南县| 涞源县| 英德市| 宁乡县| 盐城市| 平塘县| 蓬莱市| 奎屯市| 蒙阴县| 八宿县| 古田县| 太白县| 林口县| 禹州市|