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

如何在API#2中实现Azure Service Bus队列监听器?现有代码未触发

问题分析与解决方案

核心问题

你的监听器失效的直接原因是**Consume方法中注册消息处理器后立刻调用了queueClient.CloseAsync()**,客户端被直接关闭,根本没机会监听队列消息。除此之外还有几个关键问题导致功能异常:

  • 静态queueClient搭配Transient生命周期的MessageConsumer,会导致资源管理混乱,容易引发泄漏
  • 仅在DI容器注册了IMessageConsumer,但未在应用启动时触发监听逻辑,监听器从未真正运行
  • SaveChanges用了同步版本,不符合异步编程规范
  • ExceptionReceivedHandler无日志输出,出问题无法排查
  • 直接调用Update方法可能导致不必要的字段覆盖

修复后的代码示例

1. 调整MessageConsumer实现

public class MessageConsumer : IMessageConsumer, IHostedService
{
    private readonly string _connectionString = "stringTakenFromAzure";
    private readonly string _queueName = "cartqueue";
    private IQueueClient _queueClient;
    private readonly CartingDbContext _context;
    private readonly ILogger<MessageConsumer> _logger;

    public MessageConsumer(CartingDbContext context, ILogger<MessageConsumer> logger)
    {
        _context = context;
        _logger = logger;
    }

    // 应用启动时自动触发监听
    public Task StartAsync(CancellationToken cancellationToken)
    {
        _queueClient = new QueueClient(_connectionString, _queueName);

        var options = new MessageHandlerOptions(ExceptionReceivedHandler)
        {
            MaxConcurrentCalls = 1,
            AutoComplete = false
        };

        _queueClient.RegisterMessageHandler(ProcessMessageAsync, options);
        _logger.LogInformation("Azure Service Bus监听器已启动,监听队列: {QueueName}", _queueName);
        return Task.CompletedTask;
    }

    // 应用关闭时安全停止监听
    public async Task StopAsync(CancellationToken cancellationToken)
    {
        if (_queueClient != null)
        {
            await _queueClient.CloseAsync();
            _logger.LogInformation("Azure Service Bus监听器已停止");
        }
    }

    private async Task ProcessMessageAsync(Microsoft.Azure.ServiceBus.Message message, CancellationToken cancellationToken)
    {
        try
        {
            var jsonBody = Encoding.UTF8.GetString(message.Body);
            var categoryItem = JsonSerializer.Deserialize<CategoryItem>(jsonBody);
            _logger.LogInformation("收到队列消息,ItemId: {ItemId}", categoryItem.Id);

            var categoryItemInDb = await _context.CategoryItems.FirstOrDefaultAsync(x => x.Id == categoryItem.Id, cancellationToken);
            if (categoryItemInDb == null)
            {
                _context.CategoryItems.Add(categoryItem);
                _logger.LogInformation("新增Item到数据库: {ItemId}", categoryItem.Id);
            }
            else
            {
                // 手动更新必要字段,避免全量覆盖
                categoryItemInDb.Name = categoryItem.Name;
                categoryItemInDb.Price = categoryItem.Price;
                // 补充其他需要更新的字段
                _logger.LogInformation("更新数据库中的Item: {ItemId}", categoryItem.Id);
            }

            await _context.SaveChangesAsync(cancellationToken);
            await _queueClient.CompleteAsync(message.SystemProperties.LockToken);
            _logger.LogInformation("消息处理完成,已确认: {MessageId}", message.MessageId);
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "处理消息失败,MessageId: {MessageId}", message.MessageId);
            // 处理失败后让消息重新入队,可根据业务需求改为移入死信队列
            await _queueClient.AbandonAsync(message.SystemProperties.LockToken);
        }
    }

    private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs args)
    {
        _logger.LogError(args.Exception, "Service Bus异常: {ExceptionSource}", args.ExceptionReceivedContext.Source);
        return Task.CompletedTask;
    }
}

2. 修改Program.cs的DI注册

// 注册为HostedService,应用启动时自动启动监听
builder.Services.AddHostedService<MessageConsumer>();
// 若需通过IMessageConsumer接口访问,可额外注册
builder.Services.AddSingleton<IMessageConsumer>(provider => 
    provider.GetRequiredService<IHostedService>() as MessageConsumer);

关键修复点说明

  1. 实现IHostedService接口:ASP.NET Core会自动在应用启动/关闭时触发对应方法,无需手动调用监听逻辑
  2. 移除静态queueClient:改用实例字段配合Singleton生命周期(HostedService默认单例),避免资源泄漏
  3. 添加日志记录:覆盖启动、消息处理、异常全流程,方便排查问题
  4. 异步化数据库操作:使用异步版EF Core方法,符合异步编程模型
  5. 优化实体更新逻辑:手动更新必要字段,防止意外覆盖数据库中已有数据
  6. 完善错误处理:消息处理失败时记录日志并让消息重新入队,保证业务可靠性

内容的提问来源于stack exchange,提问作者Greg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 06:09:24