From 8e97f7ae7932ff5f769a2dccf25d9970317cd86b Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 9 Jul 2021 11:17:58 +0200 Subject: [PATCH] Refactor quartz job scheduler into async background service When a system has many suspended workflows, it takes a long time for Quartz to schedule each trigger, blocking application startup. Moving this logic to a hosted service allows this process to happen in the background, not blocking startup --- .../CommonTemporalActivityServices.cs | 5 +- .../StartJobs.cs | 54 ++++++++++++++----- .../QuartzWorkflowDefinitionScheduler.cs | 18 ++++--- .../QuartzWorkflowInstanceScheduler.cs | 21 +++++--- .../IScopedBackgroundService.cs | 10 ++++ .../HostedServices/ScopedBackgroundService.cs | 24 +++++++++ 6 files changed, 102 insertions(+), 30 deletions(-) rename src/activities/Elsa.Activities.Temporal.Common/{StartupTasks => HostedServices}/StartJobs.cs (54%) create mode 100644 src/core/Elsa.Core/HostedServices/IScopedBackgroundService.cs create mode 100644 src/core/Elsa.Core/HostedServices/ScopedBackgroundService.cs diff --git a/src/activities/Elsa.Activities.Temporal.Common/Extensions/CommonTemporalActivityServices.cs b/src/activities/Elsa.Activities.Temporal.Common/Extensions/CommonTemporalActivityServices.cs index cf762a0c1..7f13c87ff 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/Extensions/CommonTemporalActivityServices.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/Extensions/CommonTemporalActivityServices.cs @@ -1,8 +1,9 @@ using System; using Elsa.Activities.Temporal.Common.Bookmarks; using Elsa.Activities.Temporal.Common.Handlers; +using Elsa.Activities.Temporal.Common.HostedServices; using Elsa.Activities.Temporal.Common.Options; -using Elsa.Activities.Temporal.Common.StartupTasks; +using Elsa.HostedServices; using Elsa.Runtime; using Microsoft.Extensions.DependencyInjection; @@ -31,7 +32,7 @@ namespace Elsa.Activities.Temporal options.Services .AddNotificationHandlers(typeof(RemoveScheduledTriggers)) - .AddStartupTask() + .AddHostedService>() .AddBookmarkProvider() .AddBookmarkProvider() .AddBookmarkProvider(); diff --git a/src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs b/src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs similarity index 54% rename from src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs rename to src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs index c76f3d444..ba70b770e 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/StartupTasks/StartJobs.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs @@ -2,70 +2,98 @@ using System.Threading.Tasks; using Elsa.Activities.Temporal.Common.Bookmarks; using Elsa.Activities.Temporal.Common.Services; +using Elsa.HostedServices; using Elsa.Services; using Elsa.Services.Bookmarks; +using Microsoft.Extensions.Logging; +using Open.Linq.AsyncExtensions; -namespace Elsa.Activities.Temporal.Common.StartupTasks +namespace Elsa.Activities.Temporal.Common.HostedServices { /// /// Starts jobs based on workflow instances blocked on a Timer, Cron or StartAt activity. /// - public class StartJobs : IStartupTask + public class StartJobs : IScopedBackgroundService { // TODO: Figure out how to start jobs across multiple tenants / how to get a list of all tenants. private const string? TenantId = default; private readonly IBookmarkFinder _bookmarkFinder; private readonly IWorkflowInstanceScheduler _workflowScheduler; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ILogger _logger; - public StartJobs(IBookmarkFinder bookmarkFinder, IWorkflowInstanceScheduler workflowScheduler) + public StartJobs(IBookmarkFinder bookmarkFinder, IWorkflowInstanceScheduler workflowScheduler, IDistributedLockProvider distributedLockProvider, ILogger logger) { _bookmarkFinder = bookmarkFinder; _workflowScheduler = workflowScheduler; + _distributedLockProvider = distributedLockProvider; + _logger = logger; } - public int Order => 2000; - - public async Task ExecuteAsync(CancellationToken cancellationToken) + public async Task ExecuteAsync(CancellationToken stoppingToken) { - await ScheduleTimerEventWorkflowsAsync(cancellationToken); - await ScheduleCronEventWorkflowsAsync(cancellationToken); - await ScheduleStartAtWorkflowsAsync(cancellationToken); + await using var handle = await _distributedLockProvider.AcquireLockAsync(nameof(StartJobs), null, stoppingToken); + + if (handle == null) + return; + + await ScheduleTimerEventWorkflowsAsync(stoppingToken); + await ScheduleCronEventWorkflowsAsync(stoppingToken); + await ScheduleStartAtWorkflowsAsync(stoppingToken); } private async Task ScheduleStartAtWorkflowsAsync(CancellationToken cancellationToken) { // Schedule workflow instances that are blocked on a start-at. - var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(tenantId: TenantId, cancellationToken: cancellationToken); + var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(tenantId: TenantId, cancellationToken: cancellationToken).ToList(); + + _logger.LogDebug("Found {BookmarkResultCount} bookmarks for StartAt", bookmarkResults.Count); + var index = 0; foreach (var result in bookmarkResults) { var bookmark = (StartAtBookmark) result.Bookmark; await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, bookmark.ExecuteAt, null, cancellationToken); + + index++; + _logger.LogDebug("Scheduled {CurrentBookmarkIndex} of {BookmarkResultCount}", index, bookmarkResults.Count); } } private async Task ScheduleTimerEventWorkflowsAsync(CancellationToken cancellationToken) { // Schedule workflow instances that are blocked on a timer. - var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(tenantId: TenantId, cancellationToken: cancellationToken); + var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(tenantId: TenantId, cancellationToken: cancellationToken).ToList(); + + _logger.LogDebug("Found {BookmarkResultCount} bookmarks for Timer", bookmarkResults.Count); + var index = 0; foreach (var result in bookmarkResults) { var bookmark = (TimerBookmark) result.Bookmark; await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, bookmark.ExecuteAt, null, cancellationToken); + + index++; + _logger.LogDebug("Scheduled {CurrentBookmarkIndex} of {BookmarkResultCount}", index, bookmarkResults.Count); } } private async Task ScheduleCronEventWorkflowsAsync(CancellationToken cancellationToken) { // Schedule workflow instances blocked on a cron event. - var cronEventTriggers = await _bookmarkFinder.FindBookmarksAsync(tenantId: TenantId, cancellationToken: cancellationToken); + var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(tenantId: TenantId, cancellationToken: cancellationToken).ToList(); - foreach (var result in cronEventTriggers) + _logger.LogDebug("Found {BookmarkResultCount} bookmarks for StartAt", bookmarkResults.Count); + var index = 0; + + foreach (var result in bookmarkResults) { var trigger = (CronBookmark) result.Bookmark; await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, trigger.ExecuteAt!.Value, null, cancellationToken); + + index++; + _logger.LogDebug("Scheduled {CurrentBookmarkIndex} of {BookmarkResultCount}", index, bookmarkResults.Count); } } } diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs index c8d6fd7e3..6b9e624eb 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowDefinitionScheduler.cs @@ -2,6 +2,7 @@ using System.Threading.Tasks; using Elsa.Activities.Temporal.Common.Services; using Elsa.Activities.Temporal.Quartz.Jobs; +using Elsa.Services; using Microsoft.Extensions.Logging; using NodaTime; using Quartz; @@ -13,12 +14,15 @@ namespace Elsa.Activities.Temporal.Quartz.Services { private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowDefinitionJob); private readonly QuartzSchedulerProvider _schedulerProvider; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ElsaOptions _elsaOptions; private readonly ILogger _logger; - private readonly SemaphoreSlim _semaphore = new(1); - public QuartzWorkflowDefinitionScheduler(QuartzSchedulerProvider schedulerProvider, ILogger logger) + public QuartzWorkflowDefinitionScheduler(QuartzSchedulerProvider schedulerProvider, IDistributedLockProvider distributedLockProvider, ElsaOptions elsaOptions, ILogger logger) { _schedulerProvider = schedulerProvider; + _distributedLockProvider = distributedLockProvider; + _elsaOptions = elsaOptions; _logger = logger; } @@ -69,7 +73,11 @@ namespace Elsa.Activities.Temporal.Quartz.Services private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) { var scheduler = await _schedulerProvider.GetSchedulerAsync(cancellationToken); - await _semaphore.WaitAsync(cancellationToken); + var sharedResource = $"{nameof(QuartzWorkflowInstanceScheduler)}:{trigger.Key}"; + await using var handle = await _distributedLockProvider.AcquireLockAsync(sharedResource, _elsaOptions.DistributedLockTimeout, cancellationToken); + + if (handle == null) + return; try { @@ -83,10 +91,6 @@ namespace Elsa.Activities.Temporal.Quartz.Services { _logger.LogWarning(e, "Failed to schedule trigger {TriggerKey}", trigger.Key.ToString()); } - finally - { - _semaphore.Release(); - } } private TriggerBuilder CreateTrigger(string workflowDefinitionId, string activityId) => diff --git a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs index 60d9562cc..525985aca 100644 --- a/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs +++ b/src/activities/Elsa.Activities.Temporal.Quartz/Services/QuartzWorkflowInstanceScheduler.cs @@ -2,10 +2,12 @@ using System.Threading.Tasks; using Elsa.Activities.Temporal.Common.Services; using Elsa.Activities.Temporal.Quartz.Jobs; +using Medallion.Threading; using Microsoft.Extensions.Logging; using NodaTime; using Quartz; using Quartz.Impl.Matchers; +using IDistributedLockProvider = Elsa.Services.IDistributedLockProvider; namespace Elsa.Activities.Temporal.Quartz.Services { @@ -13,12 +15,15 @@ namespace Elsa.Activities.Temporal.Quartz.Services { private static readonly string RunWorkflowJobKey = nameof(RunQuartzWorkflowInstanceJob); private readonly QuartzSchedulerProvider _schedulerProvider; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ElsaOptions _elsaOptions; private readonly ILogger _logger; - private readonly SemaphoreSlim _semaphore = new(1); - public QuartzWorkflowInstanceScheduler(QuartzSchedulerProvider schedulerProvider, ILogger logger) + public QuartzWorkflowInstanceScheduler(QuartzSchedulerProvider schedulerProvider, IDistributedLockProvider distributedLockProvider, ElsaOptions elsaOptions, ILogger logger) { _schedulerProvider = schedulerProvider; + _distributedLockProvider = distributedLockProvider; + _elsaOptions = elsaOptions; _logger = logger; } @@ -68,8 +73,12 @@ namespace Elsa.Activities.Temporal.Quartz.Services private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) { - var scheduler = await _schedulerProvider. GetSchedulerAsync(cancellationToken); - await _semaphore.WaitAsync(cancellationToken); + var scheduler = await _schedulerProvider.GetSchedulerAsync(cancellationToken); + var sharedResource = $"{nameof(QuartzWorkflowInstanceScheduler)}:{trigger.Key}"; + await using var handle = await _distributedLockProvider.AcquireLockAsync(sharedResource, _elsaOptions.DistributedLockTimeout, cancellationToken); + + if (handle == null) + return; try { @@ -84,10 +93,6 @@ namespace Elsa.Activities.Temporal.Quartz.Services { _logger.LogWarning(e, "Failed to schedule trigger {TriggerKey}", trigger.Key.ToString()); } - finally - { - _semaphore.Release(); - } } private TriggerBuilder CreateTrigger(string workflowInstanceId, string activityId) => diff --git a/src/core/Elsa.Core/HostedServices/IScopedBackgroundService.cs b/src/core/Elsa.Core/HostedServices/IScopedBackgroundService.cs new file mode 100644 index 000000000..e401f52f1 --- /dev/null +++ b/src/core/Elsa.Core/HostedServices/IScopedBackgroundService.cs @@ -0,0 +1,10 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.HostedServices +{ + public interface IScopedBackgroundService + { + Task ExecuteAsync(CancellationToken cancellationToken); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/HostedServices/ScopedBackgroundService.cs b/src/core/Elsa.Core/HostedServices/ScopedBackgroundService.cs new file mode 100644 index 000000000..45a101b6c --- /dev/null +++ b/src/core/Elsa.Core/HostedServices/ScopedBackgroundService.cs @@ -0,0 +1,24 @@ +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; + +namespace Elsa.HostedServices +{ + /// + /// Executed the specified worker within a scoped-lifetime scope. + /// + public class ScopedBackgroundService : BackgroundService where TWorker:IScopedBackgroundService + { + private readonly IServiceScopeFactory _scopeFactory; + + public ScopedBackgroundService(IServiceScopeFactory scopeFactory) => _scopeFactory = scopeFactory; + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + using var scope = _scopeFactory.CreateScope(); + var worker = (IScopedBackgroundService)ActivatorUtilities.GetServiceOrCreateInstance(scope.ServiceProvider); + await worker.ExecuteAsync(stoppingToken); + } + } +} \ No newline at end of file