.NET 5及以上版本中Azure Service Bus Queue Listener替代方案是什么?
Azure Service Bus 队列监听器最新实现方案(.NET 5+)
背景说明
原Microsoft.Azure.ServiceBus NuGet包已正式标记为弃用,.NET 5及以上版本官方推荐使用新一代SDK Azure.Messaging.ServiceBus 实现队列监听等操作,原IQueueClient相关API已被新的处理器类替代。
实现步骤
1. 安装依赖包
通过NuGet安装最新版SDK:
dotnet add package Azure.Messaging.ServiceBus
2. 基础控制台/普通应用实现
核心使用ServiceBusProcessor类实现队列监听,代码示例:
using Azure.Messaging.ServiceBus; // 配置参数 string serviceBusConnectionString = "你的Service Bus连接字符串"; string queueName = "目标队列名称"; // 1. 创建ServiceBusClient,建议全局单例复用 await using var client = new ServiceBusClient(serviceBusConnectionString); // 2. 创建队列处理器 ServiceBusProcessor processor = client.CreateProcessor(queueName, new ServiceBusProcessorOptions()); // 3. 注册消息处理委托 processor.ProcessMessageAsync += async args => { try { // 解析消息内容 string messageBody = args.Message.Body.ToString(); Console.WriteLine($"接收到消息:{messageBody}"); // 消息处理完成后手动确认消费 await args.CompleteMessageAsync(args.Message); } catch (Exception ex) { // 处理失败,将消息放回队列/转入死信队列 Console.WriteLine($"消息处理失败:{ex.Message}"); await args.AbandonMessageAsync(args.Message); } }; // 4. 注册错误处理委托 processor.ProcessErrorAsync += args => { Console.WriteLine($"监听出错:{args.Exception.Message}"); return Task.CompletedTask; }; // 5. 启动监听 await processor.StartProcessingAsync(); Console.WriteLine("已启动队列监听,按任意键停止..."); Console.ReadKey(); // 停止监听 await processor.StopProcessingAsync();
3. ASP.NET Core 依赖注入实现
如果是ASP.NET Core项目,推荐通过依赖注入托管监听器:
- 在
Program.cs中注入相关服务:
builder.Services.AddSingleton<ServiceBusClient>(sp => new ServiceBusClient(builder.Configuration["ServiceBus:ConnectionString"])); // 注册后台托管服务运行监听器 builder.Services.AddHostedService<QueueListenerService>();
- 实现后台监听服务
QueueListenerService:
public class QueueListenerService : BackgroundService { private readonly ServiceBusClient _serviceBusClient; private ServiceBusProcessor _processor; public QueueListenerService(ServiceBusClient serviceBusClient, IConfiguration configuration) { _serviceBusClient = serviceBusClient; var queueName = configuration["ServiceBus:QueueName"]; _processor = _serviceBusClient.CreateProcessor(queueName); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // 注册消息和错误处理逻辑 _processor.ProcessMessageAsync += ProcessMessageAsync; _processor.ProcessErrorAsync += ProcessErrorAsync; // 启动监听 await _processor.StartProcessingAsync(stoppingToken); await Task.Delay(Timeout.Infinite, stoppingToken); } private async Task ProcessMessageAsync(ProcessMessageEventArgs args) { // 业务处理逻辑 var content = args.Message.Body.ToString(); // 处理完成确认 await args.CompleteMessageAsync(args.Message); } private Task ProcessErrorAsync(ProcessErrorEventArgs args) { // 错误处理逻辑 Console.WriteLine(args.Exception.ToString()); return Task.CompletedTask; } public override async Task StopAsync(CancellationToken cancellationToken) { await _processor.StopProcessingAsync(cancellationToken); await base.StopAsync(cancellationToken); } }
核心注意事项
ServiceBusClient为线程安全类,全应用生命周期内仅需初始化一个实例即可,频繁创建销毁会导致连接资源浪费。- 消息处理状态可根据业务场景选择:
CompleteMessageAsync(确认消费,消息从队列移除)、AbandonMessageAsync(放弃消费,消息放回队列重试)、DeadLetterMessageAsync(转入死信队列,不再重试)。 - 必须实现
ProcessErrorAsync委托,避免未捕获的异常导致监听进程终止。
内容的提问来源于stack exchange,提问作者Emmanouil Dagdilelis
相关产品推荐
相关产品推荐

