一、Hangfire 简介

Hangfire 是一个开源的、可持久化的.NET 后台任务处理框架,由 Sergey Odinokov 于 2013 年发起并持续维护。它允许开发者在 ASP.NET Core 应用中创建各种类型的后台任务——无论是一次性的 fire-and-forget 任务、延迟执行任务,还是周期性的定时任务,都能通过统一的 API 优雅地实现。

与传统的定时任务解决方案相比,Hangfire 具有以下核心优势:

  • 持久化存储:任务数据存储在 Redis、SQL Server、PostgreSQL 等持久化数据库中,即使应用重启也不会丢失。
  • 可视化仪表盘:内置 Hangfire Dashboard,可实时监控任务状态、执行历史和重试情况。
  • 自动重试机制:任务执行失败时自动重试,支持指数退避策略。
  • 多种任务类型:支持 Fire-and-Forget、Delayed、Recurring、Continuations 等多种任务模式。
  • 分布式执行:支持多实例分布式部署,天然支持高可用和横向扩展。

下表展示了 Hangfire 与其他主流定时任务框架的对比:

特性 Hangfire Quartz.NET FluentScheduler TimerJobScheduler
持久化存储 多种数据库 多种数据库 内存 内存
Web Dashboard 内置 需第三方
分布式支持 原生支持 支持 单实例 单实例
自动重试 内置 内置 需实现 需实现
云原生集成 优秀 一般 有限 有限
学习曲线

二、项目架构设计

本系统采用经典的分层架构,将任务调度、执行、持久化三个核心关注点进行解耦。整体架构分为调度层、执行层、存储层和监控层四个部分。

下面的 Mermaid 图展示了各组件之间的关系:

graph TD A[客户端请求 API] --> B[调度器 NeTaskHangfireScheduler] B --> C{任务类型判断} C -->|EnableSchedule| D[RecurringJob.AddOrUpdate] C -->|Enqueue| E[BackgroundJob.Enqueue] C -->|Schedule| F[BackgroundJob.Schedule] D --> G[Hangfire 存储] E --> G F --> G G --> H[Hangfire Server 后台 Worker] H --> I[执行器工厂 ITaskExecutorFactory] I --> J[具体执行器 CollectTaskExecutor] J --> K[业务逻辑处理] G --> L[Hangfire Dashboard] M[状态追踪服务] --> G M --> N[ExecutionStatus 状态机]

核心组件说明:

  • NeTaskHangfireScheduler:统一的调度入口,封装了 Hangfire 的三种调度模式。
  • ITaskExecutorFactory:执行器工厂,根据任务类型动态创建对应的执行器实例。
  • NeTaskExecutorBase:执行器抽象基类,定义了执行器的标准生命周期。
  • JobWorkerService:后台 Worker 服务,继承自 BackgroundService,负责从存储层拉取并执行任务。

三、环境配置与持久化存储

3.1 NuGet 包安装

首先需要安装以下 NuGet 包:

dotnet add package Hangfire
dotnet add package Hangfire.AspNetCore
dotnet add package Hangfire.SqlServer

3.2 服务注册

在 Program.cs 中完成 Hangfire 服务注册和配置:

using Hangfire;
using Hangfire.SqlServer;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddHangfire(config =>
{
    config.UseSimpleAssemblyNameTypeSerializer();
    config.UseRecommendedSerializerSettings();

    config.UseSqlServer(
        builder.Configuration.GetConnectionString("HangfireConnection"),
        new SqlServerStorageOptions
        {
            CommandTimeout = TimeSpan.FromMinutes(5),
            QueuePollInterval = TimeSpan.FromSeconds(15),
            JobExpirationCheckInterval = TimeSpan.FromHours(1),
            CountersAggregateInterval = TimeSpan.FromMinutes(5),
            PrepareSchemaIfNecessary = true,
            TransactionTimeout = TimeSpan.FromMinutes(1),
            DashboardTitle = "企业级任务调度中心"
        });
});

