如何使用MassTransit读取Azure死信队列(Dead Letter Queue)中的消息
实现方案
核心依赖与用到的API
首先安装必需NuGet包:
MassTransit.Azure.ServiceBus.CoreMicrosoft.Extensions.Hosting(控制台宿主用,也可以单独用依赖注入包)
核心用到的类与方法:
AddMassTransit:MassTransit服务注册入口,用来注册总线、消费者等组件UsingAzureServiceBus:指定Azure Service Bus作为传输层,配置服务总线连接信息ReceiveEndpoint:配置消息接收端点,此处需传入死信队列的完整路径IConsumer<>:消费者接口,实现该接口处理接收到的死信消息
完整代码示例
1. 主程序代码(.NET 6+ 顶级语句)
using MassTransit; using Microsoft.Extensions.Hosting; var host = Host.CreateDefaultBuilder(args) .ConfigureServices(services => { services.AddMassTransit(busConfig => { // 注册死信消息消费者 busConfig.AddConsumer<DeadLetterMessageConsumer>(); busConfig.UsingAzureServiceBus((context, sbConfig) => { // 替换为你的Azure Service Bus连接字符串 sbConfig.Host("Endpoint=sb://<你的服务总线名称>.servicebus.chinacloudapi.cn/;SharedAccessKeyName=<权限名>;SharedAccessKey=<密钥>"); // 配置死信队列接收端点 // 队列死信路径格式:<原队列名>/$DeadLetterQueue // 主题订阅死信路径格式:<主题名>/Subscriptions/<订阅名>/$DeadLetterQueue sbConfig.ReceiveEndpoint("<原业务队列名>/$DeadLetterQueue", endpointConfig => { // 关闭不必要的死信配置,避免死信队列的消息再次进入死信 endpointConfig.EnableDeadLetteringOnMessageExpiration = false; endpointConfig.AutoDeleteOnIdle = TimeSpan.MaxValue; // 注册消费者到当前端点 endpointConfig.ConfigureConsumer<DeadLetterMessageConsumer>(context); }); }); }); }) .Build(); // 启动总线开始监听死信消息 await host.StartAsync(); Console.WriteLine("死信队列监听已启动,按任意键退出程序"); Console.ReadKey(); // 优雅停止服务 await host.StopAsync();
2. 死信消息消费者实现
using MassTransit; using System.Text.Json; // 如果你明确知道死信消息的类型,可以把object替换为对应的强类型,提高性能 public class DeadLetterMessageConsumer : IConsumer<object> { public async Task Consume(ConsumeContext<object> context) { Console.WriteLine("\n========== 收到死信消息 =========="); // 输出消息基础元数据 Console.WriteLine($"消息ID:{context.MessageId}"); Console.WriteLine($"消息入队时间:{context.EnqueueTime:yyyy-MM-dd HH:mm:ss}"); // 输出消息正文 Console.WriteLine($"消息内容:{JsonSerializer.Serialize(context.Message)}"); // 读取Azure Service Bus原生死信属性 if(context.Headers.TryGetHeader("DeadLetterReason", out var deadReason)) { Console.WriteLine($"死信原因:{deadReason}"); } if (context.Headers.TryGetHeader("DeadLetterErrorDescription", out var errorDesc)) { Console.WriteLine($"死信错误详情:{errorDesc}"); } Console.WriteLine("==================================\n"); await Task.CompletedTask; } }
注意事项
- 如果仅需要一次性拉取现有死信消息而非持续监听,可以在接收端点配置中添加
endpointConfig.ConfigureConsumeTopology = false;,处理完预期数量的消息后主动调用host.StopAsync()退出即可 - 强类型消息场景下,消费者定义改为
IConsumer<你的消息类>即可直接拿到反序列化完成的消息对象,不需要手动处理JSON - 权限配置需确保连接字符串对应的访问密钥有死信队列的读取权限
内容的提问来源于stack exchange,提问作者ash
相关产品推荐
相关产品推荐

