在企业级后台系统中,定时任务几乎无处不在:定时采集日志、定时解析文件、定时同步数据。简单的 BackgroundServiceTimer 可以跑通 demo,但一到生产环境就会遇到一堆问题:任务重复执行怎么办?执行到一半崩溃了如何恢复?手动触发和定时触发如何统一?新加一种任务类型要不要改调度代码?

这篇文章结合一个真实的三层 .NET 项目实践,介绍如何基于 Hangfire 搭建一套带幂等控制、状态机、执行器工厂的任务调度系统。代码分层清晰,扩展新任务类型时只需要新增一个执行器类,不用碰核心调度逻辑。

一、整体架构:三层协作

整个调度系统由三层协作完成:

  • 调度层NeTaskHangfireScheduler,封装 Hangfire 的 IRecurringJobManagerIBackgroundJobClient,负责任务的注册、移除、立即触发、延迟触发。
  • 执行层NeTaskExecutor,负责把 Hangfire 的作业转换为一次真正的业务执行,包括前置校验、状态流转、异常捕获。
  • 持久层NeTaskNeTaskExecution,前者保存任务定义(Cron 表达式、启用状态),后者保存每一次执行的明细(状态、起止时间、幂等键)。
flowchart TB
    subgraph 调度层
        A[IRecurringJobManager] --> B[NeTaskHangfireScheduler]
        C[IBackgroundJobClient] --> B
    end
    B -->|Enqueue / Schedule / AddOrUpdate| D[Hangfire Server]
    D -->|调用静态方法| E[NeTaskExecutor]
    E --> F{任务存在?}
    F -->|否| G[记录日志并返回]
    F -->|是| H[NeTaskExecutionService]
    H --> I[创建执行记录]
    I --> J[工厂获取执行器]
    J --> K[执行业务逻辑]
    K --> L[更新状态为 Success/Failed]

这套架构的核心思想是:Hangfire 只负责“什么时候触发”,业务系统自己负责“怎么执行、执行结果是什么”。这样即使将来替换调度引擎,执行层也可以原样保留。

二、任务模型:定义与执行分离

系统中存在两类实体。

2.1 任务定义 NeTask

[DbTable("NeTask")]
public sealed class NeTask : DbObject<NeTask>
{
    public int Id { get; set; }
    public string Name { get; set; }
    public NeTaskType TaskType { get; set; }
    public bool Enabled { get; set; }
    public string CronExpression { get; set; }
    public string PayloadJson { get; set; }
    public string JobId { get; set; }
    public int? LastExecutionId { get; set; }
}

NeTask 只描述“要做什么、多久做一次”,不保存执行结果。JobId 是 Hangfire 中的 recurring job 标识,用于后续禁用或更新任务。

2.2 执行记录 NeTaskExecution

[DbTable("NeTaskExecution")]
public sealed class NeTaskExecution : DbObject<NeTaskExecution>
{
    public int Id { get; set; }
    public string IdempotencyKey { get; set; }
    public ExecutionStatus Status { get; set; }
    public int? TaskId { get; set; }
    public NeTaskType TaskType { get; set; }
    public NeTaskTriggerSource TriggerSource { get; set; }
    public string PayloadJson { get; set; }
    public DateTime? StartTime { get; set; }
    public DateTime? EndTime { get; set; }
    public string Message { get; set; }
    public string ExecutionJobId { get; set; }
}

IdempotencyKey 是幂等控制的关键。ExecutionJobId 则记录 Hangfire 本次执行的 BackgroundJob Id,方便在 Hangfire Dashboard 中反向追查。

2.3 枚举:让状态与来源自解释

public enum NeTaskType
{
    [Description("采集日志")] Collect = 1,
    [Description("解析入库")] Parse = 2,
    [Description("UPF地址解析入库")] ParseUpfAddress = 3,
    [Description("中兴PCF统计文件解析入库")] ParseZxPcfStatFile = 4
}

public enum NeTaskScheduleType
{
    [Description("立即执行")] Immediate = 1,
    [Description("定时执行")] OneTime = 2,
    [Description("周期执行")] Recurring = 3
}

public enum ExecutionStatus
{
    Pending, Running, Success, Failed
}