builder.Services.AddHangfireServer(options =>
{
    options.WorkerCount = Environment.ProcessorCount * 2;
    options.Queues = new[] { "default", "critical", "low" };
    options.ServerName = $"hangfire-{Environment.MachineName}";
});

var app = builder.Build();

app.UseHangfireDashboard("/hangfire", new DashboardOptions
{
    Authorization = new[]
    {
        new HangfireAuthorizationFilter(isAuthenticated: true)
    },
    AppPath = "/",
    DashboardTitle = "任务调度监控面板"
});

app.Run();

3.3 数据库连接配置

在 appsettings.json 中添加连接字符串:

{
  "ConnectionStrings": {
    "HangfireConnection": "Server=.;Database=HangfireDemo;Trusted_Connection=True;TrustServerCertificate=True;"
  }
}

首次运行时,Hangfire 会自动在 SQL Server 中创建所需的数据库表结构,包括 HangFire.Job、HangFire.State、HangFire.RecurringJob 等核心表。

四、调度器实现

调度器 NeTaskHangfireScheduler 是整个系统的入口点,它封装了 Hangfire 的三种核心调度模式:EnableSchedule(周期性任务)、Enqueue(即时任务)和 Schedule(延迟任务)。

using Hangfire;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;

namespace GMCC.NeTask.Scheduler;

public class NeTaskHangfireScheduler : INeTaskScheduler
{
    private readonly ILogger<NeTaskHangfireScheduler> _logger;
    private readonly HangfireOptions _options;

    public NeTaskHangfireScheduler(
        ILogger<NeTaskHangfireScheduler> logger,
        IOptions<HangfireOptions> options)
    {
        _logger = logger;
        _options = options.Value;
    }

    public string EnableSchedule<TTask>(string cronExpression, string? jobId = null)
        where TTask : ITaskExecutor
    {
        var typeName = typeof(TTask).FullName;
        var finalJobId = jobId ?? $"{typeName}_{Guid.NewGuid():N}";

        RecurringJob.AddOrUpdate<NeTaskExecutorWrapper>(
            finalJobId,
            wrapper => wrapper.ExecuteAsync(typeName!),
            cronExpression,
            TimeZoneInfo.FindSystemTimeZoneById(_options.TimeZone));

        _logger.LogInformation(
            "已注册周期性任务:{JobId}, Cron: {Cron}, 类型: {TypeName}",
            finalJobId, cronExpression, typeName);

        return finalJobId;
    }

    public string Enqueue<TTask>(params KeyValuePair<string, object>[] parameters)
        where TTask : ITaskExecutor
    {
        var typeName = typeof(TTask).FullName;

        var jobId = BackgroundJob.Enqueue<NeTaskExecutorWrapper>(
            wrapper => wrapper.ExecuteAsync(typeName!, parameters));

        _logger.LogInformation(
            "已入队即时任务:{JobId}, 类型: {TypeName}",
            jobId, typeName);

        return jobId;
    }

    public string Schedule<TTask>(TimeSpan delay, params KeyValuePair<string, object>[] parameters)
        where TTask : ITaskExecutor
    {
        var typeName = typeof(TTask).FullName;

        var jobId = BackgroundJob.Schedule<NeTaskExecutorWrapper>(
            wrapper => wrapper.ExecuteAsync(typeName!, parameters),
            delay);

        _logger.LogInformation(
            "已调度延迟任务:{JobId}, 延迟: {Delay}, 类型: {TypeName}",
            jobId, delay, typeName);

        return jobId;
    }

    public bool RemoveSchedule(string jobId)
    {
        try
        {
            RecurringJob.RemoveIfExists(jobId);
            _logger.LogInformation("已移除周期性任务:{JobId}", jobId);
            return true;
        }
        catch (Exception ex)
        {
            _logger.LogWarning(ex, "移除周期性任务失败:{JobId}", jobId);
            return false;
        }
    }
}

五、执行器工厂模式

为了实现任务执行的可扩展性和可维护性,我们采用工厂模式来创建执行器实例。

5.1 执行器接口与基类

