You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用MassTransit读取Azure死信队列(Dead Letter Queue)中的消息

实现方案

核心依赖与用到的API

首先安装必需NuGet包:

  • MassTransit.Azure.ServiceBus.Core
  • Microsoft.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.03 12:24:03