.NET Core實現(xiàn)RabbitMQ消息隊列的示例代碼
RabbitMQ 是一個流行的消息隊列中間件,它允許應用程序通過異步消息的方式進行通信。RabbitMQ 支持 AMQP 協(xié)議,可以通過多種方式與應用程序交互。在本教程中,我們將深入探討如何在 .NET Core 環(huán)境中使用 RabbitMQ 來實現(xiàn)消息隊列。我們將學習如何在生產(chǎn)者端發(fā)送消息,消費者端接收消息,并確保消息的可靠性。
1. 安裝和配置 RabbitMQ
在開始使用 RabbitMQ 之前,首先需要確保你的機器上已經(jīng)安裝并運行 RabbitMQ??梢酝ㄟ^以下方式安裝 RabbitMQ:
使用 Docker 安裝 RabbitMQ
RabbitMQ 提供了官方的 Docker 鏡像,這使得在本地機器上運行 RabbitMQ 非常簡單。
docker pull rabbitmq:management docker run -d -p 5672:5672 -p 15672:15672 rabbitmq:management
5672是 RabbitMQ 的默認消息隊列端口。15672是 RabbitMQ 管理插件的 Web 界面端口。通過瀏覽器訪問http://localhost:15672可以登錄 RabbitMQ 管理界面,默認的用戶名和密碼都是guest。
安裝并啟動 RabbitMQ 后,您可以繼續(xù)進行開發(fā)。
2. 安裝 RabbitMQ 客戶端庫
在 .NET Core 中與 RabbitMQ 進行交互,我們需要使用 RabbitMQ.Client NuGet 包??梢酝ㄟ^以下命令在項目中添加這個依賴:
dotnet add package RabbitMQ.Client
這個庫提供了與 RabbitMQ 服務進行交互所需的所有工具。
3. 創(chuàng)建生產(chǎn)者(Producer)
生產(chǎn)者是負責將消息發(fā)送到 RabbitMQ 的應用程序。它通過連接到 RabbitMQ 服務器、創(chuàng)建一個隊列和交換機,將消息發(fā)布到隊列中。
創(chuàng)建消息生產(chǎn)者代碼
下面是一個基本的生產(chǎn)者示例代碼,展示了如何連接到 RabbitMQ,聲明隊列,并發(fā)送一條簡單的消息:
using RabbitMQ.Client;
using System;
using System.Text;
class Program
{
static void Main(string[] args)
{
// 創(chuàng)建連接工廠
var factory = new ConnectionFactory() { HostName = "localhost" };
// 創(chuàng)建連接和通道
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
// 聲明一個隊列(確保隊列存在)
channel.QueueDeclare(queue: "hello_queue", durable: false, exclusive: false, autoDelete: false, arguments: null);
// 創(chuàng)建消息
string message = "Hello, RabbitMQ!";
var body = Encoding.UTF8.GetBytes(message);
// 發(fā)送消息到隊列
channel.BasicPublish(exchange: "", routingKey: "hello_queue", basicProperties: null, body: body);
Console.WriteLine(" [x] Sent {0}", message);
}
Console.WriteLine(" Press [enter] to exit.");
Console.ReadLine();
}
}
在上面的代碼中:
ConnectionFactory用來創(chuàng)建連接到 RabbitMQ 服務器的連接。QueueDeclare用來聲明一個隊列,確保隊列存在。如果隊列已經(jīng)存在,聲明將被忽略。BasicPublish用來將消息發(fā)送到隊列。
參數(shù)說明:
queue: 隊列的名稱(此例中是hello_queue)。durable: 是否將隊列標記為持久化。如果設置為true,即使 RabbitMQ 重啟,隊列也會存在。exclusive: 是否使隊列只對當前連接可用。autoDelete: 是否在最后一個消費者斷開連接時自動刪除隊列。
4. 創(chuàng)建消費者(Consumer)
消費者從隊列中獲取并處理消息。消費者通常是另一個應用程序,它會連接到 RabbitMQ,并持續(xù)地從隊列中取出消息進行處理。
創(chuàng)建消息消費者代碼
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System;
using System.Text;
class Program
{
static void Main(string[] args)
{
// 創(chuàng)建連接工廠
var factory = new ConnectionFactory() { HostName = "localhost" };
// 創(chuàng)建連接和通道
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
// 聲明隊列,確保消費者能夠連接到相同的隊列
channel.QueueDeclare(queue: "hello_queue", durable: false, exclusive: false, autoDelete: false, arguments: null);
// 創(chuàng)建消費者對象
var consumer = new EventingBasicConsumer(channel);
// 消息處理邏輯
consumer.Received += (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
Console.WriteLine(" [x] Received {0}", message);
};
// 開始消費消息
channel.BasicConsume(queue: "hello_queue", autoAck: true, consumer: consumer);
Console.WriteLine(" Press [enter] to exit.");
Console.ReadLine();
}
}
}
在上面的代碼中:
QueueDeclare用來確保消費者連接到相同的隊列。EventingBasicConsumer是消費者的實現(xiàn),用于異步接收消息。BasicConsume用于開始消費消息,autoAck設置為true,表示自動確認消息。
參數(shù)說明:
autoAck: 如果設置為true,消費者會自動確認消息。如果設置為false,需要手動確認消息。
5. 持久化消息
如果您希望在 RabbitMQ 重啟后保持消息的持久性,可以在生產(chǎn)者和消費者中啟用消息的持久化。
消息持久化設置
在生產(chǎn)者端發(fā)送持久化消息:
// 設置消息持久化 var properties = channel.CreateBasicProperties(); properties.Persistent = true; // 設置消息為持久化 channel.BasicPublish(exchange: "", routingKey: "hello_queue", basicProperties: properties, body: body);
此外,聲明隊列時也需要設置 durable: true,確保隊列本身是持久化的。
channel.QueueDeclare(queue: "hello_queue", durable: true, exclusive: false, autoDelete: false, arguments: null);
6. 消息確認機制
在消息傳遞過程中,為了確保消息被成功消費并避免丟失,可以啟用消息確認機制。在這種情況下,消費者需要顯式確認消息。
啟用手動消息確認
在消費者端禁用自動確認,并手動確認每條已成功處理的消息:
channel.BasicConsume(queue: "hello_queue", autoAck: false, consumer: consumer);
consumer.Received += (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
Console.WriteLine(" [x] Received {0}", message);
// 手動確認消息
channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
};
BasicAck 用于確認消息已經(jīng)被處理。deliveryTag 是消息的標識符,multiple 參數(shù)表示是否確認多個消息。
7. 運行和測試
- 啟動消費者應用程序,確保它可以連接到 RabbitMQ 并等待消息。
- 啟動生產(chǎn)者應用程序,它將發(fā)送消息到 RabbitMQ 隊列。
- 消費者將從隊列中接收到消息,并進行處理。
如果一切配置正確,您將在控制臺中看到生產(chǎn)者發(fā)送的消息以及消費者處理的消息。
8. 總結(jié)
通過本教程,我們學習了如何在 .NET Core 中使用 RabbitMQ 實現(xiàn)一個簡單的消息隊列系統(tǒng)。關鍵步驟包括:
- 安裝 RabbitMQ 客戶端庫。
- 在生產(chǎn)者中聲明隊列并發(fā)送消息。
- 在消費者中聲明隊列并處理消息。
- 配置消息持久化和確認機制,確保消息的可靠性。
RabbitMQ 是一個強大的消息隊列中間件,適用于各種需要解耦和異步通信的應用程序。通過靈活的交換機和隊列配置,您可以實現(xiàn)不同的消息傳遞模式,以滿足不同的業(yè)務需求。
到此這篇關于.NET Core實現(xiàn)RabbitMQ消息隊列的示例代碼的文章就介紹到這了,更多相關.NET Core RabbitMQ消息隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
.NET Core 1.0創(chuàng)建Self-Contained控制臺應用
這篇文章主要為大家詳細介紹了.NET Core 1.0創(chuàng)建Self-Contained控制臺應用的相關資料,具有一定的參考價值,感興趣的小伙伴們可以參考一下2017-04-04
.Net RabbitMQ實現(xiàn)HTTP API接口調(diào)用
RabbitMQ Management插件還提供了基于RESTful風格的HTTP API接口來方便調(diào)用。本文就主要介紹了.Net RabbitMQ實現(xiàn)HTTP API接口調(diào)用,感興趣的可以了解一下2021-06-06
.net core版 文件上傳/ 支持批量上傳拖拽及預覽功能(bootstrap fileinput上傳文件)
本篇內(nèi)容主要解決.net core中文件上傳的問題 開發(fā)環(huán)境:ubuntu+vscode.本文給大家介紹的非常詳細,感興趣的朋友一起看看吧2017-03-03
Entity?Framework?Core關聯(lián)刪除
關聯(lián)刪除通常是一個數(shù)據(jù)庫術語,用于描述在刪除行時允許自動觸發(fā)刪除關聯(lián)行的特征;即當主表的數(shù)據(jù)行被刪除時,自動將關聯(lián)表中依賴的數(shù)據(jù)行進行刪除,或者將外鍵更新為NULL或默認值。本文將為大家具體介紹一下Entity?Framework?Core關聯(lián)刪除,需要的可以參考一下2021-12-12
DropDownList獲取的SelectIndex一直為0的問題
由于初始化判斷出錯導致每次傳到服務器的時候會初始化一次,這就導致每次獲取DropDownList的SelectIndex的時候只能是02014-06-06
asp.net不同頁面間數(shù)據(jù)傳遞的多種方法
這篇文章主要介紹了asp.net不同頁面間數(shù)據(jù)傳遞的多種方法,包括使用QueryString顯式傳遞、頁面對象的屬性、cookie、Cache等9種方法2014-01-01