首先定义统一的执行器接口 INeTaskExecutor 和抽象基类 NeTaskExecutorBase:

namespace GMCC.NeTask.Executors;

public interface INeTaskExecutor
{
    Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default);
    string TaskName { get; }
}

public abstract class NeTaskExecutorBase : INeTaskExecutor
{
    protected readonly ILogger Logger;
    protected readonly ICurrentUserService CurrentUser;

    public abstract string TaskName { get; }

    protected NeTaskExecutorBase(ILogger logger, ICurrentUserService currentUser)
    {
        Logger = logger;
        CurrentUser = currentUser;
    }

    public async Task ExecuteAsync(TaskContext context, CancellationToken cancellationToken = default)
    {
        Logger.LogInformation("开始执行任务:{TaskName}, 任务ID: {TaskId}", TaskName, context.JobId);

        try
        {
            await SetRunningStatusAsync(context.JobId, cancellationToken);
            await ExecuteCoreAsync(context, cancellationToken);
            await SetCompletedStatusAsync(context.JobId, cancellationToken);

            Logger.LogInformation("任务执行成功:{TaskName}, 任务ID: {TaskId}", TaskName, context.JobId);
        }
        catch (TaskCanceledException)
        {
            Logger.LogWarning("任务被取消:{TaskName}, 任务ID: {TaskId}", TaskName, context.JobId);
            await SetCancelledStatusAsync(context.JobId, cancellationToken);
        }
        catch (Exception ex)
        {
            Logger.LogError(ex, "任务执行失败:{TaskName}, 任务ID: {TaskId}", TaskName, context.JobId);
            await SetFailedStatusAsync(context.JobId, ex, cancellationToken);
            throw;
        }
    }

    protected abstract Task ExecuteCoreAsync(TaskContext context, CancellationToken cancellationToken);

    protected virtual Task SetRunningStatusAsync(string jobId, CancellationToken ct)
        => Task.CompletedTask;

    protected virtual Task SetCompletedStatusAsync(string jobId, CancellationToken ct)
        => Task.CompletedTask;

    protected virtual Task SetFailedStatusAsync(string jobId, Exception ex, CancellationToken ct)
        => Task.CompletedTask;

    protected virtual Task SetCancelledStatusAsync(string jobId, CancellationToken ct)
        => Task.CompletedTask;
}

5.2 具体执行器示例

以数据采集任务为例,展示如何实现一个具体的执行器:

namespace GMCC.NeTask.Executors;

public class CollectTaskExecutor : NeTaskExecutorBase
{
    private readonly ICollectService _collectService;
    private readonly IDataRepository _repository;

    public override string TaskName => "数据采集执行器";

    public CollectTaskExecutor(
        ILogger<CollectTaskExecutor> logger,
        ICurrentUserService currentUser,
        ICollectService collectService,
        IDataRepository repository)
        : base(logger, currentUser)
    {
        _collectService = collectService;
        _repository = repository;
    }

    protected override async Task ExecuteCoreAsync(TaskContext context, CancellationToken cancellationToken)
    {
        var deviceId = context.GetParameter<string>("DeviceId");
        var dataType = context.GetParameter<string>("DataType");

        Logger.LogInformation("开始采集设备数据:Device={DeviceId}, Type={DataType}", deviceId, dataType);

        var data = await _collectService.FetchDataAsync(deviceId!, dataType!, cancellationToken);

        await _repository.SaveBatchAsync(data, cancellationToken);

        Logger.LogInformation("设备数据采集完成:Device={DeviceId}, 记录数={Count}",
            deviceId, data.Count);
    }

    protected override async Task SetRunningStatusAsync(string jobId, CancellationToken ct)
    {
        await _repository.UpdateTaskStatusAsync(jobId, ExecutionStatus.Running, ct);
    }

    protected override async Task SetCompletedStatusAsync(string jobId, CancellationToken ct)
    {
        await _repository.UpdateTaskStatusAsync(jobId, ExecutionStatus.Completed, ct);
    }