把来源(手动、定时、系统)和类型都建模成枚举,而不是用 magic number,可读性和可维护性会好很多。

三、调度层:对 Hangfire 做薄封装

NeTaskHangfireScheduler 把 Hangfire 的 API 收敛到业务语义上,避免业务代码直接调用 Hangfire:

public class NeTaskHangfireScheduler
{
    private readonly IBackgroundJobClient _backgroundJobClient;
    private readonly IRecurringJobManager _recurringJobManager;

    public NeTaskHangfireScheduler(
        IBackgroundJobClient backgroundJobClient,
        IRecurringJobManager recurringJobManager)
    {
        _backgroundJobClient = backgroundJobClient;
        _recurringJobManager = recurringJobManager;
    }

    public void EnableSchedule(int taskId, string jobId, string cronExpression)
    {
        ArgumentException.ThrowIfNullOrWhiteSpace(jobId);
        ArgumentException.ThrowIfNullOrWhiteSpace(cronExpression);

        _recurringJobManager.AddOrUpdate(
            jobId,
            () => NeTaskExecutor.ExecuteRecurringAsync(taskId, NeTaskTriggerSource.Schedule, null),
            cronExpression,
            new RecurringJobOptions { TimeZone = TimeZoneInfo.Local });
    }

    public void DisableSchedule(string jobId)
    {
        _recurringJobManager.RemoveIfExists(jobId);
    }

    public void TriggerSchedule(int taskId)
    {
        _backgroundJobClient.Enqueue(
            () => NeTaskExecutor.ExecuteRecurringAsync(taskId, NeTaskTriggerSource.Manual, null));
    }

    public void Enqueue(int executionId)
    {
        _backgroundJobClient.Enqueue(
            () => NeTaskExecutor.ExecuteAsync(executionId, null));
    }

    public void Schedule(int executionId, DateTime executeAt)
    {
        _backgroundJobClient.Schedule(
            () => NeTaskExecutor.ExecuteAsync(executionId, null), executeAt);
    }
}

注意两个设计点:

  1. 时区使用 TimeZoneInfo.Local,避免 Hangfire 默认 UTC 导致 Cron 表达式在国内服务器上错位一小时。
  2. 触发目标都是静态方法。Hangfire 序列化的是方法调用表达式,静态方法不依赖实例,可以避免 DI 作用域、序列化陷阱等问题。

四、执行引擎:工厂 + 泛型基类

业务任务类型会不断新增,如果直接在 NeTaskExecutor 里写 switch,代码会迅速膨胀。项目采用执行器工厂 + 泛型基类来解耦。

4.1 执行器接口与泛型基类

public interface INeTaskExecutor
{
    NeTaskType TaskType { get; }
    Type PayloadType { get; }
    Task<JsonDataResult> ExecuteAsync(NeTaskExecution execution, object payload,
        PerformContext context, CancellationToken cancellationToken = default);
    string CreateIdempotencyKey(int? taskId, object payload);
    bool TryGetPayload(string payloadJson, out object payload);
}

public interface INeTaskExecutor<TPayload> : INeTaskExecutor
{
    Task<JsonDataResult> ExecuteAsync(NeTaskExecution execution, TPayload payload,
        PerformContext context, CancellationToken cancellationToken = default);
}

public abstract class NeTaskExecutorBase<TPayload> : INeTaskExecutor<TPayload>
{
    public abstract NeTaskType TaskType { get; }
    public Type PayloadType => typeof(TPayload);

    public abstract Task<JsonDataResult> ExecuteAsync(NeTaskExecution execution,
        TPayload payload, PerformContext context, CancellationToken cancellationToken = default);

    public abstract string CreateIdempotencyKey(int? taskId, TPayload payload);

    async Task<JsonDataResult> INeTaskExecutor.ExecuteAsync(NeTaskExecution execution,
        object payload, PerformContext context, CancellationToken cancellationToken)
    {
        return await ExecuteAsync(execution, (TPayload)payload, context, cancellationToken);
    }

    bool INeTaskExecutor.TryGetPayload(string payloadJson, out object payload)
    {
        try
        {
            payload = JsonHelper.Deserialize<TPayload>(payloadJson);
            return payload != null;
        }
        catch
        {
            payload = default;
            return false;
        }
    }

