接RabbitMQ完全指南:隊(duì)列、死信與重試錯(cuò)誤處理實(shí)戰(zhàn))
SlimMessageBus對(duì)接RabbitMQ完全指南隊(duì)列、死信與重試錯(cuò)誤處理實(shí)戰(zhàn)【免費(fèi)下載鏈接】SlimMessageBusLightweight message bus interface for .NET (pub/sub and request-response) with transport plugins for popular message brokers.項(xiàng)目地址: https://gitcode.com/gh_mirrors/sl/SlimMessageBusSlimMessageBus是一個(gè)輕量級(jí)的 .NET 消息總線message bus支持發(fā)布/訂閱pub/sub和請(qǐng)求/響應(yīng)request-response兩種通信模式。搭配 RabbitMQ 傳輸插件后你可以用統(tǒng)一的 API 對(duì)接 RabbitMQ 的 Exchange、Queue 和 Binding并輕松實(shí)現(xiàn)死信隊(duì)列DLX與重試錯(cuò)誤處理是 .NET 微服務(wù)架構(gòu)中對(duì)接 RabbitMQ 的簡(jiǎn)單可靠方案。 5分鐘搭建 RabbitMQ 環(huán)境官方倉(cāng)庫(kù)自帶 Docker Compose 編排文件其中包含帶管理界面的 RabbitMQ 服務(wù)端口 5672管理界面 15672默認(rèn)賬號(hào) guest/guestrabbitmq: container_name: slim.rabbitmq image: rabbitmq:4.2.3-management-alpine ports: - 5672:5672 - 15672:15672參考文件src/Infrastructure/docker-compose.yml在 NuGet 中安裝SlimMessageBus.Host.RabbitMQ和SlimMessageBus.Host.Serialization.Json兩個(gè)包即可開(kāi)始。對(duì)接核心一行配置連接 RabbitMQRabbitMQ 模型圍繞三個(gè)概念Exchange交換機(jī)生產(chǎn)者發(fā)消息的入口、Queue隊(duì)列消費(fèi)者的收件箱、Binding綁定規(guī)則決定消息如何路由到隊(duì)列。SlimMessageBus 通過(guò)WithProviderRabbitMQ配置連接并自動(dòng)完成拓?fù)浣粨Q、隊(duì)列、綁定的創(chuàng)建無(wú)需手動(dòng)寫 AMQP 代碼services.AddSlimMessageBus(mbb { mbb.WithProviderRabbitMQ(cfg cfg.ConnectionString configuration[RabbitMQ:ConnectionString]); mbb.ProduceOrderEvent(x x.Exchange(orders, exchangeType: ExchangeType.Fanout)); mbb.ConsumeOrderEvent(x x .Queue(orders-queue) .ExchangeBinding(orders) .DeadLetterExchange(orders-dlq, exchangeType: ExchangeType.Direct) .WithConsumerOrderCreatedConsumer()); mbb.AddJsonSerializer(); }); 關(guān)鍵點(diǎn)SlimMessageBus 會(huì)自動(dòng)聲明交換機(jī)、隊(duì)列和綁定——生產(chǎn)者聲明的 Exchange、消費(fèi)者聲明的 Queue 及其 Binding都會(huì)被自動(dòng)在 RabbitMQ 中創(chuàng)建。SlimMessageBus RabbitMQ 傳輸中多種消息類型共用同一交換機(jī)的示意圖隊(duì)列與路由鍵支持通配符匹配Topic 類型交換機(jī)下SlimMessageBus 完整支持 RabbitMQ 的通配符路由鍵模式匹配規(guī)則示例*恰好一個(gè)片段regions.na.cities.*匹配regions.na.cities.toronto#零個(gè)或多個(gè)片段audit.events.#匹配audit.events.orders.placed#單獨(dú)使用匹配所有路由鍵全量訂閱mbb.ConsumeRegionEvent(x x .Queue(na-cities-queue) .ExchangeBinding(regions, routingKey: regions.na.cities.*)); 性能優(yōu)化內(nèi)部的路由鍵匹配器優(yōu)先使用精確匹配僅在無(wú)精確命中時(shí)才做通配符模式匹配并對(duì)模式做了緩存。詳見(jiàn) src/SlimMessageBus.Host.RabbitMQ/Services/RoutingKeyMatcherService.cs。若消息的路由鍵未匹配到任何消費(fèi)者可通過(guò)MessageUnrecognizedRoutingKeyHandler自定義行為默認(rèn) Ack 丟棄也可改為 Nack 送入死信隊(duì)列、或 Requeue 稍后重試適合滾動(dòng)發(fā)布場(chǎng)景。死信隊(duì)列DLX失敗消息不再丟失消費(fèi)者處理失敗時(shí)SlimMessageBus 默認(rèn)會(huì)向 RabbitMQ 發(fā)送Nack消息將被路由到隊(duì)列上配置的死信交換機(jī)或丟棄。推薦做法是為消費(fèi)隊(duì)列配置 DLXmbb.ConsumePingMessage(x x .Queue(subscriber, autoDelete: false) .ExchangeBinding(ping) // 隊(duì)列將引用死信交換機(jī)指定 exchangeType 后 DLX 也會(huì)被自動(dòng)創(chuàng)建 .DeadLetterExchange(subscriber-dlq, exchangeType: ExchangeType.Direct) .WithConsumerPingConsumer());也可以在總線級(jí)別為所有死信交換機(jī)設(shè)置默認(rèn)值如統(tǒng)一使用 Direct 類型mbb.WithProviderRabbitMQ(cfg cfg.UseDeadLetterExchangeDefaults(durable: false, autoDelete: false, exchangeType: ExchangeType.Direct, routingKey: string.Empty));這樣失敗消息會(huì)進(jìn)入subscriber-dlq對(duì)應(yīng)的隊(duì)列方便后續(xù)排查與補(bǔ)償處理。重試與錯(cuò)誤處理三種確認(rèn)模式 自定義錯(cuò)誤處理器確認(rèn)模式?jīng)Q定至少一次還是至多一次SlimMessageBus 提供三種 Ack 確認(rèn)模式默認(rèn)是最安全的ConfirmAfterMessageProcessingWhenNoManualConfirmMade模式行為投遞保證ConfirmAfterMessageProcessingWhenNoManualConfirmMade默認(rèn)處理成功 Ack出錯(cuò) Nack支持手動(dòng)干預(yù)至少一次at-least-onceAckAutomaticByRabbit由 RabbitMQ 協(xié)議層自動(dòng) Ack至多一次at-most-onceAckMessageBeforeProcessing處理前立即 Ack至多一次吞吐量略低定義見(jiàn) src/SlimMessageBus.Host.RabbitMQ/Config/RabbitMqMessageAcknowledgementMode.cs自定義錯(cuò)誤處理器實(shí)現(xiàn)重試 N 次后進(jìn)入死信實(shí)現(xiàn)IRabbitMqConsumerErrorHandlerT接口即可完全掌控失敗消息的命運(yùn)——例如瞬時(shí)故障則重入隊(duì)列Requeue持久故障則 Nack 進(jìn) DLQpublic class CustomRabbitMqConsumerErrorHandlerT : IRabbitMqConsumerErrorHandlerT { public TaskProcessResult OnHandleError(T message, IConsumerContext ctx, Exception exception, int attempts) { if (exception is TransientException) return Task.FromResultProcessResult(RabbitMqProcessResult.Requeue); // 重試 return Task.FromResultProcessResult(ProcessResult.Failure); // 送死信 } } // 注冊(cè)到 DI對(duì)任意消息類型生效 services.AddTransient(typeof(IRabbitMqConsumerErrorHandler), typeof(CustomRabbitMqConsumerErrorHandler)); 源碼src/SlimMessageBus.Host.RabbitMQ/Consumers/IRabbitMqConsumerErrorHandler.cs。消費(fèi)者內(nèi)部還能通過(guò)ConsumerContext手動(dòng)調(diào)用Ack()/Nack()精確控制單條消息。生產(chǎn)環(huán)境可靠性發(fā)布確認(rèn)與斷線自動(dòng)恢復(fù)Publisher Confirms可選調(diào)用cfg.UsePublisherConfirms()開(kāi)啟發(fā)布確認(rèn)Broker 拒絕NACK消息時(shí)會(huì)拋異常適合金融、訂單等關(guān)鍵場(chǎng)景默認(rèn)關(guān)閉以保吞吐也支持按生產(chǎn)者單獨(dú)開(kāi)啟/退出。斷線自動(dòng)恢復(fù)RabbitMqChannelManager會(huì)無(wú)限次后臺(tái)重連默認(rèn)每 5 秒連接恢復(fù)后自動(dòng)重建通道、重建拓?fù)洳⒆屗邢M(fèi)者重新注冊(cè)重啟 RabbitMQ 容器也無(wú)需人工干預(yù)。消費(fèi)者并發(fā)通過(guò)cfg.ConnectionFactory.ConsumerDispatchConcurrency調(diào)高單實(shí)例并發(fā)需要保證順序時(shí)保持為 1??偨Y(jié)需求SlimMessageBus 方案自動(dòng)創(chuàng)建 Exchange/Queue/Binding生產(chǎn)者/消費(fèi)者聲明時(shí)自動(dòng)拓?fù)涔┙o通配符路由鍵*/#模式 內(nèi)置路由鍵匹配緩存失敗消息兜底.DeadLetterExchange(...)一行配置 DLX重試邏輯自定義IRabbitMqConsumerErrorHandlerT Requeue投遞保證默認(rèn) at-least-once可切換 at-most-once斷線恢復(fù)自動(dòng)重連 消費(fèi)者自動(dòng)重新注冊(cè)延伸閱讀docs/provider_rabbitmq.md 包含請(qǐng)求/響應(yīng)、發(fā)布確認(rèn)超時(shí)、多消費(fèi)者同隊(duì)列等完整內(nèi)容更多示例可參考 src/Tests/SlimMessageBus.Host.RabbitMQ.Test/ 下的集成測(cè)試?!久赓M(fèi)下載鏈接】SlimMessageBusLightweight message bus interface for .NET (pub/sub and request-response) with transport plugins for popular message brokers.項(xiàng)目地址: https://gitcode.com/gh_mirrors/sl/SlimMessageBus創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考