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

Hangfire工作线程被占用致关键任务阻塞的解决方案咨询

问题描述

我们通过Azure Communication Services向用户发送邮件和短信,每条消息对应Hangfire的EmailNotificationJob/SmsNotificationJob任务。目前遇到两个核心问题:

  1. 频繁触发Azure邮件每小时100封的限制,导致Hangfire任务陷入长达1小时的停滞。
  2. 工作线程优先处理notifications队列的低优先级任务,这些任务停滞时会占满所有线程,导致critical队列的关键任务无法得到处理。

当前Hangfire配置代码:

services.AddHangfire((serviceProvider, globalConfiguration) => globalConfiguration
    .SetDataCompatibilityLevel(CompatibilityLevel.Version_170)
    .UseRecommendedSerializerSettings()
    .UseRedisStorage(redisConnectionString, redisStorageOptions)
    .UseBatches()
    .UseThrottling(ThrottlingAction.RetryJob, 1.Seconds())
    .UseTagsWithRedis(new() { TagsListStyle = TagsListStyle.Dropdown }, redisStorageOptions)
    .UseApplicationInsightsTelemetry(serviceProvider));

services.AddHangfireServer((provider, options) =>
{
    options.Activator = new HangfireJobActivator(provider);
    options.WorkerCount = Environment.ProcessorCount * 5;
    options.HeartbeatInterval = 10.Seconds();
    options.SchedulePollingInterval = 1.Seconds();
    options.Queues = new[] { "critical", "default", "notifications" };
    options.ServerName = Environment.MachineName;
});

邮件任务及Azure客户端代码:

backgroundJobClient.Enqueue<EmailNotificationJob>(x => x.Perform(new(notificationMessage.Email, notificationMessage.Subject, notificationMessage.Body), null));

...
[Queue("notifications")]
[AutomaticRetry(Attempts = 3)]
public class EmailNotificationJob
{
    readonly IServiceProvider _serviceProvider;

    public EmailNotificationJob(IServiceProvider serviceProvider) => _serviceProvider = serviceProvider;

    public async Task Perform(EmailNotificationJobOptions jobOptions, PerformContext? performContext)
    {
        using var scope = _serviceProvider.CreateScope();
        var emailClient =  scope.ServiceProvider.GetRequiredService<AzureEmailClient>();
        
        if (string.IsNullOrWhiteSpace(jobOptions.Email))
        {
            throw new InvalidArgumentException("Email address is empty.");
        }

        var content = jobOptions.Body.ToHtml();
        await emailClient.SendAsync(jobOptions.Subject, content, jobOptions.Email);
    }
}

public sealed class AzureEmailClient
{
    readonly ILogger<AzureEmailClient> _logger;
    readonly string _sender;
    readonly EmailClient _emailClient;

    public AzureEmailClient(ILogger<AzureEmailClient> logger, IConfiguration configuration)
    {
        _logger = logger;
        _sender = configuration[ConfigurationKeys.AzureCommunicationServicesFromEmail];
        _emailClient = configuration.CreateEmailClient();
    }

    public async Task SendAsync(string subject, string htmlContent, string recipient, CancellationToken cancellationToken = default)
    {
        try
        {
            _logger.LogInformation("{ClientName}: Initiating sending of an email to {Recipient} with subject: {Subject}",
                nameof(AzureEmailClient), recipient, subject);
            
            var emailSendOperation = await _emailClient.SendAsync(WaitUntil.Completed, _sender, recipient, subject, htmlContent, null, cancellationToken);
            
            var operationId = emailSendOperation.Id;
            _logger.LogInformation("{ClientName}: Email Sent. Status = {Status}. Operation ID = {OperationId}",
                nameof(AzureEmailClient), emailSendOperation.Value.Status, operationId);
        }
        catch (RequestFailedException ex)
        {
            _logger.LogError("{ClientName}: Email send operation failed with error code: {ErrorCode}, message: {Message}",
                nameof(AzureEmailClient), ex.ErrorCode, ex.Message);
            throw;
        }
    }
}
解决方案

不要单纯增加线程数,这只会让更多线程陷入停滞,反而加剧资源浪费。推荐从以下几个维度解决:

1. 拆分Hangfire服务器,隔离队列优先级

为不同优先级的队列配置独立的Hangfire服务器实例,确保critical队列的任务始终有专属线程处理:

// 专门处理critical队列的服务器
services.AddHangfireServer((provider, options) =>
{
    options.Activator = new HangfireJobActivator(provider);
    options.WorkerCount = 2; // 根据关键任务量调整
    options.Queues = new[] { "critical" };
    options.ServerName = $"{Environment.MachineName}-critical";
});

// 处理default和notifications队列的服务器
services.AddHangfireServer((provider, options) =>
{
    options.Activator = new HangfireJobActivator(provider);
    options.WorkerCount = Environment.ProcessorCount * 3; // 原线程数分配给低优队列
    options.Queues = new[] { "default", "notifications" };
    options.ServerName = $"{Environment.MachineName}-notifications";
});

