.net平臺(tái)的rabbitmq使用封裝demo詳解
前言
RabbitMq大家再熟悉不過(guò),這篇文章主要針對(duì)rabbitmq學(xué)習(xí)后封裝RabbitMQ.Client的一個(gè)分享。文章最后,我會(huì)把封裝組件和demo奉上。
什么是rabbitMQ
RabbitMQ是一個(gè)由erlang開發(fā)的AMQP(Advanced Message Queue 高級(jí)消息隊(duì)列協(xié)議 )的開源實(shí)現(xiàn),能夠?qū)崿F(xiàn)異步消息處理
RabbitMQ是一個(gè)消息代理:它接受和轉(zhuǎn)發(fā)消息。
你可以把它想象成一個(gè)郵局:當(dāng)你把你想要發(fā)布的郵件放在郵箱中時(shí),你可以確定郵差先生最終將郵件發(fā)送給你的收件人。在這個(gè)比喻中,RabbitMQ是郵政信箱,郵局和郵遞員。
RabbitMQ和郵局的主要區(qū)別在于它不處理紙張,而是接受,存儲(chǔ)和轉(zhuǎn)發(fā)二進(jìn)制數(shù)據(jù)塊
優(yōu)點(diǎn):異步消息處理
業(yè)務(wù)解耦(下訂單操作:扣減庫(kù)存、生成訂單、發(fā)紅包、發(fā)短信),將下單操作主流程:扣減庫(kù)存、生成訂單,然后通過(guò)MQ消息隊(duì)列完成通知,發(fā)紅包、發(fā)短信,錯(cuò)峰流控 (通知量 消息量 訂單量大的情況實(shí)現(xiàn)MQ消息隊(duì)列機(jī)制,淡季情況下訪問(wèn)量會(huì)少)
靈活的路由(Flexible Routing)
在消息進(jìn)入隊(duì)列之前,通過(guò) Exchange 來(lái)路由消息的。對(duì)于典型的路由功能,RabbitMQ 已經(jīng)提供了一些內(nèi)置的 Exchange 來(lái)實(shí)現(xiàn)。針對(duì)更復(fù)雜的路由功能,可以將多個(gè) Exchange 綁定在一起,也通過(guò)插件機(jī)制實(shí)現(xiàn)自己的 Exchange 。
RabbitMQ網(wǎng)站端口號(hào):15672
程序里面實(shí)現(xiàn)的端口為:5672
Rabbitmq的關(guān)鍵術(shù)語(yǔ)
1、綁定器(Binding):根據(jù)路由規(guī)則綁定Queue和Exchange。
2、路由鍵(Routing Key):Exchange根據(jù)關(guān)鍵字進(jìn)行消息投遞。
3、交換機(jī)(Exchange):指定消息按照路由規(guī)則進(jìn)入指定隊(duì)列
4、消息隊(duì)列(Queue):消息的存儲(chǔ)載體
5、生產(chǎn)者(Producer):消息發(fā)布者。
6、消費(fèi)者(Consumer):消息接收者。
Rabbitmq的運(yùn)作
從下圖可以看出,發(fā)布者(Publisher)是把消息先發(fā)送到交換器(Exchange),再?gòu)慕粨Q器發(fā)送到指定隊(duì)列(Queue),而先前已經(jīng)聲明交換器與隊(duì)列綁定關(guān)系,最后消費(fèi)者(Customer)通過(guò)訂閱或者主動(dòng)取指定隊(duì)列消息進(jìn)行消費(fèi)。