    protected override async Task SetFailedStatusAsync(string jobId, Exception ex, CancellationToken ct)
    {
        await _repository.UpdateTaskStatusAsync(jobId, ExecutionStatus.Failed, ct, ex.Message);
    }
}

5.3 执行器工厂注册

namespace GMCC.NeTask.Executors;

public interface ITaskExecutorFactory
{
    INeTaskExecutor Create(string typeName);
}

public class TaskExecutorFactory : ITaskExecutorFactory
{
    private readonly IServiceProvider _serviceProvider;
    private readonly IDictionary<string, Type> _executorMap;

    public TaskExecutorFactory(IServiceProvider serviceProvider)
    {
        _serviceProvider = serviceProvider;
        _executorMap = new Dictionary<string, Type>();
        ScanExecutors();
    }

    private void ScanExecutors()
    {
        var assembly = typeof(Program).Assembly;
        var executorTypes = assembly
            .GetTypes()
            .Where(t => typeof(INeTaskExecutor).IsAssignableFrom(t) && !t.IsAbstract);

        foreach (var type in executorTypes)
        {
            _executorMap[type.FullName!] = type;
        }
    }

    public INeTaskExecutor Create(string typeName)
    {
        if (!_executorMap.TryGetValue(typeName, out var executorType))
        {
            throw new InvalidOperationException(
                $"未找到类型为 {typeName} 的执行器,请确认该执行器已注册。");
        }

        return (INeTaskExecutor)_serviceProvider.GetRequiredService(executorType);
    }
}

六、执行引擎:幂等性与状态追踪

6.1 执行状态机

为了确保任务执行的幂等性和可追溯性,我们设计了一个完整的执行状态追踪体系。任务状态定义如下:

namespace GMCC.NeTask;

public enum ExecutionStatus
{
    Pending = 0,
    Queued = 1,
    Running = 2,
    Completed = 3,
    Failed = 4,
    Retrying = 5,
    Cancelled = 6,
    Timeout = 7
}

状态流转遵循以下规则:

  • Pending → Queued:任务创建后进入队列等待执行。
  • Queued → Running:Worker 抢到任务并开始执行。
  • Running → Completed:任务成功完成。
  • Running → Failed:任务执行出错,进入失败状态。
  • Failed → Retrying:Hangfire 自动重试,进入重试状态。
  • Retrying → Running:重试开始,重新进入执行状态。
  • Running → Cancelled:任务被外部取消。
  • Running → Timeout:任务执行超时。

6.2 状态追踪时序图

下面的时序图展示了任务从创建到完成的完整生命周期:

sequenceDiagram participant Client as 客户端 participant Scheduler as NeTaskHangfireScheduler participant Storage as Hangfire 存储 participant Worker as JobWorkerService participant Factory as ITaskExecutorFactory participant Executor as CollectTaskExecutor participant Repo as IDataRepository Client->>Scheduler: Enqueue<CollectTaskExecutor>(parameters) Scheduler->>Storage: BackgroundJob.Enqueue(typeName, params) Storage-->>Scheduler: 返回 JobId Scheduler-->>Client: 返回 JobId loop 每次轮询 Worker->>Storage: 获取待执行任务 alt 有任务 Storage-->>Worker: 返回任务信息 Worker->>Factory: Create(typeName) Factory-->>Worker: 返回执行器实例 Worker->>Executor: ExecuteAsync(context) Executor->>Repo: UpdateTaskStatusAsync(Running) Repo-->>Executor: OK Executor->>Executor: ExecuteCoreAsync Executor->>Repo: SaveBatchAsync(data) Repo-->>Executor: OK Executor->>Repo: UpdateTaskStatusAsync(Completed) Repo-->>Executor: OK Executor-->>Worker: 完成 end end

6.3 SetRunningStatusAsync 的实现细节

public async Task SetRunningStatusAsync(string jobId, CancellationToken ct)
{
    var existing = await _repository.GetByJobIdAsync(jobId, ct);

    if (existing != null && existing.Status == ExecutionStatus.Running)
    {
        Logger.LogWarning("任务 {JobId} 已在执行中,跳过重复执行", jobId);
        throw new IdempotentConflictException(jobId);
    }

    await _repository.UpdateStatusAsync(jobId, ExecutionStatus.Running, ct);
}

