Hangfire工作线程被占用致关键任务阻塞的解决方案咨询
我们通过Azure Communication Services向用户发送邮件和短信,每条消息对应Hangfire的EmailNotificationJob/SmsNotificationJob任务。目前遇到两个核心问题:
- 频繁触发Azure邮件每小时100封的限制,导致Hangfire任务陷入长达1小时的停滞。
- 工作线程优先处理
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