    string INeTaskExecutor.CreateIdempotencyKey(int? taskId, object payload)
    {
        return CreateIdempotencyKey(taskId, (TPayload)payload);
    }
}

泛型基类做了三件事:

  • 固定 PayloadType
  • 统一 JSON 反序列化与错误兜底。
  • 显式实现接口,把 object 强转成 TPayload,让子类只关心强类型逻辑。

4.2 工厂集中管理执行器

public static class NeTaskExecutorFactory
{
    private static readonly Dictionary<NeTaskType, INeTaskExecutor> _executors;

    static NeTaskExecutorFactory()
    {
        _executors = new INeTaskExecutor[]
        {
            new CollectTaskExecutor(),
            new ParseTaskExecutor(),
            new ParseUpfAddressTaskExecutor(),
            new ParseZxPcfStatFileTaskExecutor()
        }.ToDictionary(e => e.TaskType);
    }

    public static INeTaskExecutor GetExecutor(NeTaskType taskType)
    {
        if (_executors.TryGetValue(taskType, out var executor))
            return executor;
        throw new NotSupportedException($"不支持任务类型:{taskType}");
    }
}

新增任务类型只需:定义 Payload → 实现 NeTaskExecutorBase<TPayload> → 在工厂里注册。核心执行器与调度层完全不用改。

五、幂等控制:同一时刻只执行一次

周期性任务最怕的就是“上一次还没跑完,下一次又开始了”。项目用数据库唯一约束 + 幂等键来解决。

创建执行记录时,执行器负责生成 IdempotencyKey

public override string CreateIdempotencyKey(int? taskId, ParseTaskPayload payload)
{
    if (taskId.HasValue)
        return $"TaskId={taskId}";
    if (payload.NeGroupId.HasValue)
        return $"TaskType={(int)TaskType}:NeGroupId={payload.NeGroupId}";
    if (payload.DeviceId.HasValue)
        return $"TaskType={(int)TaskType}:DeviceId={payload.DeviceId}";
    throw new InvalidOperationException("网元组ID与设备ID不能同时为空值");
}

周期任务以 TaskId 作为幂等键;一次性任务则以任务类型 + 业务维度组合作为幂等键。数据库对 IdempotencyKey 加唯一索引后,重复插入会抛异常。

在调度入口用 try-catch 吞掉重复异常,本次触发直接跳过:

public static Task<NeTaskExecution> TryCreateAsync(int taskId, NeTaskType taskType,
    NeTaskTriggerSource triggerSource, string payloadJson, CancellationToken cancellationToken)
{
    try
    {
        return CreateAsync(taskId, taskType, triggerSource, payloadJson, null, cancellationToken);
    }
    catch
    {
        return null;
    }
}

注意:这里返回 null 表示“正在执行,跳过本次触发”,而不是失败。UI 上可以把这种情况显示为“被忽略”而不是“报错”。

六、状态机:Pending → Running → Success/Failed

执行一次任务的完整流程如下:

stateDiagram-v2
    [*] --> Pending: 创建执行记录
    Pending --> Running: SetRunningStatusAsync
    Running --> Success: 业务执行成功
    Running --> Failed: 业务异常或返回失败
    Success --> [*]
    Failed --> [*]

核心执行代码:

private static async Task ExecuteAsync(NeTaskExecution execution,
    PerformContext context, CancellationToken cancellationToken)
{
    execution.ExecutionJobId = context?.BackgroundJob.Id;
    execution.StartTime = DateTime.Now;

    try
    {
        if (execution.TaskId.HasValue)
        {
            await NeTaskService.UpdateLastExecutionIdAsync(
                execution.TaskId.Value, execution.Id, cancellationToken);
        }

        var executor = NeTaskExecutorFactory.GetExecutor(execution.TaskType);

        if (!executor.TryGetPayload(execution.PayloadJson, out var payload))
            throw new InvalidOperationException($"任务类型 {execution.TaskType} 的负载JSON不合法");

        var result = await executor.ExecuteAsync(execution, payload, context, cancellationToken);

        execution.Status = result.success ? ExecutionStatus.Success : ExecutionStatus.Failed;
        execution.Message = result.message;
    }
    catch (Exception ex)
    {
        execution.Status = ExecutionStatus.Failed;
        execution.Message = ex.Message;
    }

    execution.EndTime = DateTime.Now;
    await NeTaskExecutionService.UpdateAsync(execution, cancellationToken);
}

