如何让ASP.NET Core Web API订阅并接收Azure Service Bus消息
实现方案
你可以用ASP.NET Core自带的*托管后台服务(BackgroundService)*来承载Azure Service Bus的消息监听逻辑,和你控制台的实现逻辑几乎一致,只需要把逻辑放到后台服务里,随WebAPI启动运行、停止时优雅释放资源即可。
步骤1:安装依赖NuGet包
在WebAPI_2项目中安装Microsoft.Azure.ServiceBus和Microsoft.Extensions.Hosting.Abstractions两个依赖包。
步骤2:实现消息监听后台服务
新建ServiceBusListenerBackgroundService类,继承BackgroundService,把控制台的监听逻辑迁移进去,同时增加优雅停机的资源释放处理:
using Microsoft.Azure.ServiceBus; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Hosting; using System.Text; using System.Text.Json; using System.Threading; using System.Threading.Tasks; public class ServiceBusListenerBackgroundService : BackgroundService { private readonly IConfiguration _config; private IQueueClient _queueClient; private readonly IServiceProvider _serviceProvider; // 构造函数注入配置,如需调用其他业务服务也可以在这里注入IServiceProvider public ServiceBusListenerBackgroundService(IConfiguration config, IServiceProvider serviceProvider) { _config = config; _serviceProvider = serviceProvider; } protected override Task ExecuteAsync(CancellationToken stoppingToken) { _queueClient = new QueueClient(_config.GetConnectionString("AzureServiceBus"), "myqueue"); var messageHandlerOptions = new MessageHandlerOptions(ExceptionReceivedHandler) { MaxConcurrentCalls = 1, AutoComplete = false }; // 注册消息处理逻辑,和你控制台的实现完全一致 _queueClient.RegisterMessageHandler(ProcessMessagesAsync, messageHandlerOptions); // 监听应用停止信号,停止时关闭队列客户端释放资源 stoppingToken.Register(async () => { if (!_queueClient.IsClosedOrClosing) { await _queueClient.CloseAsync(); } }); return Task.CompletedTask; } private async Task ProcessMessagesAsync(Message message, CancellationToken token) { var jsonString = Encoding.UTF8.GetString(message.Body); // 替换为你自己的业务模型类 var obj = JsonSerializer.Deserialize<Model>(jsonString); // 如需调用Scoped生命周期的服务(比如DbContext),按下面的方式手动创建Scope获取 // using var scope = _serviceProvider.CreateScope(); // var dbContext = scope.ServiceProvider.GetRequiredService<YourDbContext>(); // 此处写入你的业务处理逻辑 Console.WriteLine($"Person Received: { obj.Field1} { obj.Field2}"); await _queueClient.CompleteAsync(message.SystemProperties.LockToken); } private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs exceptionReceivedEventArgs) { // 自行实现异常处理逻辑,比如打印日志、触发告警 var exception = exceptionReceivedEventArgs.Exception; Console.WriteLine($"消息处理异常:{exception.Message}"); return Task.CompletedTask; } }
步骤3:注册后台服务
在WebAPI_2的启动配置中注册刚才写的后台服务:
- 如果是.NET 6+的顶级语句Program.cs:
var builder = WebApplication.CreateBuilder(args); // 其他服务注册逻辑... builder.Services.AddHostedService<ServiceBusListenerBackgroundService>(); var app = builder.Build(); // 中间件配置逻辑... app.Run();
- 如果是.NET 5及更早版本的Startup.cs,在
ConfigureServices方法中添加:
public void ConfigureServices(IServiceCollection services) { // 其他服务注册逻辑... services.AddHostedService<ServiceBusListenerBackgroundService>(); }
注意事项
- 确保WebAPI_2的配置文件中已经正确配置了Azure Service Bus的连接字符串
- 可以根据业务需求调整
MessageHandlerOptions的参数,比如并发处理数、异常重试策略等 - 不要直接把Scoped生命周期的服务注入到BackgroundService的构造函数中,需要手动创建Scope获取,避免生命周期不匹配的问题
内容的提问来源于stack exchange,提问作者Solomon Levit
相关产品推荐
相关产品推荐

