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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 22:45:09