Maomi.MQ?功能強(qiáng)大的?.NET?RabbitMQ?消息隊(duì)列通訊模型框架詳解
快速開始
在本篇教程中,將介紹 Maomi.MQ.RabbitMQ 的使用方法,以便讀者能夠快速了解該框架的使用方式和特點(diǎn)。
Maomi.MQ.RabbitMQ 是一個基于 RabbitMQ 的消息隊(duì)列封裝框架,提供了很多開箱即用的功能,通過簡單靈活的方式簡化消息傳輸流程,提供一系列可靠的消息傳輸保障機(jī)制,降低開發(fā)者使用難度,減少開發(fā)時間。
主要功能:
- 簡化的消息定義和消費(fèi)者,無需復(fù)雜配置即可上手。
- 豐富靈活的配置,發(fā)揮 RabbitMQ 的強(qiáng)大力量,自動創(chuàng)建隊(duì)列、死信隊(duì)列、廣播模式、Qos 并發(fā)控制、動態(tài)路由、動態(tài)消費(fèi)者。
- 自定義序列化器,支持 Json、Protobuf、Thrift、MessagePack 二進(jìn)制消息傳輸,高性能壓縮消息減少內(nèi)存、提高并發(fā)量、跨微服務(wù)傳輸。
- 支持自動和自定義發(fā)布消息推送交換器,支持 RabbitMQ 事務(wù)推送模式。
- 簡化的消費(fèi)者模式、事件總線模式、廣播模式、動態(tài)消費(fèi)者,還支持結(jié)合 MeditaR、FastEndpoints 框架使用,進(jìn)一步減少使用負(fù)擔(dān),減少代碼侵入性。
- 自定義重試策略,從容應(yīng)對服務(wù)錯誤、強(qiáng)一致性消息、高并發(fā)流量。
- 支持本地消息表模式,強(qiáng)一致性保證業(yè)務(wù)消息不丟失,解決空懸掛、重復(fù)處理等異步消息異常問題,保證業(yè)務(wù)可靠性。
- 自由的消費(fèi)者模式,除了 MeditaR、FastEndpoints,還可以自由接入其它框架,充分利用第三方框架的優(yōu)質(zhì)能力。
快速配置
創(chuàng)建一個 Web 項(xiàng)目(可參考 WebDemo 項(xiàng)目),引入 Maomi.MQ.RabbitMQ 包,在 Web 配置中注入服務(wù):
// using Maomi.MQ;
// using RabbitMQ.Client;
builder.Services.AddMaomiMQ((MqOptionsBuilder options) =>
{
options.WorkId = 1;
options.AppName = "myapp";
options.Rabbit = (ConnectionFactory options) =>
{
options.HostName = Environment.GetEnvironmentVariable("RABBITMQ")!;
options.Port = 5672;
};
}, [typeof(Program).Assembly]);
var app = builder.Build();- WorkId: 指定用于生成分布式雪花 id 的節(jié)點(diǎn) id,默認(rèn)為 0。
每條消息生成一個唯一的 id,便于追蹤。如果不設(shè)置雪花id,在分布式服務(wù)中,多實(shí)例并行工作時,可能會產(chǎn)生相同的 id。
- AppName:用于標(biāo)識在日志和鏈路追蹤中標(biāo)識消息的生產(chǎn)者或消費(fèi)者。
- Rabbit:RabbitMQ 客戶端配置,請參考 ConnectionFactory。
定義消息模型類,模型類是 MQ 通訊的消息基礎(chǔ),該模型類將會被序列化為二進(jìn)制內(nèi)容傳遞到 RabbitMQ 服務(wù)器中。
public class TestEvent
{
public int Id { get; set; }
public override string ToString()
{
return Id.ToString();
}
}定義消費(fèi)者,消費(fèi)者需要實(shí)現(xiàn) IConsumer<TEvent> 接口,以及使用 [Consumer] 特性注解配置消費(fèi)者屬性,如下所示,[Consumer("test")] 表示該消費(fèi)者訂閱的隊(duì)列名稱是 test。
IConsumer<TEvent> 接口有三個方法,ExecuteAsync 方法用于處理消息,FaildAsync 會在 ExecuteAsync 異常時立即執(zhí)行,如果代碼一直異常,最終會調(diào)用 FallbackAsync 方法,Maomi.MQ 框架會根據(jù) ConsumerState 值確定是否將消息放回隊(duì)列重新消費(fèi),或者做其它處理動作。
[Consumer("test")]
public class MyConsumer : IConsumer<TestEvent>
{
// 消費(fèi)
public async Task ExecuteAsync(MessageHeader messageHeader, TestEvent message)
{
Console.WriteLine($"事件 id: {message.Id} {DateTime.Now}");
await Task.CompletedTask;
}
// 每次消費(fèi)失敗時執(zhí)行
public Task FaildAsync(MessageHeader messageHeader, Exception ex, int retryCount, TestEvent message)
=> Task.CompletedTask;
// 補(bǔ)償
public Task<ConsumerState> FallbackAsync(MessageHeader messageHeader, TestEvent? message, Exception? ex)
=> Task.FromResult( ConsumerState.Ack);
}Maomi.MQ 還具有多種消費(fèi)者模式,代碼寫法不一樣,后續(xù)會詳細(xì)講解不同的消費(fèi)者模式。
如果要發(fā)布消息,只需要注入 IMessagePublisher 服務(wù)即可。
private readonly IMessagePublisher _messagePublisher;
public IndexController(IMessagePublisher messagePublisher)
{
_messagePublisher = messagePublisher;
}
[HttpGet("publish")]
public async Task<string> Publisher()
{
// 發(fā)布消息
await _messagePublisher.PublishAsync(exchange: string.Empty, routingKey: "test", message: new TestEvent
{
Id = 123
});
return "ok";
}啟動 Web 服務(wù),在 swagger 頁面上請求 API 接口,MyConsumer 服務(wù)會立即接收到發(fā)布的消息。

