.NET6+中使用RabbitMQ詳細指南
解釋:
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離線安裝包的制作方法,具有一定的參考價值,感興趣的小伙伴們可以參考一下2017-03-03
在Code First模式中自動創(chuàng)建Entity模型
這篇文章介紹了在Code First模式中自動創(chuàng)建Entity模型的方法,文中通過示例代碼介紹的非常詳細。對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下2022-06-06
asp.net 通過httpModule計算頁面的執(zhí)行時間
有時候我們想檢測一下網(wǎng)頁的執(zhí)行效率。記錄下開始請求時的時間和頁面執(zhí)行完畢后的時間點,這段時間差就是頁面的執(zhí)行時間了。要實現(xiàn)這個功能,通過HttpModule來實現(xiàn)是最方便而且準確的。2011-02-02
ASP.NET動態(tài)加載用戶控件的實現(xiàn)方法
動態(tài)加載用戶控件的方法,用asp.net的朋友推薦2008-10-10

