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

部署在VM上的Azure Service Bus消息处理服务突然挂起求助

Azure Service Bus生产环境消息处理挂起问题求助

我们通过API将消息排入Azure Service Bus队列,由4台不同VM上的「Task Engine Windows服务」处理消息。每台VM单独建立连接,租户ID、客户端ID和命名空间一致,未使用分区线程。测试环境运行正常,但生产环境(每日约700万消息)下,VM会突然停止处理消息并挂起。站点可靠性团队未发现VM存在CPU或内存问题,每次必须重启「Task Engine Windows服务」才能恢复消息处理。

Azure Service Bus指标

Azure Service Bus 指标

配置

  • AutoCompleteMessages = false
  • PrefetchCount = 50
  • ReceiveMode = PeekLock
  • MaxConcurrentCalls = 5
  • MaxAutoLockRenewalDuration = 15分钟(900秒)
  • MessageLockDuration = 5分钟

代码

初始化Azure连接(服务启动时执行一次)

private static void initializeAzureServiceBusClientConnection()
{
    if (_isAzureEnabled)
    {
        var _tenantId = GetFromConfig.AppSettings["tenant-id"];
        var _clientId = GetFromConfig.AppSettings["client-id"];
        var _clientSecret = GetFromConfig.AppSettings["client-secret"];
        var _servicebusNamespace = GetFromConfig.AppSettings["servicebus-namespace"];

        var _token = new ClientSecretCredential(_tenantId, _clientId, _clientSecret);
    
        var clientOptions = new ServiceBusClientOptions()
        {
            TransportType = ServiceBusTransportType.AmqpWebSockets,
        };

        ServiceBusClient = new ServiceBusClient(_servicebusNamespace, _token, clientOptions);
    }
}

消息处理逻辑

using Azure.Messaging.ServiceBus;

namespace TaskEngine
{
    internal sealed class AzureManager
    {
        private Thread _threadQueue;

        //Called once when Task Engine starts
        internal void Startup()
        {
            if (TaskEngineManager.IsAzureBusinessEventsEnabled)
            {
                Logger.Log(_isRunning == false, "The AzureManager startup method was called multiple times without a stop.");
                _isRunning = true;
                _messageQueueName = System.Configuration.ConfigurationManager.AppSettings["servicebus-queue"];
                _invalidMessageQueueName = System.Configuration.ConfigurationManager.AppSettings["servicebus-invalid-queue"];

                getAzureServiceBusAccess();

                // Start up background thread to manage queue
                _threadQueue = new Thread(queueLoop);                       //We are running in Thread because there is another process that runs in parallel
                _threadQueue.Name = "Azure Business Event Queue Thread";
                _threadQueue.Start();
            }
            else
            {
                Tracer.RaiseInfo(TraceSourceNames.AzureManager, "Azure Service Bus is disabled.");
            }
        }


        private void queueLoop()
        {
            bool isConnectionReset;

            do
            {
                isConnectionReset = false; //Reset in each loop
                isConnectionReset = AsyncGatekeeper.Run(() => processUsingAzureServiceBusMQAsync());
            } while (_isWorking && isConnectionReset);
        }



        private async System.Threading.Tasks.Task<bool> processUsingAzureServiceBusMQAsync()
        {
            if (_isWorking)
            {
                try
                {
                    // add handler to process messages
                    _serviceBusReceiverProcessor.ProcessMessageAsync += MessageHandlerAsync;

                    // add handler to process any errors
                    _serviceBusReceiverProcessor.ProcessErrorAsync += ErrorHandlerAsync;

                    await _serviceBusReceiverProcessor.StartProcessingAsync();

                }
                catch (Exception ex)
                {
                    Tracer.RaiseError(TraceSourceNames.AzureManager, "Azure Service Bus unhandled exception occurred. Attempting to reset connection.", ex);
                    try
                    {
                        bool isTaskCompleted = await closeAndDisposeConnectionAsync();
                        if (isTaskCompleted)
                        {
                            tryResetConnections(ex);
                            cooldownReceiveRequests();
                        }
                    }
                    catch (Exception)
                    {
                        Tracer.RaiseError(TraceSourceNames.AzureManager, "Azure Service Bus unhandled exception occurred. Reset connection failed.", ex);
                    }
                    return true;    //reset retry loop
                }
            }
            return false;
        }


        async System.Threading.Tasks.Task MessageHandlerAsync(ProcessMessageEventArgs args)
        {
            var message = args.Message;
            try
            {
                // start processing 
                Tracer.RaiseInfo(TraceSourceNames.AzureManager, $"Azure Service Bus Event Queue thread received message {message.MessageId}.");
                var _messageBody = message.Body;

                bool isSuccess = processMessage(message.MessageId, _messageBody.ToStream());

                if (!isSuccess)
                {
                    await WriteToInvalidQueueAsync(_messageBody);
                }

                args.MessageLockLostAsync += OnMessageLockLostAsync;

                // complete the message. message is deleted from the queue. 
                try
                {
                    await args.CompleteMessageAsync(message);
                }
                catch (Exception e)
                {
                    Tracer.RaiseWarning(TraceSourceNames.AzureManager, "Azure service bus message lock already released when completing message. Continuing processing messages.", e);
                }

                Tracer.RaiseInfo(TraceSourceNames.AzureManager, $"Azure Service Bus Event Queue thread completed processing message {message.MessageId}.");
            }
            catch (Exception ex)
            {
                Tracer.RaiseError(TraceSourceNames.AzureManager, "Unexpected error occured moving task from Azure Service Bus to database; attempting to re-queue message.", ex);
                try
                {
                    await args.AbandonMessageAsync(args.Message);
                }
                catch (Exception e)
                {
                    Tracer.RaiseWarning(TraceSourceNames.AzureManager, "Azure service bus message lock already released when abandoning message.", e);
                }
            }
        }