七、后台 Worker 服务

JobWorkerService 是一个自定义的后台 Worker 服务,它与 Hangfire Server 协同工作,提供额外的任务处理能力。

using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;

namespace GMCC.NeTask.Workers;

public class JobWorkerService : BackgroundService
{
    private readonly ILogger<JobWorkerService> _logger;
    private readonly ITaskExecutorFactory _executorFactory;
    private readonly IServiceProvider _serviceProvider;

    protected override int WorkerCount { get; } = Environment.ProcessorCount;

    public JobWorkerService(
        ILogger<JobWorkerService> logger,
        ITaskExecutorFactory executorFactory,
        IServiceProvider serviceProvider)
    {
        _logger = logger;
        _executorFactory = executorFactory;
        _serviceProvider = serviceProvider;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("JobWorkerService 已启动,WorkerCount: {Count}", WorkerCount);

        var workerTasks = Enumerable.Range(0, WorkerCount)
            .Select(i => RunWorkerLoopAsync(i, stoppingToken))
            .ToArray();

        await Task.WhenAll(workerTasks);
    }

    private async Task RunWorkerLoopAsync(int workerId, CancellationToken ct)
    {
        while (!ct.IsCancellationRequested)
        {
            try
            {
                await using var scope = _serviceProvider.CreateAsyncScope();
                var queueService = scope.ServiceProvider.GetRequiredService<IJobQueueService>();

                var job = await queueService.DequeueAsync(ct);

                if (job == null)
                {
                    await Task.Delay(TimeSpan.FromSeconds(5), ct);
                    continue;
                }

                _logger.LogDebug("Worker-{WorkerId} 开始处理任务:{JobId}", workerId, job.JobId);

                var executor = _executorFactory.Create(job.ExecutorType);
                var context = new TaskContext(job);

                await executor.ExecuteAsync(context, ct);
            }
            catch (OperationCanceledException) when (ct.IsCancellationRequested)
            {
                _logger.LogInformation("Worker-{WorkerId} 正在退出...", workerId);
                break;
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Worker-{WorkerId} 发生未处理异常", workerId);
                await Task.Delay(TimeSpan.FromSeconds(10), ct);
            }
        }
    }

    public override Task StopAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("JobWorkerService 正在停止...");
        return base.StopAsync(cancellationToken);
    }
}

八、Cron 表达式速查

Hangfire 使用标准的 Cron 表达式来定义周期性任务的触发时间。Cron 表达式由 5 个字段组成:分 时 日 月 周。

表达式 含义 常用场景
* * * * * 每分钟执行一次 高频轮询
*/5 * * * * 每 5 分钟执行一次 数据同步
0 * * * * 每小时整点执行 定时巡检
0 */2 * * * 每 2 小时执行一次 批量处理
0 0 * * * 每天零点执行 日终结算
0 2 * * * 每天凌晨 2 点执行 数据备份
0 0 1 * * 每月 1 号零点执行 月度报表
0 9 * * 1-5 工作日早 9 点执行 晨会提醒
0 0 0 1 1 * 每年 1 月 1 日零点 年度初始化
*/10 * * * * * 每 10 秒执行一次 实时监控

除了标准的 5 段式 Cron 表达式,Hangfire 还支持以下快捷别名:

RecurringJob.AddOrUpdate<MyTask>("job-id", 
    task => task.Run(), 
    Cron.Minutely);
RecurringJob.AddOrUpdate<MyTask>("job-id", 
    task => task.Run(), 
    Cron.Hourly);
RecurringJob.AddOrUpdate<MyTask>("job-id", 
    task => task.Run(), 
    Cron.Daily);
RecurringJob.AddOrUpdate<MyTask>("job-id", 
    task => task.Run(), 
    Cron.Weekly);
RecurringJob.AddOrUpdate<MyTask>("job-id", 
    task => task.Run(), 
    Cron.Monthly);

