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)