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

部署于Azure App Service的.NET服务无法接收全部Service Bus队列消息

Azure Service Bus处理器处理消息后停止接收的解决方案

问题描述

我将一个.NET Windows Service项目部署到Azure App Service,目标是从Azure Service Bus队列接收并处理100,000条消息。但处理约800条消息后,服务停止接收新消息并进入空闲状态,如何确保服务能持续接收并处理完所有消息?

原代码

public void ReceiveMessages(string queueName, string namespaceConnectionString, string appInsightInstrumentationKey, string domainUrl, SecretClient client)
{
  var clientOptions = new ServiceBusClientOptions()
  {
      TransportType = ServiceBusTransportType.AmqpWebSockets,
  };
  var serviceBusClient = new ServiceBusClient(namespaceConnectionString, clientOptions);
  var adminClient = new ServiceBusAdministrationClient(namespaceConnectionString);
  var processorOptions = new ServiceBusProcessorOptions
  {
      AutoCompleteMessages = false,
      MaxAutoLockRenewalDuration = TimeSpan.FromHours(1)
  };
  var serviceBusProcessor = serviceBusClient.CreateProcessor(queueName, new ServiceBusProcessorOptions());
  _serviceBusProcessor = serviceBusProcessor;

  var handlerParams = new MessageHandlerParameter
  {
      AppInsightInstrumentationKey = appInsightInstrumentationKey,
      DomainUrl = domainUrl,
      Client = client
  };
  try
  {
      // Add handler to process messages;
      serviceBusProcessor.ProcessMessageAsync += (args) => MessageHandlerAsync(args, handlerParams);

      // Add handler to process any errors;
      serviceBusProcessor.ProcessErrorAsync += ErrorHandlerAsync;

      serviceBusProcessor.StartProcessingAsync();

      Console.WriteLine("Wait sometime for processing to start");
      Task.Delay(Timeout.InfiniteTimeSpan);
  }
  catch (Exception ex)
  {
  }
  finally
  {
      _ = serviceBusProcessor.DisposeAsync();
      _ = serviceBusClient.DisposeAsync();
  }
}

核心问题分析

  1. 处理器配置未生效:创建ServiceBusProcessor时传入了新的空ServiceBusProcessorOptions,而非提前配置好的processorOptions,导致自动锁续期、手动完成消息的设置完全失效。
  2. 异步方法未正确等待:方法为同步void类型,但调用了StartProcessingAsync()和Task.Delay()等异步方法却未await,导致代码直接进入finally块,提前释放了处理器和客户端资源,处理器刚启动就被终止。
  3. 空异常处理:catch块无任何日志记录,无法排查运行中出现的错误。
  4. 默认并发限制:默认MaxConcurrentCalls=1,处理吞吐量极低,加上App Service的空闲回收机制,容易触发进程休眠。

解决方案

1. 修正代码逻辑

将方法改为异步类型,确保配置生效并正确等待异步操作:

public async Task ReceiveMessages(string queueName, string namespaceConnectionString, string appInsightInstrumentationKey, string domainUrl, SecretClient client)
{
    var clientOptions = new ServiceBusClientOptions()
    {
        TransportType = ServiceBusTransportType.AmqpWebSockets,
    };
    var serviceBusClient = new ServiceBusClient(namespaceConnectionString, clientOptions);
    
    var processorOptions = new ServiceBusProcessorOptions
    {
        AutoCompleteMessages = false,
        MaxAutoLockRenewalDuration = TimeSpan.FromHours(1),
        // 根据业务处理能力调整并发数,提升吞吐量
        MaxConcurrentCalls = 10,
        // 设置预取数量,减少网络请求次数
        PrefetchCount = 50
    };
    // 使用配置好的processorOptions创建处理器
    var serviceBusProcessor = serviceBusClient.CreateProcessor(queueName, processorOptions);
    _serviceBusProcessor = serviceBusProcessor;

    var handlerParams = new MessageHandlerParameter
    {
        AppInsightInstrumentationKey = appInsightInstrumentationKey,
        DomainUrl = domainUrl,
        Client = client
    };
    try
    {
        serviceBusProcessor.ProcessMessageAsync += (args) => MessageHandlerAsync(args, handlerParams);
        serviceBusProcessor.ProcessErrorAsync += ErrorHandlerAsync;

        // 异步启动处理器并等待完成
        await serviceBusProcessor.StartProcessingAsync();

        Console.WriteLine("消息处理已启动,持续运行中...");
        // 异步等待无限延迟,避免线程阻塞
        await Task.Delay(Timeout.InfiniteTimeSpan);
    }
    catch (Exception ex)
    {
        // 添加日志记录,便于排查问题
        Console.WriteLine($"运行时异常: {ex.Message}", ex);
        // 可选:触发告警或尝试重启处理器
    }
    finally
    {
        // 确保处理器先停止再释放
        if (serviceBusProcessor != null)
        {
            await serviceBusProcessor.StopProcessingAsync();
            await serviceBusProcessor.DisposeAsync();
        }
        await serviceBusClient.DisposeAsync();
    }
}

2. 正确处理消息状态

在MessageHandlerAsync中必须明确标记消息的处理结果,避免处理器卡住:

private async Task MessageHandlerAsync(ProcessMessageEventArgs args, MessageHandlerParameter parameters)
{
    try
    {
        // 业务处理逻辑
        var messageContent = args.Message.Body.ToString();
        // ... 你的处理代码 ...

        // 处理完成后标记消息为已完成
        await args.CompleteMessageAsync(args.Message);
    }
    catch (Exception ex)
    {
        // 记录错误日志
        Console.WriteLine($"消息处理失败: {ex.Message}", ex);
        // 可恢复错误:放弃消息,让其重新排队
        await args.AbandonMessageAsync(args.Message);
        // 不可恢复错误:将消息移入死信队列
        // await args.DeadLetterMessageAsync(args.Message, "处理失败", ex.Message);
    }
}

3. 配置Azure App Service

  • 开启始终开启(Always On):在App Service的配置页面中启用该选项,避免服务因空闲被Azure回收。
  • 调整应用服务计划:根据消息处理的资源需求,选择合适的实例大小,确保CPU、内存资源充足。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 04:40:56