最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Maomi.MQ?功能強(qiáng)大的?.NET?RabbitMQ?消息隊(duì)列通訊模型框架詳解

 更新時間:2026年03月09日 08:49:03   作者:癡者工良  
Maomi.MQ.RabbitMQ 是一個基于 RabbitMQ 的消息隊(duì)列封裝框架,提供了很多開箱即用的功能,通過簡單靈活的方式簡化消息傳輸流程,在本篇教程中,將介紹 Maomi.MQ.RabbitMQ 的使用方法,以便讀者能夠快速了解該框架的使用方式和特點(diǎn),感興趣的朋友跟隨小編一起看看吧

快速開始

在本篇教程中,將介紹 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)文章

最新評論

土默特左旗| 铜川市| 乐业县| 理塘县| 砀山县| 陆河县| 遂川县| 丁青县| 桦南县| 遵义县| 闵行区| 丁青县| 泊头市| 石泉县| 遂昌县| 中牟县| 安仁县| 平潭县| 洪湖市| 怀化市| 绍兴县| 烟台市| 宁阳县| 北安市| 密云县| 蕲春县| 长泰县| 石景山区| 乐亭县| 蓝山县| 焉耆| 灵石县| 阿拉善右旗| 临潭县| 高清| 阿荣旗| 清新县| 临猗县| 揭阳市| 航空| 巴林右旗|