这样即使notifications队列的任务全部停滞,critical队列的服务器线程也不会被占用,关键任务能正常执行。

2. 针对Azure速率限制做任务流控

(1)实现任务延迟重试,避免立即占满线程

在邮件任务中捕获速率限制异常(Azure的RequestFailedException会返回429 Too Many Requests状态码),手动将任务延迟到限制解除后再执行,而非依赖Hangfire的自动重试:

public async Task Perform(EmailNotificationJobOptions jobOptions, PerformContext? performContext)
{
    using var scope = _serviceProvider.CreateScope();
    var emailClient = scope.ServiceProvider.GetRequiredService<AzureEmailClient>();
    
    if (string.IsNullOrWhiteSpace(jobOptions.Email))
    {
        throw new InvalidArgumentException("Email address is empty.");
    }

    var content = jobOptions.Body.ToHtml();
    try
    {
        await emailClient.SendAsync(jobOptions.Subject, content, jobOptions.Email);
    }
    catch (RequestFailedException ex) when (ex.Status == (int)HttpStatusCode.TooManyRequests)
    {
        // 提取重试间隔(Azure响应头通常包含Retry-After)
        var retryAfter = ex.Headers?.RetryAfter?.Delta ?? TimeSpan.FromHours(1);
        if (performContext != null)
        {
            // 延迟重新排队,避免当前线程被阻塞
            BackgroundJob.Schedule<EmailNotificationJob>(x => x.Perform(jobOptions, null), retryAfter);
        }
        // 标记当前任务为成功,避免自动重试
        return;
    }
}

同时修改AutomaticRetry属性,禁用速率限制异常的自动重试:

[Queue("notifications")]
[AutomaticRetry(Attempts = 3, LogEvents = true, OnAttemptsExceeded = AttemptsExceededAction.Delete,
    ExceptionsFilter = typeof(NonRateLimitExceptionFilter))]
public class EmailNotificationJob
{
    // ...
}

// 自定义异常过滤器,只重试非速率限制的异常
public class NonRateLimitExceptionFilter : IRetryExceptionFilter
{
    public bool Filter(Exception exception)
    {
        return exception is not RequestFailedException rateLimitEx || rateLimitEx.Status != (int)HttpStatusCode.TooManyRequests;
    }
}

(2)预限流任务入队

在添加邮件任务到队列前,先做流量控制,避免短时间内入队超过Azure限制的任务数量。比如用Redis计数器实现每小时100封的限流:

public async Task EnqueueEmailNotification(NotificationMessage notificationMessage, IDistributedCache cache)
{
    var counterKey = $"email-send-counter:{DateTime.UtcNow:yyyyMMddHH}";
    var currentCount = await cache.GetAsync(counterKey);
    var count = currentCount != null ? int.Parse(Encoding.UTF8.GetString(currentCount)) : 0;
    
    if (count < 100)
    {
        backgroundJobClient.Enqueue<EmailNotificationJob>(x => x.Perform(new(notificationMessage.Email, notificationMessage.Subject, notificationMessage.Body), null));
        await cache.SetStringAsync(counterKey, (count + 1).ToString(), new DistributedCacheEntryOptions
        {
            AbsoluteExpirationRelativeToNow = TimeSpan.FromHours(1)
        });
    }
    else
    {
        // 计算到下一小时的延迟时间
        var now = DateTime.UtcNow;
        var nextHour = now.AddHours(1).Date.AddHours(now.Hour + 1);
        var delay = nextHour - now;
        backgroundJobClient.Schedule<EmailNotificationJob>(x => x.Perform(new(notificationMessage.Email, notificationMessage.Subject, notificationMessage.Body), null), delay);
    }
}

3. 优化Hangfire重试策略

当前配置的UseThrottling(ThrottlingAction.RetryJob, 1.Seconds())会让失败任务每秒重试,遇到速率限制时会导致大量无效重试,浪费资源。建议调整为指数退避重试:

// 替换原Throttling配置,改用自动重试的指数退避
services.AddHangfire((serviceProvider, globalConfiguration) => globalConfiguration
    // ... 其他配置
    .UseAutomaticRetry(new AutomaticRetrySettings
    {
        Attempts = 5,
        DelayInSeconds = 60, // 初始延迟1分钟
        MaxDelayInSeconds = 3600, // 最大延迟1小时
        DelaysInSeconds = new[] { 60, 300, 900, 1800, 3600 } // 自定义退避间隔
    })
    // ... 其他配置
);
总结

优先采用拆分Hangfire服务器隔离队列的方案解决关键任务被阻塞的问题,同时结合速率限制感知的任务延迟重试和入队前限流来处理Azure的邮件发送限制,避免线程资源被无效占用。单纯增加线程数无法从根本上解决问题,反而会加剧资源浪费。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 18:57:34