那么剛剛提到的訂閱和主動(dòng)取可以理解成,推(被動(dòng)),拉(主動(dòng))。
推,只要隊(duì)列增加一條消息,就會(huì)通知空閑的消費(fèi)者進(jìn)行消費(fèi)。(我不找你,就等你找我,觀察者模式)
拉,不會(huì)通知消費(fèi)者,而是由消費(fèi)者主動(dòng)輪循或者定時(shí)去取隊(duì)列消息。(我需要才去找你)
使用場(chǎng)景我舉個(gè)例子,假如有兩套系統(tǒng) 訂單系統(tǒng)和發(fā)貨系統(tǒng),從訂單系統(tǒng)發(fā)起發(fā)貨消息指令,為了及時(shí)發(fā)貨,發(fā)貨系統(tǒng)需要訂閱隊(duì)列,只要有指令就處理。
可是程序偶爾會(huì)出異常,例如網(wǎng)絡(luò)或者DB超時(shí)了,把消息丟到失敗隊(duì)列,這個(gè)時(shí)候需要重發(fā)機(jī)制。但是我又不想while(IsPostSuccess == True),因?yàn)橹灰霎惓A?,?huì)在某個(gè)時(shí)間段內(nèi)都會(huì)有異常,這樣的重試是沒(méi)意義的。
這個(gè)時(shí)候不需要及時(shí)的去處理消息,有個(gè)JOB定時(shí)或者每隔幾分鐘(失敗次數(shù)*間隔分鐘)去取失敗隊(duì)列消息,進(jìn)行重發(fā)。
Publish(發(fā)布)的封裝
步驟:初始化鏈接->聲明交換器->聲明隊(duì)列->換機(jī)器與隊(duì)列綁定->發(fā)布消息。注意的是,我將Model存到了ConcurrentDictionary里面,因?yàn)槁暶髋c綁定是非常耗時(shí)的,其次,往重復(fù)的隊(duì)列發(fā)送消息是不需要重新初始化的。
/// <summary>
/// 交換器聲明
/// </summary>
/// <param name="iModel"></param>
/// <param name="exchange">交換器</param>
/// <param name="type">交換器類型:
/// 1、Direct Exchange – 處理路由鍵。需要將一個(gè)隊(duì)列綁定到交換機(jī)上,要求該消息與一個(gè)特定的路由鍵完全
/// 匹配。這是一個(gè)完整的匹配。如果一個(gè)隊(duì)列綁定到該交換機(jī)上要求路由鍵 “dog”,則只有被標(biāo)記為“dog”的
/// 消息才被轉(zhuǎn)發(fā),不會(huì)轉(zhuǎn)發(fā)dog.puppy,也不會(huì)轉(zhuǎn)發(fā)dog.guard,只會(huì)轉(zhuǎn)發(fā)dog
/// 2、Fanout Exchange – 不處理路由鍵。你只需要簡(jiǎn)單的將隊(duì)列綁定到交換機(jī)上。一個(gè)發(fā)送到交換機(jī)的消息都
/// 會(huì)被轉(zhuǎn)發(fā)到與該交換機(jī)綁定的所有隊(duì)列上。很像子網(wǎng)廣播,每臺(tái)子網(wǎng)內(nèi)的主機(jī)都獲得了一份復(fù)制的消息。Fanout
/// 交換機(jī)轉(zhuǎn)發(fā)消息是最快的。
/// 3、Topic Exchange – 將路由鍵和某模式進(jìn)行匹配。此時(shí)隊(duì)列需要綁定要一個(gè)模式上。符號(hào)“#”匹配一個(gè)或多
/// 個(gè)詞,符號(hào)“*”匹配不多不少一個(gè)詞。因此“audit.#”能夠匹配到“audit.irs.corporate”,但是“audit.*”
/// 只會(huì)匹配到“audit.irs”。</param>
/// <param name="durable">持久化</param>
/// <param name="autoDelete">自動(dòng)刪除</param>
/// <param name="arguments">參數(shù)</param>
private static void ExchangeDeclare(IModel iModel, string exchange, string type = ExchangeType.Direct,
bool durable = true,
bool autoDelete = false, IDictionary<string, object> arguments = null)
{
exchange = exchange.IsNullOrWhiteSpace() ? "" : exchange.Trim();
iModel.ExchangeDeclare(exchange, type, durable, autoDelete, arguments);
}
/// <summary>
/// 隊(duì)列聲明
/// </summary>
/// <param name="channel"></param>
/// <param name="queue">隊(duì)列</param>
/// <param name="durable">持久化</param>
/// <param name="exclusive">排他隊(duì)列,如果一個(gè)隊(duì)列被聲明為排他隊(duì)列,該隊(duì)列僅對(duì)首次聲明它的連接可見,
/// 并在連接斷開時(shí)自動(dòng)刪除。這里需要注意三點(diǎn):其一,排他隊(duì)列是基于連接可見的,同一連接的不同信道是可
/// 以同時(shí)訪問(wèn)同一個(gè)連接創(chuàng)建的排他隊(duì)列的。其二,“首次”,如果一個(gè)連接已經(jīng)聲明了一個(gè)排他隊(duì)列,其他連
/// 接是不允許建立同名的排他隊(duì)列的,這個(gè)與普通隊(duì)列不同。其三,即使該隊(duì)列是持久化的,一旦連接關(guān)閉或者
/// 客戶端退出,該排他隊(duì)列都會(huì)被自動(dòng)刪除的。這種隊(duì)列適用于只限于一個(gè)客戶端發(fā)送讀取消息的應(yīng)用場(chǎng)景。</param>
/// <param name="autoDelete">自動(dòng)刪除</param>
/// <param name="arguments">參數(shù)</param>
private static void QueueDeclare(IModel channel, string queue, bool durable = true, bool exclusive = false,
bool autoDelete = false, IDictionary<string, object> arguments = null)
{
queue = queue.IsNullOrWhiteSpace() ? "UndefinedQueueName" : queue.Trim();
channel.QueueDeclare(queue, durable, exclusive, autoDelete, arguments);
}
/// <summary>
/// 獲取Model
/// </summary>
/// <param name="exchange">交換機(jī)名稱</param>
/// <param name="queue">隊(duì)列名稱</param>
/// <param name="routingKey"></param>
/// <param name="isProperties">是否持久化</param>
/// <returns></returns>
private static IModel GetModel(string exchange, string queue, string routingKey, bool isProperties = false)
{
return ModelDic.GetOrAdd(queue, key =>
{
var model = _conn.CreateModel();
ExchangeDeclare(model, exchange, ExchangeType.Fanout, isProperties);
QueueDeclare(model, queue, isProperties);
model.QueueBind(queue, exchange, routingKey);
ModelDic[queue] = model;
return model;
});
}
/// <summary>
/// 發(fā)布消息
/// </summary>
/// <param name="routingKey">路由鍵</param>
/// <param name="body">隊(duì)列信息</param>
/// <param name="exchange">交換機(jī)名稱</param>
/// <param name="queue">隊(duì)列名</param>
/// <param name="isProperties">是否持久化</param>
/// <returns></returns>
public void Publish(string exchange, string queue, string routingKey, string body, bool isProperties = false)
{
var channel = GetModel(exchange, queue, routingKey, isProperties);
try
{
channel.BasicPublish(exchange, routingKey, null, body.SerializeUtf8());
}
catch (Exception ex)
{
throw ex.GetInnestException();
}
}
下次是本機(jī)測(cè)試的發(fā)布速度截圖:

4.2W/S屬于穩(wěn)定速度,把反序列化(ToJson)會(huì)稍微快一些。
Subscribe(訂閱)的封裝
發(fā)布的時(shí)候是申明了交換器和隊(duì)列并綁定,然而訂閱的時(shí)候只需要聲明隊(duì)列就可。從下面代碼能看到,捕獲到異常的時(shí)候,會(huì)把消息送到自定義的“死信隊(duì)列”里,由另外的JOB進(jìn)行定時(shí)重發(fā),因此,finally是應(yīng)答成功的。
/// <summary>
/// 獲取Model
/// </summary>
/// <param name="queue">隊(duì)列名稱</param>
/// <param name="isProperties"></param>
/// <returns></returns>
private static IModel GetModel(string queue, bool isProperties = false)
{
return ModelDic.GetOrAdd(queue, value =>
{
var model = _conn.CreateModel();
QueueDeclare(model, queue, isProperties);
//每次消費(fèi)的消息數(shù)
model.BasicQos(0, 1, false);
ModelDic[queue] = model;
return model;
});
}
/// <summary>
/// 接收消息
/// </summary>
/// <typeparam name="T"></typeparam>
/// <param name="queue">隊(duì)列名稱</param>
/// <param name="isProperties"></param>
/// <param name="handler">消費(fèi)處理</param>
/// <param name="isDeadLetter"></param>
public void Subscribe<T>(string queue, bool isProperties, Action<T> handler, bool isDeadLetter) where T : class
{
//隊(duì)列聲明
var channel = GetModel(queue, isProperties);
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body;
var msgStr = body.DeserializeUtf8();
var msg = msgStr.FromJson<T>();
try
{
handler(msg);
}
catch (Exception ex)
{
ex.GetInnestException().WriteToFile("隊(duì)列接收消息", "RabbitMq");
if (!isDeadLetter)
PublishToDead<DeadLetterQueue>(queue, msgStr, ex);
}
finally
{
channel.BasicAck(ea.DeliveryTag, false);
}
};
channel.BasicConsume(queue, false, consumer);
}
下次是本機(jī)測(cè)試的發(fā)布速度截圖:

快的時(shí)候有1.9K/S,慢的時(shí)候也有1.7K/S
Pull(拉)的封裝
直接上代碼:
/// <summary>
/// 獲取消息
/// </summary>
/// <typeparam name="T"></typeparam>
/// <param name="exchange"></param>
/// <param name="queue"></param>
/// <param name="routingKey"></param>
/// <param name="handler">消費(fèi)處理</param>
private void Poll<T>(string exchange, string queue, string routingKey, Action<T> handler) where T : class
{
var channel = GetModel(exchange, queue, routingKey);
var result = channel.BasicGet(queue, false);
if (result.IsNull())
return;
var msg = result.Body.DeserializeUtf8().FromJson<T>();
try
{
handler(msg);
}
catch (Exception ex)
{
ex.GetInnestException().WriteToFile("隊(duì)列接收消息", "RabbitMq");
}
finally
{
channel.BasicAck(result.DeliveryTag, false);
}
}

快的時(shí)候有1.8K/s,穩(wěn)定是1.5K/S
Rpc(遠(yuǎn)程調(diào)用)的封裝
首先說(shuō)明下,RabbitMq只是提供了這個(gè)RPC的功能,但是并不是真正的RPC,為什么這么說(shuō):
1、傳統(tǒng)Rpc隱藏了調(diào)用細(xì)節(jié),像調(diào)用本地方法一樣傳參、拋出異常
2、RabbitMq的Rpc是基于消息的,消費(fèi)者消費(fèi)后,通過(guò)新隊(duì)列返回響應(yīng)結(jié)果。
/// <summary>
/// RPC客戶端
/// </summary>
/// <param name="exchange"></param>
/// <param name="queue"></param>
/// <param name="routingKey"></param>
/// <param name="body"></param>
/// <param name="isProperties"></param>
/// <returns></returns>
public string RpcClient(string exchange, string queue, string routingKey, string body, bool isProperties = false)
{
var channel = GetModel(exchange, queue, routingKey, isProperties);
var consumer = new QueueingBasicConsumer(channel);
channel.BasicConsume(queue, true, consumer);
try
{
var correlationId = Guid.NewGuid().ToString();
var basicProperties = channel.CreateBasicProperties();
basicProperties.ReplyTo = queue;
basicProperties.CorrelationId = correlationId;
channel.BasicPublish(exchange, routingKey, basicProperties, body.SerializeUtf8());
var sw = Stopwatch.StartNew();
while (true)
{
var ea = consumer.Queue.Dequeue();
if (ea.BasicProperties.CorrelationId == correlationId)
{
return ea.Body.DeserializeUtf8();
}
if (sw.ElapsedMilliseconds > 30000)
throw new Exception("等待響應(yīng)超時(shí)");
}
}
catch (Exception ex)
{
throw ex.GetInnestException();
}
}
/// <summary>
/// RPC服務(wù)端
/// </summary>
/// <typeparam name="T"></typeparam>
/// <param name="exchange"></param>
/// <param name="queue"></param>
/// <param name="isProperties"></param>
/// <param name="handler"></param>
/// <param name="isDeadLetter"></param>
public void RpcService<T>(string exchange, string queue, bool isProperties, Func<T, T> handler, bool isDeadLetter)
{
//隊(duì)列聲明
var channel = GetModel(queue, isProperties);
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body;
var msgStr = body.DeserializeUtf8();
var msg = msgStr.FromJson<T>();
var props = ea.BasicProperties;
var replyProps = channel.CreateBasicProperties();
replyProps.CorrelationId = props.CorrelationId;
try
{
msg = handler(msg);
}
catch (Exception ex)
{
ex.GetInnestException().WriteToFile("隊(duì)列接收消息", "RabbitMq");
}
finally
{
channel.BasicPublish(exchange, props.ReplyTo, replyProps, msg.ToJson().SerializeUtf8());
channel.BasicAck(ea.DeliveryTag, false);
}
};
channel.BasicConsume(queue, false, consumer);
}
可以用,但不建議去用。可以考慮其他的RPC框架。grpc、thrift等。
結(jié)尾
本篇文章,沒(méi)有過(guò)多的寫RabbitMq的知識(shí)點(diǎn),因?yàn)閳@子的學(xué)習(xí)筆記實(shí)在太多了。下面把我的代碼奉上 https://github.com/SkyChenSky/Sikiro.Mq.Rabbit。如果有發(fā)現(xiàn)寫得不對(duì)的地方麻煩在評(píng)論指出,我會(huì)及時(shí)修改以免誤導(dǎo)別人。
到此這篇關(guān)于.net平臺(tái)的rabbitmq使用封裝的文章就介紹到這了,更多相關(guān).net使用rabbitmq內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
ASP.NET實(shí)現(xiàn)數(shù)據(jù)的添加(第10節(jié))
這篇文章主要介紹了ASP.NET如何實(shí)現(xiàn)數(shù)據(jù)的添加,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2015-08-08
.net c# gif動(dòng)畫如何添加圖片水印實(shí)現(xiàn)思路及代碼
本文將詳細(xì)介紹下c#實(shí)現(xiàn)gif動(dòng)畫添加圖片水印,思路很清晰,感興趣的你可以參考下哈,希望可以幫助到你2013-03-03
asp.net中的check與uncheck關(guān)鍵字用法解析
這篇文章主要介紹了asp.net中的check與uncheck關(guān)鍵字用法,以實(shí)例形式較為詳細(xì)的分析了check與uncheck關(guān)鍵字的各種常見用法與使用時(shí)的注意事項(xiàng),非常具有實(shí)用價(jià)值,需要的朋友可以參考下2014-10-10
基于Cookie使用過(guò)濾器實(shí)現(xiàn)客戶每次訪問(wèn)只登錄一次
這篇文章主要介紹了基于Cookie使用過(guò)濾器實(shí)現(xiàn)客戶每次訪問(wèn)只登錄一次,需要的朋友可以參考下2017-06-06
如何使用ASP.NET制作簡(jiǎn)單的驗(yàn)證碼
當(dāng)用戶進(jìn)行注冊(cè)、登陸的時(shí)候都會(huì)遇到輸入驗(yàn)證碼的情況,那驗(yàn)證碼到底是怎么產(chǎn)生的吶,本文就是介紹了如何使用ASP.NET制作簡(jiǎn)單的驗(yàn)證碼,感興趣的朋友可以參考一下2015-07-07
ASP.NET MVC小結(jié)之基礎(chǔ)篇(一)
本文是ASP.NET MVC系列的第一篇文章,跟其他學(xué)習(xí)系列一樣,咱們先來(lái)點(diǎn)基礎(chǔ)知識(shí),之后再循序漸進(jìn)。我們先從asp.net mvc的概念開始吧。2014-11-11
ASP.NET?Core?6框架揭秘實(shí)例演示之如何承載你的后臺(tái)服務(wù)
這篇文章主要介紹了ASP.NET?Core?6框架揭秘實(shí)例演示之如何承載你的后臺(tái)服務(wù),主要包括利用承載服務(wù)收集性能指標(biāo)、依賴注入的應(yīng)用、配置選項(xiàng)的應(yīng)用等知識(shí)點(diǎn),本文給大家介紹的非常詳細(xì),需要的朋友可以參考下2022-03-03
asp.net使用DataTable構(gòu)造Json字符串的方法
這篇文章主要介紹了asp.net使用DataTable構(gòu)造Json字符串的方法,涉及asp.net字符串序列化、遍歷及構(gòu)造等操作技巧,需要的朋友可以參考下2015-12-12

