From bbdea76459c3141b30fd302d36047f858b21a44e Mon Sep 17 00:00:00 2001 From: maliming Date: Fri, 3 Jul 2026 14:13:54 +0800 Subject: [PATCH] Add background job worker enhancements (dedicated workers, parallel execution, job retention) --- .../AbpBackgroundJobWorkerOptions.cs | 155 +++++++- .../BackgroundJobs/AbpBackgroundJobsModule.cs | 12 +- .../BackgroundJobCleanupWorker.cs | 66 ++++ .../Abp/BackgroundJobs/BackgroundJobInfo.cs | 8 + .../BackgroundJobs/BackgroundJobNameFilter.cs | 71 ++++ .../BackgroundJobNameFilterMode.cs | 19 + .../Abp/BackgroundJobs/BackgroundJobWorker.cs | 353 +++++++++++++++--- .../BackgroundJobWorkerConfiguration.cs | 37 ++ .../BackgroundJobWorkerManager.cs | 131 +++++++ .../DedicatedWorkerDefinition.cs | 27 ++ .../Abp/BackgroundJobs/IBackgroundJobStore.cs | 26 +- .../BackgroundJobs/IBackgroundJobWorker.cs | 23 +- .../InMemoryBackgroundJobStore.cs | 35 +- 13 files changed, 893 insertions(+), 70 deletions(-) create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker.cs create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilter.cs create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilterMode.cs create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerConfiguration.cs create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerManager.cs create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/DedicatedWorkerDefinition.cs diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerOptions.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerOptions.cs index b339508028..f0dc0aa500 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerOptions.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerOptions.cs @@ -1,4 +1,8 @@ -namespace Volo.Abp.BackgroundJobs; +using System; +using System.Collections.Generic; +using System.Linq; + +namespace Volo.Abp.BackgroundJobs; public class AbpBackgroundJobWorkerOptions { @@ -15,6 +19,7 @@ public class AbpBackgroundJobWorkerOptions /// /// Maximum count of jobs to fetch from data store in one loop. + /// Also used as the batch size for the retention cleanup deletions (see ). /// Default: 1000. /// public int MaxJobFetchCount { get; set; } @@ -42,7 +47,61 @@ public class AbpBackgroundJobWorkerOptions /// Distributed lock name for the worker. /// Default value: "AbpBackgroundJobWorker". /// - public string DistributedLockName { get; set; } + public string DistributedLockName { get; set; } + + /// + /// When set to true, a successfully completed job is kept (its is set) + /// instead of being deleted, so it can be used for auditing/history. Completed jobs are excluded from the + /// waiting jobs query and are removed by the retention cleanup (see ). + /// Default value: false (successful jobs are deleted). + /// + public bool StoreSuccessfulJobs { get; set; } + + /// + /// How long a kept (successfully completed) job is retained before the cleanup deletes it. + /// Only relevant when is true. Set to null to keep them forever (no automatic cleanup). + /// Default value: 7 days. + /// + public TimeSpan? SuccessfulJobRetentionTime { get; set; } + + /// + /// Interval (as milliseconds) between cleanup runs that delete retained successful jobs older than . + /// Default value: 3,600,000 (1 hour). + /// + public int CleanSuccessfulJobsPeriod { get; set; } + + /// + /// Distributed lock name for the cleanup worker. + /// Default value: "AbpBackgroundJobCleanup". + /// + public string CleanupDistributedLockName { get; set; } + + /// + /// Dedicated workers, each processing only the configured job types with its own distributed lock. + /// When this list is not empty, an additional default worker is started that processes all + /// other job types. When it is empty, a single default worker processes all job types (the default behavior). + /// Use to add a dedicated worker. + /// + public List WorkerConfigurations { get; } + + /// + /// Maximum number of jobs that a single worker executes in parallel within one poll cycle. + /// When it is 1 (default), jobs are executed sequentially under a single worker distributed lock + /// (the default behavior). When greater than 1, the worker distributed lock is not used; instead + /// each job is claimed with its own distributed lock so multiple application instances can execute + /// different jobs concurrently. + /// This value should be configured consistently across all application instances: mixing sequential + /// (worker lock) and parallel (per-job lock) instances removes the common mutual exclusion and may + /// let the same job run on more than one instance. + /// Default value: 1. + /// + public int MaxParallelJobExecutionCount { get; set; } + + /// + /// Prefix of the per-job distributed lock name used when is greater than 1. + /// Default value: "AbpBackgroundJob:". + /// + public string PerJobDistributedLockPrefix { get; set; } public AbpBackgroundJobWorkerOptions() { @@ -52,5 +111,97 @@ public class AbpBackgroundJobWorkerOptions DefaultTimeout = 172800; DefaultWaitFactor = 2.0; DistributedLockName = "AbpBackgroundJobWorker"; + WorkerConfigurations = new List(); + MaxParallelJobExecutionCount = 1; + PerJobDistributedLockPrefix = "AbpBackgroundJob:"; + SuccessfulJobRetentionTime = TimeSpan.FromDays(7); + CleanSuccessfulJobsPeriod = 3600000; + CleanupDistributedLockName = "AbpBackgroundJobCleanup"; + } + + /// + /// Adds a dedicated worker that processes only the given job argument types. + /// + /// A unique distributed lock name for this worker. + /// The job argument types handled exclusively by this worker. + public AbpBackgroundJobWorkerOptions AddDedicatedWorker(string lockName, params Type[] jobArgsTypes) + { + Check.NotNullOrEmpty(jobArgsTypes, nameof(jobArgsTypes)); + + var configuration = new BackgroundJobWorkerConfiguration(lockName, jobArgsTypes.Distinct().ToArray()); + + var duplicateType = configuration.JobArgsTypes.FirstOrDefault( + type => WorkerConfigurations.Any(c => c.JobArgsTypes.Contains(type))); + if (duplicateType != null) + { + throw new AbpException( + $"The background job args type '{duplicateType.FullName}' is already assigned to a dedicated worker. " + + $"Each job type can be handled by only one dedicated worker."); + } + + if (lockName == DistributedLockName || WorkerConfigurations.Any(c => c.LockName == lockName)) + { + throw new AbpException( + $"The distributed lock name '{lockName}' is already used by another background job worker. " + + $"Each worker must have a unique lock name."); + } + + WorkerConfigurations.Add(configuration); + return this; + } + + public AbpBackgroundJobWorkerOptions AddDedicatedWorker(string lockName) + { + return AddDedicatedWorker(lockName, typeof(TArgs)); + } + + public AbpBackgroundJobWorkerOptions AddDedicatedWorker(string lockName) + { + return AddDedicatedWorker(lockName, typeof(TArgs1), typeof(TArgs2)); + } + + public AbpBackgroundJobWorkerOptions AddDedicatedWorker(string lockName) + { + return AddDedicatedWorker(lockName, typeof(TArgs1), typeof(TArgs2), typeof(TArgs3)); + } + + /// + /// Adds a dedicated worker with a lock name derived from the given job argument type names. + /// Use the AddDedicatedWorker(string lockName, ...) overloads to set an explicit lock name. + /// + /// The job argument types handled exclusively by this worker. + public AbpBackgroundJobWorkerOptions AddDedicatedWorker(params Type[] jobArgsTypes) + { + return AddDedicatedWorker(GetDedicatedWorkerLockName(jobArgsTypes), jobArgsTypes); + } + + public AbpBackgroundJobWorkerOptions AddDedicatedWorker() + { + return AddDedicatedWorker(typeof(TArgs)); + } + + public AbpBackgroundJobWorkerOptions AddDedicatedWorker() + { + return AddDedicatedWorker(typeof(TArgs1), typeof(TArgs2)); + } + + public AbpBackgroundJobWorkerOptions AddDedicatedWorker() + { + return AddDedicatedWorker(typeof(TArgs1), typeof(TArgs2), typeof(TArgs3)); + } + + protected virtual string GetDedicatedWorkerLockName(Type[] jobArgsTypes) + { + Check.NotNullOrEmpty(jobArgsTypes, nameof(jobArgsTypes)); + + if (jobArgsTypes.Any(t => t == null)) + { + throw new ArgumentException("Job args types cannot contain null.", nameof(jobArgsTypes)); + } + + // Hash the (stable, sorted) full type names so the derived lock name stays short and within the + // length limits of distributed lock providers (e.g. SQL Server sp_getapplock is limited to 255 chars). + var key = string.Join(",", jobArgsTypes.Select(t => t.FullName).Distinct().OrderBy(n => n, StringComparer.Ordinal)); + return "AbpBackgroundJobDedicatedWorker:" + key.ToMd5(); } } diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs index b3601ad1be..0f4b20cd65 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs @@ -1,4 +1,4 @@ -using System.Threading.Tasks; +using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using Volo.Abp.BackgroundWorkers; @@ -23,9 +23,15 @@ public class AbpBackgroundJobsModule : AbpModule { public override async Task OnApplicationInitializationAsync(ApplicationInitializationContext context) { - if (context.ServiceProvider.GetRequiredService>().Value.IsJobExecutionEnabled) + // The manager decides (based on options) whether and which workers to run. + await context.AddBackgroundWorkerAsync(); + + // Only register the cleanup worker when it has something to do (retained successful jobs). + var jobOptions = context.ServiceProvider.GetRequiredService>().Value; + var workerOptions = context.ServiceProvider.GetRequiredService>().Value; + if (jobOptions.IsJobExecutionEnabled && workerOptions.StoreSuccessfulJobs && workerOptions.SuccessfulJobRetentionTime != null) { - await context.AddBackgroundWorkerAsync(); + await context.AddBackgroundWorkerAsync(); } } diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker.cs new file mode 100644 index 0000000000..72c36b8ffb --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker.cs @@ -0,0 +1,66 @@ +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Volo.Abp.BackgroundWorkers; +using Volo.Abp.DistributedLocking; +using Volo.Abp.Threading; +using Volo.Abp.Timing; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// Periodically deletes retained successfully completed jobs older than +/// . +/// Only relevant when is enabled. +/// +public class BackgroundJobCleanupWorker : AsyncPeriodicBackgroundWorkerBase +{ + protected AbpBackgroundJobOptions JobOptions { get; } + + protected AbpBackgroundJobWorkerOptions WorkerOptions { get; } + + protected IAbpDistributedLock DistributedLock { get; } + + public BackgroundJobCleanupWorker( + AbpAsyncTimer timer, + IServiceScopeFactory serviceScopeFactory, + IOptions jobOptions, + IOptions workerOptions, + IAbpDistributedLock distributedLock) + : base(timer, serviceScopeFactory) + { + JobOptions = jobOptions.Value; + WorkerOptions = workerOptions.Value; + DistributedLock = distributedLock; + Timer.Period = WorkerOptions.CleanSuccessfulJobsPeriod; + } + + protected override async Task DoWorkAsync(PeriodicBackgroundWorkerContext workerContext) + { + if (!JobOptions.IsJobExecutionEnabled || + !WorkerOptions.StoreSuccessfulJobs || + WorkerOptions.SuccessfulJobRetentionTime == null) + { + return; + } + + var store = workerContext.ServiceProvider.GetRequiredService(); + var clock = workerContext.ServiceProvider.GetRequiredService(); + var completedBefore = clock.Now.Subtract(WorkerOptions.SuccessfulJobRetentionTime.Value); + + await using (var handle = await DistributedLock.TryAcquireAsync(WorkerOptions.CleanupDistributedLockName, cancellationToken: StoppingToken)) + { + if (handle == null) + { + return; + } + + int deletedCount; + do + { + deletedCount = await store.DeleteAsync(WorkerOptions.ApplicationName, completedBefore, WorkerOptions.MaxJobFetchCount, StoppingToken); + } + while (deletedCount > 0 && deletedCount >= WorkerOptions.MaxJobFetchCount && !StoppingToken.IsCancellationRequested); + } + } +} diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobInfo.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobInfo.cs index 4a9cec922c..51f73ae7b8 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobInfo.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobInfo.cs @@ -50,6 +50,14 @@ public class BackgroundJobInfo /// public virtual bool IsAbandoned { get; set; } + /// + /// The time this job was completed successfully. + /// When set, the job is kept as history and excluded from the waiting jobs query. + /// It is only set when is enabled; + /// otherwise successfully completed jobs are deleted. + /// + public virtual DateTime? CompletionTime { get; set; } + /// /// Priority of this job. /// diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilter.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilter.cs new file mode 100644 index 0000000000..06c0ebf398 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilter.cs @@ -0,0 +1,71 @@ +using System; +using System.Collections.Generic; +using System.Linq; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// Filters the waiting jobs of a background job worker by job name. +/// A worker is exactly one of: no filter (), include-only (a dedicated worker) or +/// exclude-only (the default worker in a multi-worker setup) — the two can never be combined. +/// +public class BackgroundJobNameFilter +{ + /// + /// A filter that matches every job name. + /// + public static BackgroundJobNameFilter None { get; } = new(BackgroundJobNameFilterMode.None); + + public BackgroundJobNameFilterMode Mode { get; } + + public IReadOnlyList JobNames { get; } + + public BackgroundJobNameFilter(BackgroundJobNameFilterMode mode, IReadOnlyList? jobNames = null) + { + if (!Enum.IsDefined(typeof(BackgroundJobNameFilterMode), mode)) + { + throw new ArgumentException($"Invalid background job name filter mode: {mode}", nameof(mode)); + } + + var names = jobNames?.Where(x => !x.IsNullOrWhiteSpace()).Distinct(StringComparer.Ordinal).ToList() ?? new List(); + + if (mode == BackgroundJobNameFilterMode.None && names.Count > 0) + { + throw new ArgumentException("Job names must be empty when the filter mode is None.", nameof(jobNames)); + } + + if (mode != BackgroundJobNameFilterMode.None && names.Count == 0) + { + throw new ArgumentException("Job names cannot be empty when the filter mode is Include or Exclude.", nameof(jobNames)); + } + + Mode = mode; + JobNames = names.AsReadOnly(); + } + + public static BackgroundJobNameFilter Include(IReadOnlyList jobNames) + { + return new BackgroundJobNameFilter(BackgroundJobNameFilterMode.Include, jobNames); + } + + public static BackgroundJobNameFilter Exclude(IReadOnlyList jobNames) + { + return new BackgroundJobNameFilter(BackgroundJobNameFilterMode.Exclude, jobNames); + } + + /// + /// Whether the given job name passes this filter, using an ordinal (case-sensitive) comparison for the + /// in-memory eligibility re-check. The persistent stores translate and + /// into a database query instead, so their filtering follows the database collation. + /// Job names are expected to be unique beyond case (they are derived from the type name by default). + /// + public virtual bool IsMatch(string jobName) + { + return Mode switch + { + BackgroundJobNameFilterMode.Include => JobNames.Contains(jobName, StringComparer.Ordinal), + BackgroundJobNameFilterMode.Exclude => !JobNames.Contains(jobName, StringComparer.Ordinal), + _ => true + }; + } +} diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilterMode.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilterMode.cs new file mode 100644 index 0000000000..6ed2801b44 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameFilterMode.cs @@ -0,0 +1,19 @@ +namespace Volo.Abp.BackgroundJobs; + +public enum BackgroundJobNameFilterMode : byte +{ + /// + /// No filter; all job names match. + /// + None = 0, + + /// + /// Only the job names in the filter match. + /// + Include = 1, + + /// + /// All job names except those in the filter match. + /// + Exclude = 2 +} diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs index a015e32d66..3053bb426f 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs @@ -1,17 +1,22 @@ -using System; +using System; +using System.Collections.Generic; using System.Linq; +using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; using Volo.Abp.BackgroundWorkers; +using Volo.Abp.DependencyInjection; using Volo.Abp.DistributedLocking; +using Volo.Abp.ExceptionHandling; using Volo.Abp.Threading; using Volo.Abp.Timing; namespace Volo.Abp.BackgroundJobs; -public class BackgroundJobWorker : AsyncPeriodicBackgroundWorkerBase, IBackgroundJobWorker +public class BackgroundJobWorker : IBackgroundJobWorker, ITransientDependency { protected AbpBackgroundJobOptions JobOptions { get; } @@ -19,95 +24,321 @@ public class BackgroundJobWorker : AsyncPeriodicBackgroundWorkerBase, IBackgroun protected IAbpDistributedLock DistributedLock { get; } + protected IServiceScopeFactory ServiceScopeFactory { get; } + + protected AbpAsyncTimer Timer { get; } + + public ILogger Logger { get; set; } + + protected string DistributedLockName { get; set; } = default!; + + protected BackgroundJobNameFilter JobNameFilter { get; set; } = BackgroundJobNameFilter.None; + + protected CancellationTokenSource StoppingTokenSource { get; } + + protected CancellationToken StoppingToken { get; } + public BackgroundJobWorker( AbpAsyncTimer timer, IOptions jobOptions, IOptions workerOptions, IServiceScopeFactory serviceScopeFactory, IAbpDistributedLock distributedLock) - : base( - timer, - serviceScopeFactory) { + Timer = timer; DistributedLock = distributedLock; + ServiceScopeFactory = serviceScopeFactory; WorkerOptions = workerOptions.Value; JobOptions = jobOptions.Value; + Logger = NullLogger.Instance; + Timer.Period = WorkerOptions.JobPollPeriod; + Timer.Elapsed = TimerOnElapsed; + + StoppingTokenSource = new CancellationTokenSource(); + StoppingToken = StoppingTokenSource.Token; + } + + public virtual Task StartAsync( + string? distributedLockName = null, + BackgroundJobNameFilter? jobNameFilter = null, + CancellationToken cancellationToken = default) + { + DistributedLockName = distributedLockName ?? WorkerOptions.DistributedLockName; + JobNameFilter = jobNameFilter ?? BackgroundJobNameFilter.None; + + Timer.Start(cancellationToken); + + return Task.CompletedTask; + } + + public virtual Task StopAsync(CancellationToken cancellationToken = default) + { + StoppingTokenSource.Cancel(); + Timer.Stop(cancellationToken); + StoppingTokenSource.Dispose(); + + return Task.CompletedTask; + } + + private async Task TimerOnElapsed(AbpAsyncTimer timer) + { + await RunAsync(); + } + + protected virtual async Task RunAsync() + { + using var scope = ServiceScopeFactory.CreateScope(); + + try + { + var workerContext = new PeriodicBackgroundWorkerContext(scope.ServiceProvider, StoppingToken); + + if (WorkerOptions.MaxParallelJobExecutionCount > 1) + { + await ExecuteJobsInParallelAsync(workerContext); + } + else + { + await ExecuteJobsWithWorkerLockAsync(workerContext); + } + } + catch (Exception ex) + { + await scope.ServiceProvider + .GetRequiredService() + .NotifyAsync(new ExceptionNotificationContext(ex)); + + Logger.LogException(ex); + } } - protected override async Task DoWorkAsync(PeriodicBackgroundWorkerContext workerContext) + protected virtual async Task ExecuteJobsWithWorkerLockAsync(PeriodicBackgroundWorkerContext workerContext) { - await using (var handler = await DistributedLock.TryAcquireAsync(WorkerOptions.DistributedLockName, cancellationToken: StoppingToken)) + await using (var handler = await DistributedLock.TryAcquireAsync(DistributedLockName, cancellationToken: StoppingToken)) { if (handler != null) { - var store = workerContext.ServiceProvider.GetRequiredService(); + await ExecuteWaitingJobsAsync(workerContext); + } + else + { + await WaitForNextTryAsync(); + } + } + } + + protected virtual async Task ExecuteWaitingJobsAsync(PeriodicBackgroundWorkerContext workerContext) + { + var store = workerContext.ServiceProvider.GetRequiredService(); + + var waitingJobs = await GetWaitingJobsAsync(workerContext, store); - var waitingJobs = await store.GetWaitingJobsAsync(WorkerOptions.ApplicationName, WorkerOptions.MaxJobFetchCount); + if (!waitingJobs.Any()) + { + return; + } - if (!waitingJobs.Any()) + var jobExecuter = workerContext.ServiceProvider.GetRequiredService(); + var clock = workerContext.ServiceProvider.GetRequiredService(); + var serializer = workerContext.ServiceProvider.GetRequiredService(); + + foreach (var jobInfo in waitingJobs) + { + await TryExecuteJobAsync(workerContext, store, jobInfo, jobExecuter, clock, serializer); + } + } + + /// + /// Executes waiting jobs in parallel across application instances, up to + /// jobs per cycle. + /// + protected virtual async Task ExecuteJobsInParallelAsync(PeriodicBackgroundWorkerContext workerContext) + { + var store = workerContext.ServiceProvider.GetRequiredService(); + + var waitingJobs = await GetWaitingJobsAsync(workerContext, store); + + if (!waitingJobs.Any()) + { + return; + } + + var runningTasks = new List(); + + // Await already-started jobs even if acquiring a lock for a later job throws, + // so no claimed job is left running detached from this cycle. + try + { + foreach (var jobInfo in waitingJobs) + { + if (runningTasks.Count >= WorkerOptions.MaxParallelJobExecutionCount || StoppingToken.IsCancellationRequested) { - return; + break; } - var jobExecuter = workerContext.ServiceProvider.GetRequiredService(); - var clock = workerContext.ServiceProvider.GetRequiredService(); - var serializer = workerContext.ServiceProvider.GetRequiredService(); - - foreach (var jobInfo in waitingJobs) + var handle = await DistributedLock.TryAcquireAsync(GetPerJobDistributedLockName(jobInfo), cancellationToken: StoppingToken); + if (handle == null) { - jobInfo.TryCount++; - jobInfo.LastTryTime = clock.Now; - - try - { - var jobConfiguration = JobOptions.GetJob(jobInfo.JobName); - var jobArgs = serializer.Deserialize(jobInfo.JobArgs, jobConfiguration.ArgsType); - var context = new JobExecutionContext( - workerContext.ServiceProvider, - jobConfiguration.JobType, - jobArgs, - workerContext.CancellationToken); - - try - { - await jobExecuter.ExecuteAsync(context); - - await store.DeleteAsync(jobInfo.Id); - } - catch (BackgroundJobExecutionException) - { - var nextTryTime = CalculateNextTryTime(jobInfo, clock); - - if (nextTryTime.HasValue) - { - jobInfo.NextTryTime = nextTryTime.Value; - } - else - { - jobInfo.IsAbandoned = true; - } - - await TryUpdateAsync(store, jobInfo); - } - } - catch (Exception ex) - { - Logger.LogException(ex); - jobInfo.IsAbandoned = true; - await TryUpdateAsync(store, jobInfo); - } + // Another instance is already processing this job. + continue; } + + runningTasks.Add(ExecuteClaimedJobAsync(jobInfo, handle)); } - else + } + finally + { + await Task.WhenAll(runningTasks); + } + } + + protected virtual async Task ExecuteClaimedJobAsync(BackgroundJobInfo jobInfo, IAbpDistributedLockHandle handle) + { + await using (handle) + { + // Each concurrently executed job runs in its own service scope so that scoped services + // (e.g. the DbContext and the unit of work) are not shared across parallel jobs. + using var scope = ServiceScopeFactory.CreateScope(); + + try { - try + var workerContext = new PeriodicBackgroundWorkerContext(scope.ServiceProvider, StoppingToken); + var store = scope.ServiceProvider.GetRequiredService(); + var clock = scope.ServiceProvider.GetRequiredService(); + + // Re-read under the lock: another instance may have completed, abandoned or rescheduled this job + // between fetching the waiting list and acquiring the per-job lock. + var currentJobInfo = await store.FindAsync(jobInfo.Id); + if (!IsJobEligible(currentJobInfo, clock)) { - await Task.Delay(WorkerOptions.JobPollPeriod * 12, StoppingToken); + return; } - catch (TaskCanceledException) { } + + var jobExecuter = scope.ServiceProvider.GetRequiredService(); + var serializer = scope.ServiceProvider.GetRequiredService(); + + await TryExecuteJobAsync(workerContext, store, currentJobInfo, jobExecuter, clock, serializer); + } + catch (Exception ex) + { + await scope.ServiceProvider + .GetRequiredService() + .NotifyAsync(new ExceptionNotificationContext(ex)); + + Logger.LogException(ex); + } + } + } + + protected virtual bool IsJobEligible(BackgroundJobInfo? jobInfo, IClock clock) + { + return jobInfo != null && + jobInfo.ApplicationName == WorkerOptions.ApplicationName && + !jobInfo.IsAbandoned && + jobInfo.CompletionTime == null && + jobInfo.NextTryTime <= clock.Now && + JobNameFilter.IsMatch(jobInfo.JobName); + } + + protected virtual string GetPerJobDistributedLockName(BackgroundJobInfo jobInfo) + { + return WorkerOptions.PerJobDistributedLockPrefix + jobInfo.Id; + } + + protected virtual async Task> GetWaitingJobsAsync( + PeriodicBackgroundWorkerContext workerContext, + IBackgroundJobStore store) + { + return await store.GetWaitingJobsAsync( + WorkerOptions.ApplicationName, + WorkerOptions.MaxJobFetchCount, + JobNameFilter); + } + + protected virtual async Task TryExecuteJobAsync( + PeriodicBackgroundWorkerContext workerContext, + IBackgroundJobStore store, + BackgroundJobInfo jobInfo, + IBackgroundJobExecuter jobExecuter, + IClock clock, + IBackgroundJobSerializer serializer) + { + jobInfo.TryCount++; + jobInfo.LastTryTime = clock.Now; + + try + { + var jobConfiguration = JobOptions.GetJob(jobInfo.JobName); + var jobArgs = serializer.Deserialize(jobInfo.JobArgs, jobConfiguration.ArgsType); + var context = new JobExecutionContext( + workerContext.ServiceProvider, + jobConfiguration.JobType, + jobArgs, + workerContext.CancellationToken); + + try + { + await jobExecuter.ExecuteAsync(context); + + await HandleJobSuccessAsync(store, jobInfo, clock); + } + catch (BackgroundJobExecutionException) + { + await HandleJobFailureAsync(store, jobInfo, clock); } } + catch (Exception ex) + { + await HandleJobErrorAsync(store, jobInfo, ex); + } + } + + protected virtual async Task HandleJobSuccessAsync(IBackgroundJobStore store, BackgroundJobInfo jobInfo, IClock clock) + { + if (WorkerOptions.StoreSuccessfulJobs) + { + // Keep the job as history: mark it completed instead of deleting. It is then excluded from the + // waiting jobs query and removed later by the retention cleanup. + jobInfo.CompletionTime = clock.Now; + await store.UpdateAsync(jobInfo); + } + else + { + await store.DeleteAsync(jobInfo.Id); + } + } + + protected virtual async Task HandleJobFailureAsync(IBackgroundJobStore store, BackgroundJobInfo jobInfo, IClock clock) + { + var nextTryTime = CalculateNextTryTime(jobInfo, clock); + + if (nextTryTime.HasValue) + { + jobInfo.NextTryTime = nextTryTime.Value; + } + else + { + jobInfo.IsAbandoned = true; + } + + await TryUpdateAsync(store, jobInfo); + } + + protected virtual async Task HandleJobErrorAsync(IBackgroundJobStore store, BackgroundJobInfo jobInfo, Exception ex) + { + Logger.LogException(ex); + jobInfo.IsAbandoned = true; + await TryUpdateAsync(store, jobInfo); + } + + protected virtual async Task WaitForNextTryAsync() + { + try + { + await Task.Delay(WorkerOptions.JobPollPeriod * 12, StoppingToken); + } + catch (TaskCanceledException) { } } protected virtual async Task TryUpdateAsync(IBackgroundJobStore store, BackgroundJobInfo jobInfo) diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerConfiguration.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerConfiguration.cs new file mode 100644 index 0000000000..848a2a6a1c --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerConfiguration.cs @@ -0,0 +1,37 @@ +using System; +using System.Collections.Generic; +using System.Linq; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// Configuration of a dedicated that processes only specific job types. +/// +public class BackgroundJobWorkerConfiguration +{ + /// + /// A unique distributed lock name for this worker. It must be different from the names used by other workers. + /// It is used to serialize the worker across application instances when + /// is 1; in parallel mode + /// (greater than 1) jobs are claimed with per-job locks instead and this lock is not acquired. + /// + public string LockName { get; } + + /// + /// The job argument types that are processed exclusively by this worker. + /// + public IReadOnlyList JobArgsTypes { get; } + + public BackgroundJobWorkerConfiguration(string lockName, params Type[] jobArgsTypes) + { + LockName = Check.NotNullOrWhiteSpace(lockName, nameof(lockName)); + Check.NotNullOrEmpty(jobArgsTypes, nameof(jobArgsTypes)); + + if (jobArgsTypes.Any(t => t == null)) + { + throw new ArgumentException("Job args types cannot contain null.", nameof(jobArgsTypes)); + } + + JobArgsTypes = jobArgsTypes.ToList(); + } +} diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerManager.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerManager.cs new file mode 100644 index 0000000000..929c4803da --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorkerManager.cs @@ -0,0 +1,131 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Volo.Abp.BackgroundWorkers; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// Owns and controls the background job workers. +/// When no is configured, a single +/// default worker processes all jobs. Otherwise, one dedicated worker is started per configuration +/// (each with its own distributed lock and job-type filter) plus a default worker for the remaining jobs. +/// The workers are resolved from DI, so a replaced is respected. +/// +public class BackgroundJobWorkerManager : IBackgroundWorker +{ + protected AbpBackgroundJobOptions JobOptions { get; } + + protected AbpBackgroundJobWorkerOptions WorkerOptions { get; } + + protected IServiceProvider ServiceProvider { get; } + + protected List Workers { get; } + + public BackgroundJobWorkerManager( + IOptions jobOptions, + IOptions workerOptions, + IServiceProvider serviceProvider) + { + JobOptions = jobOptions.Value; + WorkerOptions = workerOptions.Value; + ServiceProvider = serviceProvider; + Workers = new List(); + } + + public virtual async Task StartAsync(CancellationToken cancellationToken = default) + { + if (!JobOptions.IsJobExecutionEnabled) + { + return; + } + + if (!WorkerOptions.WorkerConfigurations.Any()) + { + await StartWorkerAsync(cancellationToken: cancellationToken); + return; + } + + // AddDedicatedWorker already rejects duplicate job types and lock names eagerly. This is the backstop + // for what can only be known here: two different args types that resolve to the same job name. + // Validate all configurations first, so a misconfiguration does not leave already-started workers running. + var dedicatedWorkers = new List(); + var allDedicatedJobNames = new List(); + + // The default worker uses WorkerOptions.DistributedLockName, so dedicated workers must not reuse it. + var lockNames = new List { WorkerOptions.DistributedLockName }; + + foreach (var configuration in WorkerOptions.WorkerConfigurations) + { + var jobNames = configuration.JobArgsTypes + .Select(GetJobName) + .Distinct() + .ToList(); + + var alreadyConfigured = jobNames.Intersect(allDedicatedJobNames).ToList(); + if (alreadyConfigured.Any()) + { + throw new AbpException( + $"The following background job(s) are configured for more than one dedicated worker: {string.Join(", ", alreadyConfigured)}. " + + $"Each job type can be handled by only one dedicated worker."); + } + + if (lockNames.Contains(configuration.LockName)) + { + throw new AbpException( + $"The distributed lock name '{configuration.LockName}' is used by more than one background job worker " + + $"(the default worker uses '{WorkerOptions.DistributedLockName}'). Each worker must have a unique lock name to run independently."); + } + + lockNames.Add(configuration.LockName); + allDedicatedJobNames.AddRange(jobNames); + dedicatedWorkers.Add(new DedicatedWorkerDefinition(configuration.LockName, jobNames)); + } + + foreach (var dedicatedWorker in dedicatedWorkers) + { + await StartWorkerAsync(dedicatedWorker.LockName, BackgroundJobNameFilter.Include(dedicatedWorker.JobNames), cancellationToken); + } + + // Default worker processes every job that is not handled by a dedicated worker. + await StartWorkerAsync(null, BackgroundJobNameFilter.Exclude(allDedicatedJobNames), cancellationToken); + } + + protected virtual string GetJobName(Type argsType) + { + try + { + return JobOptions.GetJob(argsType).JobName; + } + catch (AbpException ex) + { + throw new AbpException( + $"No background job is registered for the args type '{argsType.FullName}' configured via AddDedicatedWorker. " + + $"Register the job before configuring a dedicated worker for it.", ex); + } + } + + protected virtual async Task StartWorkerAsync( + string? distributedLockName = null, + BackgroundJobNameFilter? jobNameFilter = null, + CancellationToken cancellationToken = default) + { + var worker = ServiceProvider.GetRequiredService(); + await worker.StartAsync(distributedLockName, jobNameFilter, cancellationToken); + Workers.Add(worker); + } + + public virtual async Task StopAsync(CancellationToken cancellationToken = default) + { + foreach (var worker in Workers) + { + await worker.StopAsync(cancellationToken); + } + + Workers.Clear(); + } +} diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/DedicatedWorkerDefinition.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/DedicatedWorkerDefinition.cs new file mode 100644 index 0000000000..ba98c2c1dc --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/DedicatedWorkerDefinition.cs @@ -0,0 +1,27 @@ +using System.Collections.Generic; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// A validated, ready-to-start dedicated worker: the distributed lock name it runs under and the +/// resolved job names it is responsible for. Built by from a +/// after all configurations have been validated. +/// +public class DedicatedWorkerDefinition +{ + /// + /// The distributed lock name this worker runs under. + /// + public string LockName { get; } + + /// + /// The resolved job names this worker is responsible for. + /// + public IReadOnlyList JobNames { get; } + + public DedicatedWorkerDefinition(string lockName, IReadOnlyList jobNames) + { + LockName = lockName; + JobNames = jobNames; + } +} diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobStore.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobStore.cs index 7fcf8cc5b6..421ec55434 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobStore.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobStore.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Threading; using System.Threading.Tasks; namespace Volo.Abp.BackgroundJobs; @@ -24,7 +25,7 @@ public interface IBackgroundJobStore /// /// Gets waiting jobs. It should get jobs based on these: - /// Conditions: ApplicationName is applicationName And !IsAbandoned And NextTryTime <= Clock.Now. + /// Conditions: ApplicationName is applicationName And !IsAbandoned And CompletionTime == null And NextTryTime <= Clock.Now. /// Order by: Priority DESC, TryCount ASC, NextTryTime ASC. /// Maximum result: . /// @@ -32,12 +33,35 @@ public interface IBackgroundJobStore /// Maximum result count. Task> GetWaitingJobsAsync(string? applicationName, int maxResultCount); + /// + /// Gets waiting jobs (see ), additionally filtered by job name + /// according to . + /// + /// Application name. + /// Maximum result count. + /// Job name filter. When null, no job name filter is applied. + Task> GetWaitingJobsAsync( + string? applicationName, + int maxResultCount, + BackgroundJobNameFilter? jobNameFilter); + /// /// Deletes a job. /// /// The Job Unique Identifier. Task DeleteAsync(Guid jobId); + /// + /// Deletes successfully completed jobs ( is set) of the given + /// application that completed before . Used by the retention cleanup. + /// + /// The number of deleted jobs. + Task DeleteAsync( + string? applicationName, + DateTime completedBefore, + int maxResultCount, + CancellationToken cancellationToken = default); + /// /// Updates a job. /// diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobWorker.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobWorker.cs index ecae732b8d..9aa4ac25a8 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobWorker.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobWorker.cs @@ -1,8 +1,27 @@ -using Volo.Abp.BackgroundWorkers; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; namespace Volo.Abp.BackgroundJobs; -public interface IBackgroundJobWorker : IBackgroundWorker +/// +/// A background job worker that polls and executes waiting jobs. +/// Instances are created, configured and started by . +/// +public interface IBackgroundJobWorker { + /// + /// Starts this worker. + /// + /// + /// Distributed lock name for this worker. When null, is used. + /// + /// Filters the jobs this worker processes by name. When null, all jobs are processed. + /// Cancellation token. + Task StartAsync( + string? distributedLockName = null, + BackgroundJobNameFilter? jobNameFilter = null, + CancellationToken cancellationToken = default); + Task StopAsync(CancellationToken cancellationToken = default); } diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/InMemoryBackgroundJobStore.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/InMemoryBackgroundJobStore.cs index 85916e8c37..e487891435 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/InMemoryBackgroundJobStore.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/InMemoryBackgroundJobStore.cs @@ -2,6 +2,7 @@ using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; +using System.Threading; using System.Threading.Tasks; using Volo.Abp.DependencyInjection; using Volo.Abp.Timing; @@ -37,9 +38,20 @@ public class InMemoryBackgroundJobStore : IBackgroundJobStore, ISingletonDepende public virtual Task> GetWaitingJobsAsync(string? applicationName, int maxResultCount) { + return GetWaitingJobsAsync(applicationName, maxResultCount, null); + } + + public virtual Task> GetWaitingJobsAsync( + string? applicationName, + int maxResultCount, + BackgroundJobNameFilter? jobNameFilter) + { + var filter = jobNameFilter ?? BackgroundJobNameFilter.None; + var waitingJobs = _jobs.Values .Where(t => t.ApplicationName == applicationName) - .Where(t => !t.IsAbandoned && t.NextTryTime <= Clock.Now) + .Where(t => !t.IsAbandoned && t.CompletionTime == null && t.NextTryTime <= Clock.Now) + .Where(t => filter.IsMatch(t.JobName)) .OrderByDescending(t => t.Priority) .ThenBy(t => t.TryCount) .ThenBy(t => t.NextTryTime) @@ -57,6 +69,27 @@ public class InMemoryBackgroundJobStore : IBackgroundJobStore, ISingletonDepende return Task.CompletedTask; } + public virtual Task DeleteAsync( + string? applicationName, + DateTime completedBefore, + int maxResultCount, + CancellationToken cancellationToken = default) + { + var idsToDelete = _jobs.Values + .Where(t => t.ApplicationName == applicationName) + .Where(t => t.CompletionTime != null && t.CompletionTime < completedBefore) + .Take(maxResultCount) + .Select(t => t.Id) + .ToList(); + + foreach (var id in idsToDelete) + { + _jobs.TryRemove(id, out _); + } + + return Task.FromResult(idsToDelete.Count); + } + public virtual Task UpdateAsync(BackgroundJobInfo jobInfo) { if (jobInfo.IsAbandoned)