周期任务与一次性任务的入口略有不同:

  • 周期任务:由 Hangfire 直接调用 ExecuteRecurringAsync,内部先查任务定义、校验启用状态,再创建执行记录,状态初始化为 Running
  • 一次性/延迟任务:创建记录时状态为 Pending,Hangfire 真正执行前调用 SetRunningStatusAsync,只有 Pending 状态才能被更新为 Running,防止任务被重复消费。
public static async Task<bool> SetRunningStatusAsync(int executionId, CancellationToken cancellationToken)
{
    // UPDATE NeTaskExecution SET Status = Running
    // WHERE Id = @id AND Status = Pending
    ...
    return result > 0;
}

这种“条件更新 + 影响行数判断”是分布式任务状态机的经典写法,天然具备并发安全。

七、手动触发与延迟触发

周期任务之外的两种常见场景也一并覆盖。

立即执行

public static async Task<JsonDataResult> CreateAsync(NeTaskType taskType, object payload,
    DateTime? executeAt = null, CancellationToken cancellationToken = default)
{
    string payloadJson = JsonHelper.Serialize(payload);
    var execution = await CreateAsync(null, taskType, NeTaskTriggerSource.Manual,
        payloadJson, executeAt, cancellationToken);
    return new JsonDataResult { success = true, data = execution };
}

CreateAsync 在内部会根据 executeAt 决定是 Enqueue 还是 Schedule

延迟执行校验

if (executeAt.HasValue && executeAt.Value < DateTime.Now.AddMinutes(delayMinutes))
{
    throw new ArgumentException($"定时执行时间 {executeAt} 比当前时间要晚 {delayMinutes} 分钟以上", nameof(executeAt));
}

这里预留 5 分钟缓冲,避免用户把“立即执行”误配成“过去的时间”。

八、扩展新任务类型:四步即可

以新增“清理过期日志”任务为例:

  1. 定义 Payload:
public class CleanupTaskPayload : ITaskPayload
{
    public int RetentionDays { get; set; }
}
  1. 在枚举中新增类型:
public enum NeTaskType
{
    ...
    [Description("清理过期日志")] Cleanup = 5
}
  1. 实现执行器:
public class CleanupTaskExecutor : NeTaskExecutorBase<CleanupTaskPayload>
{
    public override NeTaskType TaskType => NeTaskType.Cleanup;

    public override string CreateIdempotencyKey(int? taskId, CleanupTaskPayload payload)
    {
        if (!taskId.HasValue) throw new ArgumentNullException(nameof(taskId));
        return $"TaskId={taskId}";
    }

    public override async Task<JsonDataResult> ExecuteAsync(NeTaskExecution execution,
        CleanupTaskPayload payload, PerformContext context, CancellationToken cancellationToken)
    {
        // 业务逻辑
        return new JsonDataResult { success = true, message = "清理完成" };
    }
}
  1. 在工厂注册:
new CleanupTaskExecutor()

无需改动 Hangfire 调度器、执行器入口、状态机代码。

九、总结

这套设计的核心收益:

  1. 调度与执行解耦:Hangfire 只负责触发,业务执行逻辑独立演化。
  2. 幂等防重:通过 IdempotencyKey + 数据库唯一约束,避免同一任务并发或重复执行。
  3. 状态机清晰Pending → Running → Success/Failed 配合条件更新,保证并发安全。
  4. 扩展成本低:新增任务类型只需实现一个执行器类并注册到工厂。
  5. 可追溯:每次执行都有独立记录,并与 Hangfire JobId 关联,便于排查。

如果你的项目也有类似的定时任务需求,不妨参考这套模式:先抽象任务定义与执行记录,再对 Hangfire 做一层薄封装,最后用工厂把业务执行器管理起来。这样写出来的调度系统既稳定,又容易维护。