如何防止Hangfire中正在执行的同参数周期性任务重复入队?
解决Hangfire周期性任务重复入队问题
针对你遇到的长耗时任务被Hangfire每分钟重复入队导致队列堆积的问题,结合你提出的「同一方法不同systemCode可并行、同一方法同一systemCode不可并行」规则,提供以下两种可行方案:
方案一:统一调度器控制入队逻辑
替换原有的30个独立周期性任务,改为一个统一的调度任务,每分钟触发一次,由调度器负责检查每个任务实例的运行状态,仅在无运行实例时手动入队。
1. 创建统一调度方法
public void IntegrationTaskScheduler() { // 获取所有任务配置(替换为你实际获取30个任务的逻辑) var allTaskConfigs = GetAllHangfireTaskConfigs(); foreach (var taskConfig in allTaskConfigs) { // 检查当前任务实例是否正在运行 if (!IsTaskRunning(taskConfig.MethodName, taskConfig.SystemCode)) { // 手动入队任务 EnqueueTask(taskConfig.MethodName, taskConfig.SystemCode); } } } // 检查任务是否正在运行(基于分布式锁实现) private bool IsTaskRunning(string methodName, string systemCode) { var lockKey = $"Hangfire:TaskLock:{methodName}:{systemCode}"; var storage = JobStorage.Current; try { // 尝试获取锁,超时时间设为极短,仅用于判断锁是否存在 using (storage.GetDistributedLock(lockKey, TimeSpan.FromMilliseconds(100))) { return false; } } catch (TimeoutException) { // 锁已存在,说明任务正在执行 return true; } } // 根据方法名和参数入队任务(替换为你实际的任务入队逻辑) private void EnqueueTask(string methodName, string systemCode) { switch (methodName) { case nameof(CMS1Integration): BackgroundJob.Enqueue(() => CMS1Integration(systemCode)); break; case nameof(CMS2Integration): BackgroundJob.Enqueue(() => CMS2Integration(systemCode)); break; case nameof(CMS3Integration): BackgroundJob.Enqueue(() => CMS3Integration(systemCode)); break; // 其他任务方法... } }
2. 添加唯一的周期性调度任务
// 替换原有的30个RecurringJob.AddOrUpdate调用,仅添加这一个调度任务 RecurringJob.AddOrUpdate( "IntegrationTaskScheduler", () => IntegrationTaskScheduler(), "*/1 * * * *", TimeZoneInfo.Local);
3. 任务方法内添加分布式锁(双重保障)
为避免调度器检查与任务启动之间的时间差导致重复入队,在每个任务方法内添加分布式锁:
public void CMS1Integration(string systemCode) { var lockKey = $"Hangfire:TaskLock:{nameof(CMS1Integration)}:{systemCode}"; // 锁过期时间设为任务最大预计执行时间(比如1小时),防止任务异常退出导致锁残留 using (JobStorage.Current.GetDistributedLock(lockKey, TimeSpan.FromHours(1))) { // 你的任务执行逻辑 } } // CMS2Integration、CMS3Integration方法同理
方案二:自定义RecurringJob过滤器拦截重复入队
保留原有的30个周期性任务,通过自定义过滤器在RecurringJob触发时检查任务运行状态,若任务正在执行则取消本次入队。
1. 自定义跳过重复任务过滤器
public class SkipIfTaskRunningFilter : IRecurringJobFilter { public void OnRecurringJobScheduled(RecurringJobScheduledContext context) { var job = context.RecurringJob.Job; // 提取方法名和systemCode参数(根据你的任务参数结构调整) var methodName = job.Method.Name; var systemCode = job.Args.FirstOrDefault() as string; if (string.IsNullOrEmpty(systemCode)) return; var lockKey = $"Hangfire:TaskLock:{methodName}:{systemCode}"; var storage = JobStorage.Current; try { // 尝试获取锁,超时时间极短,仅用于判断是否有任务在运行 using (storage.GetDistributedLock(lockKey, TimeSpan.FromMilliseconds(50))) { // 锁可用,允许入队 } } catch (TimeoutException) { // 锁已存在,取消本次入队 context.Canceled = true; } } public void OnRecurringJobUpdated(RecurringJobUpdatedContext context) { } public void OnRecurringJobDeleted(RecurringJobDeletedContext context) { } }
2. 添加RecurringJob时绑定过滤器
修改原有的RecurringJob添加逻辑,绑定自定义过滤器:
foreach (var hangFireCmd in allHangFireCmds) { RecurringJob.AddOrUpdate( hangFireCmd.Name, GetAction(hangFireCmd), hangFireCmd.Cron, TimeZoneInfo.Local, new RecurringJobOptions { JobFilters = new JobFilterCollection { new SkipIfTaskRunningFilter() } }); }
3. 任务方法内添加分布式锁
同方案一的步骤3,在每个任务方法内添加分布式锁,确保同一实例不会并行执行。
关键注意事项
- 锁的过期时间必须设置为大于任务的最长预计执行时间,避免任务正常执行过程中锁过期,导致新任务被错误入队。
- 若使用SQL Server作为Hangfire存储,分布式锁的实现依赖于数据库事务,性能略低于Redis存储;建议优先使用Redis作为Hangfire的存储后端,提升锁操作的效率。
内容的提问来源于stack exchange,提问作者markzzz
相关产品推荐
相关产品推荐

