在企业级后台系统中,定时任务几乎无处不在:定时采集日志、定时解析文件、定时同步数据。简单的
BackgroundService 或 Timer 可以跑通
demo,但一到生产环境就会遇到一堆问题:任务重复执行怎么办?执行到一半崩溃了如何恢复?手动触发和定时触发如何统一?新加一种任务类型要不要改调度代码?
这篇文章结合一个真实的三层 .NET 项目实践,介绍如何基于 Hangfire 搭建一套带幂等控制、状态机、执行器工厂的任务调度系统。代码分层清晰,扩展新任务类型时只需要新增一个执行器类,不用碰核心调度逻辑。
一、整体架构:三层协作
整个调度系统由三层协作完成:
- 调度层:
NeTaskHangfireScheduler,封装 Hangfire 的IRecurringJobManager与IBackgroundJobClient,负责任务的注册、移除、立即触发、延迟触发。 - 执行层:
NeTaskExecutor,负责把 Hangfire 的作业转换为一次真正的业务执行,包括前置校验、状态流转、异常捕获。 - 持久层:
NeTask与NeTaskExecution,前者保存任务定义(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);
}
}
注意两个设计点:
- 时区使用
TimeZoneInfo.Local,避免 Hangfire 默认 UTC 导致 Cron 表达式在国内服务器上错位一小时。 - 触发目标都是静态方法。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 分钟缓冲,避免用户把“立即执行”误配成“过去的时间”。
八、扩展新任务类型:四步即可
以新增“清理过期日志”任务为例:
- 定义 Payload:
public class CleanupTaskPayload : ITaskPayload
{
public int RetentionDays { get; set; }
}
- 在枚举中新增类型:
public enum NeTaskType
{
...
[Description("清理过期日志")] Cleanup = 5
}
- 实现执行器:
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 = "清理完成" };
}
}
- 在工厂注册:
new CleanupTaskExecutor()
无需改动 Hangfire 调度器、执行器入口、状态机代码。
九、总结
这套设计的核心收益:
- 调度与执行解耦:Hangfire 只负责触发,业务执行逻辑独立演化。
- 幂等防重:通过
IdempotencyKey+ 数据库唯一约束,避免同一任务并发或重复执行。 - 状态机清晰:
Pending → Running → Success/Failed配合条件更新,保证并发安全。 - 扩展成本低:新增任务类型只需实现一个执行器类并注册到工厂。
- 可追溯:每次执行都有独立记录,并与 Hangfire JobId 关联,便于排查。
如果你的项目也有类似的定时任务需求,不妨参考这套模式:先抽象任务定义与执行记录,再对 Hangfire 做一层薄封装,最后用工厂把业务执行器管理起来。这样写出来的调度系统既稳定,又容易维护。