C#使用RabbitMQ發(fā)送和接收消息工具類(lèi)的實(shí)現(xiàn)
下面是一個(gè)簡(jiǎn)單的 C# RabbitMQ 發(fā)送和接收消息的封裝工具類(lèi)的示例代碼:
工具類(lèi)
通過(guò)NuGet安裝RabbitMQ.Client
using Newtonsoft.Json;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Channels;
using System.Threading.Tasks;
namespace WorkerService1
{
public class RabbitMQHelper : IDisposable
{
private readonly ConnectionFactory _factory;
private IConnection _connection;
private IModel _channel;
public RabbitMQHelper()
{
// 設(shè)置連接參數(shù)
_factory = new ConnectionFactory() { HostName = "localhost", Port = 5672, UserName = "guest", Password = "guest" };
}
/// <summary>
/// 發(fā)送消息
/// </summary>
/// <typeparam name="T"></typeparam>
/// <param name="queueName"></param>
/// <param name="message"></param>
public void SendMessage<T>(string queueName, T message)
{
try
{
InitConnection();
// 聲明隊(duì)列
_channel.QueueDeclare(queue: queueName,
durable: true,// 設(shè)置為true表示隊(duì)列是持久化的
exclusive: false,
autoDelete: false,
arguments: null);
var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(message));
_channel.BasicPublish(exchange: "", routingKey: queueName, basicProperties: null, body: body);
}
catch (Exception ex)
{
Console.WriteLine("Failed to send message: {0}", ex.Message);
}
}
/// <summary>
/// 接收消息
/// </summary>
/// <typeparam name="T"></typeparam>
/// <param name="queueName"></param>
/// <param name="messageHandler"></param>
public void ReceiveMessage<T>(string queueName, Action<T> messageHandler)
{
try
{
InitConnection();
// 聲明隊(duì)列(接收需聲明隊(duì)列,否則隊(duì)列不存在時(shí),無(wú)法接收消息)
_channel.QueueDeclare(queue: queueName,
durable: true, // 設(shè)置為true表示隊(duì)列是持久化的
exclusive: false,
autoDelete: false,
arguments: null);
//設(shè)置消費(fèi)者數(shù)量(并發(fā)度),每個(gè)消費(fèi)者每次只能處理一條消息
_channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);
// 創(chuàng)建消費(fèi)者
var consumer = new EventingBasicConsumer(_channel);
consumer.Received += (model, ea) =>
{
try
{
var message = Encoding.UTF8.GetString(ea.Body.ToArray());
var convertedMessage = JsonConvert.DeserializeObject<T>(message);
//委托方法
messageHandler.Invoke(convertedMessage);
// 消息處理成功,確認(rèn)消息
_channel.BasicAck(ea.DeliveryTag, false);
}
catch (Exception ex)
{
// 消息處理異常,確認(rèn)消息
_channel.BasicAck(ea.DeliveryTag, false);
}
};
_channel.BasicConsume(queue: queueName,
autoAck: false,// 設(shè)置為true表示自動(dòng)確認(rèn)消息
consumer: consumer);
}
catch (Exception ex)
{
Console.WriteLine("Failed to receive message: {0}", ex.Message);
}
}
/// <summary>
/// 初始化鏈接
/// </summary>
private void InitConnection()
{
if (_connection == null || !_connection.IsOpen)
{
_connection = _factory.CreateConnection();
_channel = _connection.CreateModel();
}
}
/// <summary>
/// 釋放資源
/// </summary>
public void Dispose()
{
_channel?.Close();
_channel?.Dispose();
_connection?.Close();
_connection?.Dispose();
}
}
}
使用示例
using System;
using System.Text;
using System.Threading.Tasks;
using WorkerService1;
public class Program
{
private static string QueueName = "myqueue_key";
public static void Main()
{
var rabbitMQHelper = new RabbitMQHelper();
for (long i = 0; i < 30; i++)
{
rabbitMQHelper.SendMessage(QueueName, i);
}
rabbitMQHelper.ReceiveMessage<long>(QueueName, ReceivedHandle);
Console.ReadLine();
}
/// <summary>
/// 接收處理
/// </summary>
/// <param name="index"></param>
private static void ReceivedHandle(long index)
{
try
{
Console.WriteLine($"第{index}次開(kāi)始{DateTime.Now}");
Thread.Sleep(2000);
Console.WriteLine($"第{index}次結(jié)束{DateTime.Now}");
}
catch (Exception ex)
{
Console.WriteLine(ex.Message);
}
}
}到此這篇關(guān)于C#使用RabbitMQ發(fā)送和接收消息工具類(lèi)的實(shí)現(xiàn)的文章就介紹到這了,更多相關(guān)C# RabbitMQ發(fā)送和接收內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- 在C# .NET中使用RabbitMQ實(shí)現(xiàn)發(fā)布/訂閱模式的方法
- C#使用RabbitMQ的詳細(xì)教程
- C#?RabbitMQ的使用詳解
- C#通過(guò)rabbitmq實(shí)現(xiàn)定時(shí)任務(wù)(延時(shí)隊(duì)列)
- C#用RabbitMQ實(shí)現(xiàn)消息訂閱與發(fā)布
- C#利用RabbitMQ實(shí)現(xiàn)點(diǎn)對(duì)點(diǎn)消息傳輸
- c# rabbitmq 簡(jiǎn)單收發(fā)消息的示例代碼
- C#操作RabbitMQ的完整實(shí)例
- C#中RabbitMQ的使用小結(jié)
相關(guān)文章
C#使用MQTTnet實(shí)現(xiàn)服務(wù)端與客戶端的通訊的示例
本文主要介紹了C#使用MQTTnet實(shí)現(xiàn)服務(wù)端與客戶端的通訊的示例,包括協(xié)議特性、連接管理、QoS機(jī)制和安全策略,具有一定的參考價(jià)值,感興趣的可以了解一下2025-05-05
C#中Thread.CurrentThread的用法小結(jié)
本文主要介紹了C#中Thread.CurrentThread的用法小結(jié),通過(guò)Thread.CurrentThread可以訪問(wèn)和修改當(dāng)前線程的各種屬性和方法,具有一定的參考價(jià)值,感興趣的可以了解一下2025-04-04
C#使用ffmpeg實(shí)現(xiàn)將圖片保存為mp4視頻
FFmpeg是一個(gè)開(kāi)源的跨平臺(tái)多媒體處理工具,它提供了強(qiáng)大的功能,包括頻和視頻編碼、解碼、轉(zhuǎn)碼等,本文我們將使用FFmpeg實(shí)現(xiàn)將圖片保存為mp4視頻,感興趣的可以了解下2024-11-11
C#從文件或標(biāo)準(zhǔn)輸入設(shè)備讀取指定行的方法
這篇文章主要介紹了C#從文件或標(biāo)準(zhǔn)輸入設(shè)備讀取指定行的方法,涉及C#文件及IO操作的相關(guān)技巧,具有一定參考借鑒價(jià)值,需要的朋友可以參考下2015-04-04
c++函數(shù)轉(zhuǎn)c#函數(shù)示例程序分享
這篇文章主要介紹了c++函數(shù)轉(zhuǎn)c#函數(shù)示例程序,大家參考使用吧2013-12-12
extern外部方法使用C#的實(shí)現(xiàn)方法
這篇文章主要介紹了extern外部方法使用C#的實(shí)現(xiàn)方法,較為詳細(xì)的分析了外部方法使用C#的具體步驟與實(shí)現(xiàn)技巧,具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2014-12-12
C#根據(jù)反射和特性實(shí)現(xiàn)ORM映射實(shí)例分析
這篇文章主要介紹了C#根據(jù)反射和特性實(shí)現(xiàn)ORM映射的方法,實(shí)例分析了反射的原理、特性與ORM的實(shí)現(xiàn)技巧,具有一定參考借鑒價(jià)值,需要的朋友可以參考下2015-04-04