        System.Threading.Tasks.Task ErrorHandlerAsync(ProcessErrorEventArgs args)
        {
            var errMessage = $"Azure Service Bus error handler caught exception. Message processing will retry. Error Source: {args.ErrorSource}; EntityPath: {args.EntityPath}; FullyQualifiedNamespace: {args.FullyQualifiedNamespace}";
            Tracer.RaiseError(TraceSourceNames.AzureManager, errMessage, args.Exception);
            return System.Threading.Tasks.Task.CompletedTask;
        }

        System.Threading.Tasks.Task OnMessageLockLostAsync(MessageLockLostEventArgs args)
        {
            string errMessage = $"Azure Service Bus Message lock Lost handler caught exception. Message: {args.Message}. Message Locked Until: {args.Message.LockedUntil.ToLocalTime()}";

            var messageInfo = convertMessage(args.Message.Body.ToStream());
            if (messageInfo != null)
            {
                errMessage = $"Azure Service Bus Message Lock Lost handler caught exception. Message processing will retry. Message Info: {messageInfo}. Message Locked Until: {args.Message.LockedUntil.ToLocalTime()}";
            }

            Tracer.RaiseError(TraceSourceNames.AzureManager, errMessage, args.Exception);
            return System.Threading.Tasks.Task.CompletedTask;
        }

        //Called when Task Engine is started for the first time and any time when connection is lost and needs to be reset
        private static void getAzureServiceBusAccess()
        {
            Logger.Log(!(string.IsNullOrEmpty(_messageQueueName)
                            || string.IsNullOrEmpty(_invalidMessageQueueName))
                            , "Incomplete Azure Service Bus configuration settings. Unable to proceed.");

            Tracer.RaiseInfo(TraceSourceNames.AzureManager, $"Attempting to connect to Azure Service Bus.");
            _serviceBusSender = TaskEngineManager.ServiceBusClient.CreateSender(_invalidMessageQueueName, new ServiceBusSenderOptions());

            var processorOptions = new ServiceBusProcessorOptions()
            {
                MaxConcurrentCalls = 5,
                PrefetchCount = 50,
                MaxAutoLockRenewalDuration = TimeSpan.FromSeconds(900),     //Default 15 min. This lock renewal duration should be greater than LockDuration which is set to 5 min.
                AutoCompleteMessages = false,
            };
            _serviceBusReceiverProcessor = TaskEngineManager.ServiceBusClient.CreateProcessor(_messageQueueName, processorOptions);
        }

        private void tryResetConnections(Exception exception)
        {
            Tracer.RaiseInfo(TraceSourceNames.AzureManager, "Resetting Azure Service Bus connection and retrying.");

            try
            {
                if (DateTime.Now.Subtract(LastQueueReset).TotalSeconds > 1800)
                {
                    LastQueueReset = DateTime.Now;
                    getAzureServiceBusAccess();
                    Tracer.RaiseInfo(TraceSourceNames.AzureManager, "Azure Service Bus Connection succesfully reset.");
                }
                else
                {
                    NotificationRegulator.Notify(NotificationReason.AzureServiceBusFailure, exception, Resources.AzureServiceBus_ConnectError);
                }
            }
            catch (Exception ex)
            {
                Tracer.RaiseError(TraceSourceNames.AzureManager, "Resetting Azure Service Bus connection failed.", ex);
                throw;
            }
        }

        private static async System.Threading.Tasks.Task<bool> closeAndDisposeConnectionAsync()
        {
            try
            {
                await _serviceBusReceiverProcessor.StopProcessingAsync();
                await _serviceBusReceiverProcessor.DisposeAsync();
            }
            catch (Exception ex)
            {
                //Do not throw and eat exception - Receiver may have been already disposed
                Tracer.RaiseWarning(TraceSourceNames.AzureManager, "Error disposing Azure Service Bus Receiver Connection.", ex);
            }

            try
            {
                await _serviceBusSender.DisposeAsync();
            }
            catch (Exception ex)
            {
                //Do not throw and eat exception - Receiver may have been already disposed
                Tracer.RaiseWarning(TraceSourceNames.AzureManager, "Error disposing Azure Service Bus Sender Connection.", ex);
            }

            return true;
        }
    }
}

我们添加了大量日志排查,但仅偶尔出现LockLost异常,连接未自动释放或关闭,除非手动重启。调整配置参数后问题仍未解决,每日需重启服务至少3次,恳请提供解决方案!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:02:33