就是這么簡單,就是這么方便。
怎么發(fā)布消息
自動發(fā)布
雖然發(fā)布者和消費(fèi)者共用一個模型類,但是在一個項(xiàng)目中怎么配置模型類,都不會影響消費(fèi)者。將分布者消費(fèi)者隔離簡化框架設(shè)計(jì)和微服務(wù)解耦,支持不同的編程語言服務(wù)相互通訊,共同完成業(yè)務(wù)邏輯。
例如為了簡化消息發(fā)布,我們可以在模型類指定綁定的路由鍵。
[RouterKey("scenario.quickstart")]
public sealed class QuickStartMessage
{
}
發(fā)布時是根據(jù)模型類的 [RouterKey] 自動找到要推送的交換器和路由鍵,簡化發(fā)布消息的參數(shù)。
await _publisher.AutoPublishAsync(message);
當(dāng)然手動設(shè)置也可以:
await _publisher.PublishAsync(string.Empty, request.Queue, message);
消費(fèi)者可以自由設(shè)置要消費(fèi)的隊(duì)列,即使是相同的模型類,可以自由設(shè)定隊(duì)列名稱,不會產(chǎn)生干擾。
[Consumer("scenario.quickstart")]
public sealed class QuickStartConsumer : IConsumer<QuickStartMessage>
{
}
手動發(fā)布
適用于你需要精細(xì)控制 exchange/routingKey 或消息屬性(TTL、Header、優(yōu)先級等)的場景。
[HttpPost("publish-manual")]
public async Task<IResult> PublishManual()
{
var message = new OrderCreatedMessage
{
OrderNo = "SO-20260211-002",
Amount = 299.00m
};
await _publisher.PublishAsync(
exchange: "biz.order.exchange",
routingKey: "order.created.v1",
message: message,
properties: p =>
{
p.Expiration = "60000";
p.Headers ??= new Dictionary<string, object?>();
p.Headers["tenant"] = "tenant-a";
});
return Results.Ok(message);
}但是 Maomi.MQ.RabbitMQ 提供了更為簡單易用的方式,實(shí)現(xiàn)自動處理隊(duì)列屬性。
例如,要設(shè)計(jì)一個死信隊(duì)列,你只需要在消費(fèi)者上設(shè)置屬性即可:
[Consumer(
"example.retry.main",
RetryFaildRequeue = false,
DeadExchange = "",
DeadRoutingKey = "example.retry.dead")]
public sealed class RetryConsumer : IConsumer<RetryMessage>
{
}
RetryFaildRequeue表示消費(fèi)失敗后不會放回原隊(duì)列。DeadExchange死信交換器名稱。DeadRoutingKey死信隊(duì)列路由鍵。
所以,發(fā)布消息時,只需要使用這兩行代碼:
await _publisher.PublishAsync(string.Empty, request.Queue, message); await _publisher.AutoPublishAsync(message);
普通消費(fèi)者、動態(tài)消費(fèi)者、事件總線
Maomi.MQ 支持以多種姿勢創(chuàng)建消費(fèi)者,一個項(xiàng)目里面可以同時使用多種消費(fèi)者模式,自由靈活而不會沖突。
這三種模式不是互斥關(guān)系,而是處理問題的方式不同:
- 普通消費(fèi)者:最基礎(chǔ)的消費(fèi)模式,繼承
IConsumer<TMessage>即可使用,還可以自由擴(kuò)展出不同類型的消費(fèi)者模式,例如 MediatR 。 - 動態(tài)消費(fèi)者:也是繼承
IConsumer<TMessage>, 運(yùn)行時動態(tài)創(chuàng)建/停止訂閱。 - 事件總線:一個消息觸發(fā)多個有順序的處理步驟,生成執(zhí)行鏈和回滾鏈路。
普通消費(fèi)者模式
實(shí)現(xiàn) IConsumer<TMessage> 接口即可,可以自定義消費(fèi)、重試、回滾邏輯,這種模式使用簡單,還能從容處理消息錯誤和回滾。
正常處理消息會調(diào)用 ExecuteAsync,失敗后調(diào)用 FaildAsync,重試耗盡進(jìn)入 FallbackAsync。
[Consumer("biz.order.created.v1", Qos = 10, RetryFaildRequeue = false)]
public sealed class NormalOrderConsumer : IConsumer<OrderCreatedMessage>
{
public Task ExecuteAsync(MessageHeader messageHeader, OrderCreatedMessage message)
{
Console.WriteLine($"normal consume => {message.OrderNo}");
return Task.CompletedTask;
}
public Task FaildAsync(MessageHeader messageHeader, Exception ex, int retryCount, OrderCreatedMessage message)
=> Task.CompletedTask;
public Task<ConsumerState> FallbackAsync(MessageHeader messageHeader, OrderCreatedMessage? message, Exception? ex)
=> Task.FromResult(ConsumerState.Ack);
}動態(tài)消費(fèi)者模式
你可以在運(yùn)行時將業(yè)務(wù)需要訂閱的消息隊(duì)列動態(tài)注冊,例如 SAAS 平臺新建租戶后需要動態(tài)按照租戶前綴消費(fèi)對應(yīng)的主題,也可以自行取消訂閱。
public sealed class DynamicDemoService
{
private readonly IDynamicConsumer _dynamicConsumer;
public DynamicDemoService(IDynamicConsumer dynamicConsumer)
{
_dynamicConsumer = dynamicConsumer;
}
public async Task<string> StartAsync(string queue)
{
var options = new ConsumerAttribute(queue) { Qos = 5 };
var consumerTag = await _dynamicConsumer.ConsumerAsync<OrderCreatedMessage>(
options,
execute: (header, message) =>
{
Console.WriteLine($"dynamic consume => {message.OrderNo}");
return Task.CompletedTask;
},
faild: (header, ex, retryCount, message) => Task.CompletedTask,
fallback: (header, message, ex) => Task.FromResult(ConsumerState.Ack));
return consumerTag;
}
public Task StopByQueueAsync(string queue)
=> _dynamicConsumer.StopConsumerAsync(queue);
}事件總線模式
事件總線模式可以自由編排事件執(zhí)行鏈路,框架會按照鏈路自動執(zhí)行并在執(zhí)行失敗后自動執(zhí)行回滾鏈路。
IEventMiddleware<T> 負(fù)責(zé)構(gòu)建執(zhí)行鏈,[EventOrder] 控制步驟順序,適合一個事件拆成多個業(yè)務(wù)步驟,能夠很好將業(yè)務(wù)解耦。
using Maomi.MQ.EventBus;
[RouterKey("biz.order.pipeline.v1")]
public sealed class OrderPipelineEvent
{
public Guid OrderId { get; set; } = Guid.NewGuid();
public decimal Amount { get; set; }
}
[Consumer("biz.order.pipeline.v1")]
public sealed class OrderPipelineMiddleware : IEventMiddleware<OrderPipelineEvent>
{
public Task ExecuteAsync(MessageHeader messageHeader, OrderPipelineEvent message, EventHandlerDelegate<OrderPipelineEvent> next)
{
return next(messageHeader, message, CancellationToken.None);
}
public Task FaildAsync(MessageHeader messageHeader, Exception ex, int retryCount, OrderPipelineEvent? message)
=> Task.CompletedTask;
public Task<ConsumerState> FallbackAsync(MessageHeader messageHeader, OrderPipelineEvent? message, Exception? ex)
=> Task.FromResult(ConsumerState.Ack);
}
[EventOrder(1)]
public sealed class ReserveInventoryHandler : IEventHandler<OrderPipelineEvent>
{
public Task ExecuteAsync(OrderPipelineEvent message, CancellationToken cancellationToken)
=> Task.CompletedTask;
public Task CancelAsync(OrderPipelineEvent message, CancellationToken cancellationToken)
=> Task.CompletedTask;
}
[EventOrder(2)]
public sealed class CreateBillHandler : IEventHandler<OrderPipelineEvent>
{
public Task ExecuteAsync(OrderPipelineEvent message, CancellationToken cancellationToken)
=> Task.CompletedTask;
public Task CancelAsync(OrderPipelineEvent message, CancellationToken cancellationToken)
=> Task.CompletedTask;
}廣播模式
例如為了減少延時和提供性能,在 Redis 做完緩存后,還需要做本地緩存。但是如果數(shù)據(jù)發(fā)生變更,怎么刷新本地緩存呢?
那就使用廣播模式,同一個服務(wù)的不同實(shí)例使用廣播模式時,每個實(shí)例都可以收到消息,而不是隨機(jī)分配給其中一個。
代碼非常簡單,設(shè)置 IsBroadcast = true 即可,這樣即使是同一個服務(wù)的不同實(shí)例,也會收到廣播通知,該實(shí)例下線,會自動取消訂閱,不會消耗服務(wù)器資源。
[Consumer("scenario.quickstart", IsBroadcast = true)]
public sealed class QuickStartConsumer : IConsumer<QuickStartMessage>
{
}
高性能序列化器
Maomi.MQ 默認(rèn)使用 JSON 做序列化數(shù)據(jù)傳輸,你也可以引入其它序列化器,提高壓縮消息的性能。
目前支持 System.Text.Json、Protobuf、Thrift、MessagePack 四種二進(jìn)制序列化協(xié)議,你可以自由選擇組合使用不同的序列化器到項(xiàng)目中,以便在不同微服務(wù)中傳遞消息,并且實(shí)現(xiàn)高性能傳遞消息。
例如使用 protobuf-net 框架 識別標(biāo)記了 [ProtoContract] 的模型類,那么此類型使用 Protobuf 協(xié)議壓縮消息,其它消息還是走 JSON。
using ProtoBuf;
builder.Services.AddMaomiMQ(options =>
{
options.MessageSerializers = serializers =>
{
// 添加 Protobuf 序列化器
serializers.Insert(0, new ProtobufMessageSerializer());
};
}, [typeof(Program).Assembly]);
[ProtoContract]
public sealed class PersonMessage
{
[ProtoMember(1)]
public Guid Id { get; set; } = Guid.NewGuid();
[ProtoMember(2)]
public string Name { get; set; } = string.Empty;
[ProtoMember(3)]
public int Age { get; set; }
}強(qiáng)一致性事務(wù)模式
借鑒 CAP 等框架的本地消息表模式,通過 MQ 和本地消息表,實(shí)現(xiàn)簡單的強(qiáng)一致性的分布式事務(wù),在業(yè)務(wù)不太復(fù)雜的企業(yè)項(xiàng)目中,可以簡化編寫事務(wù)的難度,不同的微服務(wù)以輕量、簡潔、不復(fù)雜的模式接入,降低了編碼和維護(hù)難度。
配置(以 MySQL 為例):
using Maomi.MQ.Transaction.Mysql;
using MySqlConnector;
builder.Services.AddMaomiMQTransactionMySql();
builder.Services.AddMaomiMQTransaction(options =>
{
options.ProviderName = TransactionProviderNames.MySql;
options.Connection = _ => new MySqlConnection(builder.Configuration.GetConnectionString("Default"));
options.AutoCreateTable = true;
});業(yè)務(wù)代碼發(fā)布消息:
public sealed class OrderAppService
{
private readonly IMessagePublisher _publisher;
private readonly string _connectionString;
public OrderAppService(IMessagePublisher publisher, IConfiguration configuration)
{
_publisher = publisher;
_connectionString = configuration.GetConnectionString("Default")!;
}
public async Task CreateOrderAsync(CancellationToken cancellationToken)
{
await using var connection = new MySqlConnection(_connectionString);
await connection.OpenAsync(cancellationToken);
await using var transaction = await connection.BeginTransactionAsync(cancellationToken);
// 1) 執(zhí)行業(yè)務(wù) SQL
// await SaveOrderAsync(connection, transaction, ...);
await transaction.CommitAsync(cancellationToken);
// 2) 發(fā)送消息(放在事務(wù)外面)
await _publisher.AutoPublishAsync(new OrderCreatedMessage
{
OrderNo = "SO-TX-001",
Amount = 520m
}, cancellationToken: cancellationToken);
}
}注意:本地事務(wù)模式不是全局分布式事務(wù)協(xié)調(diào)器,它解決的是“單庫與消息發(fā)送”的一致性。
如果事務(wù)提交了但是最后發(fā)送消息失敗,那么后臺會有一個 BackgroundService 定期掃描數(shù)據(jù)庫,將這些沒有發(fā)送成功的消息推送到 RabbitMQ。
結(jié)合 MediatR
如果你的業(yè)務(wù)已經(jīng)使用 MediatR,可以直接把 MQ 當(dāng)作 MediatR 命令輸入通道,也就是原有的代碼完全不需要改動,只需要在 Command 上加上 [MediatRConsumer] 即可。
using Maomi.MQ.MediatR;
using MediatR;
[MediatRConsumer("biz.mediatr.order", Qos = 1)]
public sealed class SyncOrderCommand : IRequest
{
public string OrderNo { get; set; } = string.Empty;
}
public sealed class SyncOrderCommandHandler : IRequestHandler<SyncOrderCommand>
{
public Task Handle(SyncOrderCommand request, CancellationToken cancellationToken)
=> Task.CompletedTask;
}
// 通過 MediatR 觸發(fā) MQ 發(fā)布
await mediator.Send(new MediatRMqCommand<SyncOrderCommand>
{
Message = new SyncOrderCommand { OrderNo = "SO-MED-001" }
});原理:MediatRTypeFilter 會把帶 MediatRConsumer 的命令類型映射成 MQ 消費(fèi)者,消費(fèi)后再轉(zhuǎn)發(fā)給 IMediator.Send(...)。
或者你可以繼續(xù)使用 IMessagePublisher 發(fā)布消息。
await _publisher.PublishAsync(string.Empty, request.Queue, message);
如果你的項(xiàng)目已經(jīng)引入了 MediatR,那么不需要為了使用 RabbitMQ 再搞出別的消費(fèi)模式代碼,使用 Maomi.MQ 可以直接簡化接入 RabbitMQ 的麻煩,繼續(xù)以 MediatR 的模式實(shí)現(xiàn)異步消費(fèi)。
結(jié)合 FastEndpoints
如果你使用 FastEndpoints,也可以通過類型過濾器把 IEvent/ICommand 接入 MQ。
using Maomi.MQ.Filters;
app.Services.RegisterGenericCommand(typeof(FastEndpointsMqCommand<>), typeof(FastEndpointsMqCommandHandler<>));
app.UseFastEndpoints();
[FastEndpointsConsumer("biz.fast.order.paid", Qos = 1)]
public sealed class OrderPaidEvent : FastEndpoints.IEvent
{
public string OrderNo { get; set; } = string.Empty;
}await _messagePublisher.AutoPublishAsync(new OrderPaidEvent
{
OrderNo = "SO-FE-001"
});
其它能力
下面這些能力這里先簡述,后續(xù)你可以按需深入:
- 動態(tài)路由配置:實(shí)現(xiàn)
IRoutingProvider,統(tǒng)一改寫Exchange/RoutingKey/Queue。 - 自定義重試機(jī)制:實(shí)現(xiàn)
IRetryPolicyFactory,按隊(duì)列/類型定制重試策略。 - 寬松自由可定制消費(fèi)者模式:實(shí)現(xiàn)
ITypeFilter,把第三方框架類型映射成 Maomi.MQ 消費(fèi)者。
到此這篇關(guān)于Maomi.MQ 功能強(qiáng)大的 .NET RabbitMQ 消息隊(duì)列通訊模型框架來了的文章就介紹到這了,更多相關(guān)Maomi.MQ 功能強(qiáng)大的 .NET RabbitMQ 消息隊(duì)列通訊模型框架來了內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
c# 操作符?? null coalescing operator
?? "null coalescing" operator 是c#新提供的一個操作符,這個操作符提供的功能是判斷左側(cè)的操作數(shù)是否是null,如果是則返回結(jié)果是右側(cè)的操作數(shù);非null則返回左側(cè)的操作數(shù)。2009-06-06
ASP.NET Core中使用MialKit實(shí)現(xiàn)郵件發(fā)送功能
這篇文章主要介紹了ASP.NET Core中使用MialKit實(shí)現(xiàn)郵件發(fā)送功能,本文給大家介紹的非常詳細(xì),具有一定的參考借鑒價值,需要的朋友可以參考下2019-10-10
Visual Studio 2017創(chuàng)建.net standard類庫編譯出錯原因及解決方法
這篇文章主要為大家詳細(xì)介紹了Visual Studio 2017創(chuàng)建.net standard類庫編譯出錯原因及解決方法,具有一定的參考價值,感興趣的小伙伴們可以參考一下2017-04-04
利用ASP.NET MVC+EasyUI+SqlServer搭建企業(yè)開發(fā)框架
本文主要介紹使用asp.net mvc4、sqlserver、jquery2.0和easyui1.4.5搭建企業(yè)級開發(fā)框架的過程,希望能夠幫到大家。2016-04-04
ASP.NET MVC 3仿Server.Transfer效果的實(shí)現(xiàn)方法
這篇文章主要介紹了ASP.NET MVC 3仿Server.Transfer效果的實(shí)現(xiàn)方法,需要的朋友可以參考下2015-10-10
CKEditor與dotnetcore實(shí)現(xiàn)圖片上傳功能
這篇文章主要為大家詳細(xì)介紹了CKEditor與dotnetcore實(shí)現(xiàn)圖片上傳功能,具有一定的參考價值,感興趣的小伙伴們可以參考一下2017-09-09
使用.NET?6開發(fā)TodoList應(yīng)用之引入數(shù)據(jù)存儲的思路詳解
在這篇文章中,我們僅討論如何實(shí)現(xiàn)數(shù)據(jù)存儲基礎(chǔ)設(shè)施的引入,具體的實(shí)體定義和操作后面專門來說。對.NET?6開發(fā)TodoList引入數(shù)據(jù)存儲相關(guān)知識感興趣的朋友一起看看吧2021-12-12
DataGridView - DataGridViewCheckBoxCell的使用介紹
Datagridview是.net中最復(fù)雜的控件,Datagridview中,用戶可以對行、列、單元格進(jìn)行編程,下面與大家分享下DataGridViewCheckBoxCell的使用,感興趣的朋友可以參考下哈2013-06-06