九、最佳实践与踩坑指南

在多年的 Hangfire 项目实践中,我们总结出以下 5 条关键最佳实践和常见陷阱:

9.1 使用 IoC 容器解析执行器

Hangfire 后台任务在执行时会从 DI 容器中解析依赖。如果执行器中使用了 Scoped 生命周期的服务,必须确保解析方式正确。推荐在 Program.cs 中通过 AddHangfire 配置 Job Activator:

builder.Services.AddHangfire(config =>
{
    config.UseSimpleAssemblyNameTypeSerializer();
    config.UseRecommendedSerializerSettings();
    config.UseSqlServer(connectionString);
});

Hangfire 默认使用 ActivatorUtilities 从 DI 容器解析任务类,确保所有执行器都已注册为 Scoped 或 Transient。

9.2 避免在任务中使用闭包捕获复杂对象

Hangfire 需要将任务方法的参数序列化到数据库中。闭包中捕获的复杂对象(如 DbContext、HttpClient)可能导致序列化异常。最佳做法是在任务方法内部通过 DI 容器获取所需服务:

// 错误做法
var data = _dbContext.Set<T>().ToList();
BackgroundJob.Enqueue(() => ProcessData(data));

// 正确做法
BackgroundJob.Enqueue<DataProcessor>(p => p.ProcessByIdAsync(recordId));

9.3 设置合理的重试策略

Hangfire 默认在任务失败后立即重试,最多 10 次。对于需要指数退避的场景,建议自定义重试策略:

GlobalJobFilters.Filters.Add(new AutomaticRetryAttribute
{
    Attempts = 5,
    DelayInSeconds = 30,
    OnAttemptsExceeded = AttemptsExceededAction.Delete
});

9.4 多实例部署时配置唯一 ServerName

在负载均衡环境下,多个 Hangfire Server 实例共用同一个数据库存储。必须为每个实例配置唯一的 ServerName,否则会导致任务竞争异常:

builder.Services.AddHangfireServer(options =>
{
    options.ServerName = $"app-{Environment.MachineName}-{Guid.NewGuid():N}";
    options.WorkerCount = 4;
    options.Queues = new[] { "critical", "default" };
});

9.5 定期清理过期任务数据

Hangfire 会自动清理过期的任务记录,但默认的清理间隔较长。建议通过 SqlServerStorageOptions 调整清理策略,并定期清理 Dashboard 中的历史数据:

config.UseSqlServer(connectionString, new SqlServerStorageOptions
{
    JobExpirationCheckInterval = TimeSpan.FromHours(1),
    CountersAggregateInterval = TimeSpan.FromMinutes(5),
    QueuePollInterval = TimeSpan.FromSeconds(10),
    TransactionTimeout = TimeSpan.FromMinutes(1)
});

十、总结

Hangfire 作为 .NET 生态中最成熟的后台任务调度框架之一,为企业级应用提供了可靠、可扩展的定时任务解决方案。通过本文介绍的架构设计和实践经验,我们可以总结出以下核心要点:

  1. 架构分层:将调度层、执行层、存储层解耦,使各模块可独立演进和测试。
  2. 工厂模式:通过执行器工厂实现任务类型的动态扩展,新增业务任务无需修改框架代码。
  3. 状态追踪:建立完整的执行状态机,配合幂等性检查,确保任务 Exactly-Once 执行。
  4. 持久化存储:利用 Hangfire 的多种存储后端(SQL Server、Redis、PostgreSQL),保证任务数据不丢失。
  5. 监控仪表盘:充分利用 Hangfire Dashboard 的可视化能力,实现任务执行的实时监控和运维诊断。

在实际项目中,还需要结合业务需求考虑任务优先级队列、分布式锁、任务依赖关系等高级特性。Hangfire 社区也在持续迭代,未来的版本将在 .NET 8+ 的 Native AOT、性能优化等方面进一步提升。

希望本文能为正在构建企业级定时任务系统的开发者们提供一个完整的参考方案,让 Hangfire 在你的项目中真正发挥价值。