Optimize workflow trigger re-indexing

This commit is contained in:
Sipke Schoorstra 2021-01-10 13:43:34 +01:00
parent b44f2b681c
commit 49d4e5bb36
13 changed files with 154 additions and 53 deletions

View file

@ -1,5 +1,4 @@
using System;
using System.Diagnostics;
using System.Diagnostics;
using System.Threading.Tasks;
using Elsa.DistributedLock;
using Elsa.Models;
@ -18,12 +17,12 @@ namespace Elsa.Activities.Timers.Quartz.Jobs
private readonly IWorkflowQueue _workflowQueue;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly ILogger _logger;
private Stopwatch _stopwatch = new();
private readonly Stopwatch _stopwatch = new();
public RunQuartzWorkflowJob(
IWorkflowRegistry workflowRegistry,
IWorkflowInstanceStore workflowInstanceStore,
IWorkflowQueue workflowQueue,
IWorkflowQueue workflowQueue,
IDistributedLockProvider distributedLockProvider,
ILogger<RunQuartzWorkflowJob> logger)
{

View file

@ -15,7 +15,7 @@ namespace Elsa.Triggers
Func<ITrigger, bool> evaluate,
CancellationToken cancellationToken = default);
Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default);
Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, string? workflowInstanceId, CancellationToken cancellationToken = default);
Task RemoveTriggerAsync(ITrigger trigger, CancellationToken cancellationToken = default);
}
}

View file

@ -1,4 +1,5 @@
using System;
using System.Diagnostics;
using System.Threading.Tasks;
using Elsa.DistributedLock;
using Elsa.Messages;
@ -17,6 +18,7 @@ namespace Elsa.Consumers
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly IEventPublisher _eventPublisher;
private readonly ILogger _logger;
private readonly Stopwatch _stopwatch = new();
public RunWorkflowInstanceConsumer(
IWorkflowRunner workflowRunner,
@ -35,10 +37,12 @@ namespace Elsa.Consumers
public async Task Handle(RunWorkflowInstance message)
{
var workflowInstanceId = message.WorkflowInstanceId;
var lockKey = workflowInstanceId;
_logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}.", workflowInstanceId);
_stopwatch.Restart();
if (!await _distributedLockProvider.AcquireLockAsync(workflowInstanceId))
if (!await _distributedLockProvider.AcquireLockAsync(lockKey))
{
// Reschedule message.
_logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Rescheduling message.", workflowInstanceId);
@ -61,7 +65,9 @@ namespace Elsa.Consumers
}
finally
{
await _distributedLockProvider.ReleaseLockAsync(workflowInstanceId);
await _distributedLockProvider.ReleaseLockAsync(lockKey);
_stopwatch.Stop();
_logger.LogDebug("Held lock on workflow instance {WorkflowInstanceId} for {ElapsedTime}.", workflowInstanceId, _stopwatch.Elapsed);
}
}

View file

@ -237,14 +237,7 @@ namespace Elsa
.Subscriptions(s => s.StoreInMemory(store))
.Transport(t => t.UseInMemoryTransport(transport, queueName))
.Routing(r => r.TypeBased().Map(context.MessageTypeMap))
.Options(options =>
{
if(ServiceBusOptions.NumberOfWorkers != null)
options.SetNumberOfWorkers(ServiceBusOptions.NumberOfWorkers.Value);
if(ServiceBusOptions.MaxParallelism != null)
options.SetMaxParallelism(ServiceBusOptions.MaxParallelism.Value);
});
.Options(options => options.Apply(ServiceBusOptions));
}
}
}

View file

@ -0,0 +1,16 @@
using Rebus.Config;
namespace Elsa
{
public static class RebusOptionsConfigurerExtensions
{
public static void Apply(this OptionsConfigurer configurer, ServiceBusOptions options)
{
if(options.NumberOfWorkers != null)
configurer.SetNumberOfWorkers(options.NumberOfWorkers.Value);
if(options.MaxParallelism != null)
configurer.SetMaxParallelism(options.MaxParallelism.Value);
}
}
}

View file

@ -0,0 +1,46 @@
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Events;
using MediatR;
using Microsoft.Extensions.Logging;
namespace Elsa.Handlers
{
public class LogWorkflowExecution : INotificationHandler<ActivityExecuting>, INotificationHandler<ActivityExecuted>, INotificationHandler<WorkflowExecutionBurstCompleted>
{
private readonly ILogger _logger;
private readonly Stopwatch _stopwatch = new();
public LogWorkflowExecution(ILogger<LogWorkflowExecution> logger)
{
_logger = logger;
}
public Task Handle(ActivityExecuting notification, CancellationToken cancellationToken)
{
_stopwatch.Start();
var activityBlueprint = notification.Activity;
var activityId = activityBlueprint.Id;
var workflowInstanceId = notification.WorkflowExecutionContext.WorkflowInstance.Id;
_logger.LogDebug("Executing activity {ActivityType} {ActivityId} for workflow {workflowInstanceId}", activityBlueprint.Type, activityId, workflowInstanceId);
return Task.CompletedTask;
}
public Task Handle(ActivityExecuted notification, CancellationToken cancellationToken)
{
_stopwatch.Stop();
var activityBlueprint = notification.Activity;
var activityId = activityBlueprint.Id;
var workflowInstanceId = notification.WorkflowExecutionContext.WorkflowInstance.Id;
_logger.LogDebug("Executed activity {ActivityType} {ActivityId} for workflow {workflowInstanceId} in {ElapsedTime}", activityBlueprint.Type, activityId, workflowInstanceId, _stopwatch.Elapsed);
return Task.CompletedTask;
}
public Task Handle(WorkflowExecutionBurstCompleted notification, CancellationToken cancellationToken)
{
_logger.LogDebug("Burst of workflow execution completed for workflow {WorkflowInstanceId}", notification.WorkflowExecutionContext.WorkflowInstance.Id);
return Task.CompletedTask;
}
}
}

View file

@ -12,7 +12,7 @@ namespace Elsa.Handlers
public class PersistWorkflow :
INotificationHandler<WorkflowExecuted>,
INotificationHandler<WorkflowSuspended>,
INotificationHandler<ActivityExecuted>,
INotificationHandler<WorkflowExecutionPassCompleted>,
INotificationHandler<WorkflowExecutionFinished>
{
private readonly IWorkflowInstanceStore _workflowInstanceStore;
@ -36,9 +36,9 @@ namespace Elsa.Handlers
await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken);
}
public async Task Handle(ActivityExecuted notification, CancellationToken cancellationToken)
public async Task Handle(WorkflowExecutionPassCompleted notification, CancellationToken cancellationToken)
{
if (notification.WorkflowExecutionContext.WorkflowBlueprint.PersistenceBehavior == WorkflowPersistenceBehavior.ActivityExecuted || notification.Activity.PersistWorkflow)
if (notification.WorkflowExecutionContext.WorkflowBlueprint.PersistenceBehavior == WorkflowPersistenceBehavior.ActivityExecuted || notification.ActivityExecutionContext.ActivityBlueprint.PersistWorkflow)
await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken);
}
@ -61,6 +61,7 @@ namespace Elsa.Handlers
{
var workflowInstance = workflowExecutionContext.WorkflowInstance;
await _workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken);
_logger.LogDebug("Committed workflow {WorkflowInstanceId} to storage", workflowInstance.Id);
}
}
}

View file

@ -23,23 +23,23 @@ namespace Elsa.Handlers
public async Task Handle(WorkflowInstanceSaved notification, CancellationToken cancellationToken)
{
var workflowInstance = notification.WorkflowInstance;
await UpdateTriggersAsync(workflowInstance.DefinitionId, workflowInstance.TenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken);
await UpdateTriggersAsync(workflowInstance.DefinitionId, workflowInstance.Id, workflowInstance.TenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken);
}
public async Task Handle(ManyWorkflowInstancesDeleted notification, CancellationToken cancellationToken)
{
var workflowInstance = notification.WorkflowInstances.First();
await UpdateTriggersAsync(workflowInstance.DefinitionId, workflowInstance.TenantId, VersionOptions.Latest, cancellationToken);
await UpdateTriggersAsync(workflowInstance.DefinitionId, null, workflowInstance.TenantId, VersionOptions.Latest, cancellationToken);
}
private async Task UpdateTriggersAsync(string workflowDefinitionId, string? tenantId, VersionOptions versionOptions, CancellationToken cancellationToken)
private async Task UpdateTriggersAsync(string workflowDefinitionId, string? workflowInstanceId, string? tenantId, VersionOptions versionOptions, CancellationToken cancellationToken)
{
var workflowBlueprint = await _workflowRegistry.GetWorkflowAsync(workflowDefinitionId, tenantId, versionOptions, cancellationToken);
if (workflowBlueprint == null)
return;
await _workflowSelector.UpdateTriggersAsync(workflowBlueprint, cancellationToken);
await _workflowSelector.UpdateTriggersAsync(workflowBlueprint, workflowInstanceId, cancellationToken);
}
}
}

View file

@ -1,8 +1,19 @@
namespace Elsa
using Rebus.Config;
namespace Elsa
{
public class ServiceBusOptions
{
public int? NumberOfWorkers { get; set; }
public int? MaxParallelism { get; set; }
public void Apply(OptionsConfigurer configurer)
{
if(NumberOfWorkers != null)
configurer.SetNumberOfWorkers(NumberOfWorkers.Value);
if(MaxParallelism != null)
configurer.SetMaxParallelism(MaxParallelism.Value);
}
}
}

View file

@ -285,15 +285,13 @@ namespace Elsa.Services
var activityBlueprint = workflowBlueprint.GetActivity(currentActivityId)!;
var activityExecutionContext = new ActivityExecutionContext(scope, workflowExecutionContext, activityBlueprint, scheduledActivity.Input, cancellationToken);
var activity = await activityExecutionContext.ActivateActivityAsync(cancellationToken);
_logger.LogDebug("Executing activity {ActivityType}: {ActivityId}", activityBlueprint.Type, currentActivityId);
var result = await activityOperation(activityExecutionContext, activity);
await _mediator.Publish(new ActivityExecuting(activityExecutionContext), cancellationToken);
var result = await activityOperation(activityExecutionContext, activity);
await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken);
await result.ExecuteAsync(activityExecutionContext, cancellationToken);
workflowExecutionContext.WorkflowInstance.Output = activityExecutionContext.Output;
workflowExecutionContext.ExecutionLog.Add(activity.Id);
workflowExecutionContext.PruneActivityData();
await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken);
activityOperation = Execute;
workflowExecutionContext.CompletePass();
await _mediator.Publish(new WorkflowExecutionPassCompleted(workflowExecutionContext, activityExecutionContext), cancellationToken);

View file

@ -1,3 +1,4 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.DistributedLock;
@ -30,15 +31,24 @@ namespace Elsa.StartupTasks
public async Task ExecuteAsync(CancellationToken cancellationToken = default)
{
if (!await _distributedLockProvider.AcquireLockAsync(GetType().Name, cancellationToken))
var lockKey = GetType().Name;
if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken))
return;
var instances = await _workflowInstanceStore.FindManyAsync(new WorkflowStatusSpecification(WorkflowStatus.Running), cancellationToken: cancellationToken);
foreach (var instance in instances)
await _workflowScheduler.RunWorkflowAsync(
instance,
cancellationToken: cancellationToken);
try
{
var instances = await _workflowInstanceStore.FindManyAsync(new WorkflowStatusSpecification(WorkflowStatus.Running), cancellationToken: cancellationToken);
foreach (var instance in instances)
await _workflowScheduler.RunWorkflowAsync(
instance,
cancellationToken: cancellationToken);
}
finally
{
await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken);
}
}
}
}

View file

@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
@ -10,6 +11,7 @@ using Elsa.Services;
using Elsa.Services.Models;
using Microsoft.Extensions.Caching.Memory;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Open.Linq.AsyncExtensions;
namespace Elsa.Triggers
@ -24,7 +26,9 @@ namespace Elsa.Triggers
private readonly IEnumerable<ITriggerProvider> _triggerProviders;
private readonly IMemoryCache _memoryCache;
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<WorkflowSelector> _logger;
private IDictionary<string, ICollection<TriggerDescriptor>>? _descriptors;
private readonly Stopwatch _stopwatch = new();
public WorkflowSelector(
IWorkflowRegistry workflowRegistry,
@ -33,7 +37,8 @@ namespace Elsa.Triggers
IWorkflowContextManager workflowContextManager,
IEnumerable<ITriggerProvider> triggerProviders,
IMemoryCache memoryCache,
IServiceProvider serviceProvider)
IServiceProvider serviceProvider,
ILogger<WorkflowSelector> logger)
{
_workflowRegistry = workflowRegistry;
_workflowFactory = workflowFactory;
@ -42,6 +47,7 @@ namespace Elsa.Triggers
_triggerProviders = triggerProviders;
_memoryCache = memoryCache;
_serviceProvider = serviceProvider;
_logger = logger;
}
public async Task<IEnumerable<WorkflowSelectorResult>> SelectWorkflowsAsync(
@ -62,11 +68,23 @@ namespace Elsa.Triggers
return SelectWorkflowsAsync(_descriptors, triggerType, evaluate).Select(result => result.WorkflowBlueprint.Activities.First(x => x.Id == result.ActivityId));
}
public async Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default)
public async Task UpdateTriggersAsync(IWorkflowBlueprint workflowBlueprint, string? workflowInstanceId, CancellationToken cancellationToken = default)
{
if (workflowInstanceId == null)
_logger.LogDebug("Updating triggers for workflows of {WorkflowDefinitionId}", workflowBlueprint.Id);
else
_logger.LogDebug("Updating triggers for workflow {WorkflowInstanceId}", workflowInstanceId);
_stopwatch.Restart();
_descriptors ??= await GetDescriptorsAsync(cancellationToken);
_descriptors[workflowBlueprint.Id] = await BuildDescriptorsForAsync(workflowBlueprint, cancellationToken).ToList();
_descriptors[workflowBlueprint.Id] = await BuildDescriptorsForAsync(workflowBlueprint, workflowInstanceId, cancellationToken).ToList();
_memoryCache.Set(CacheKey, _descriptors);
_stopwatch.Stop();
if (workflowInstanceId == null)
_logger.LogDebug("Updated triggers for workflows of {WorkflowDefinitionId} in {ElapsedTime}", workflowBlueprint.Id, _stopwatch.Elapsed);
else
_logger.LogDebug("Updated triggers for workflow {WorkflowInstanceId} in {ElapsedTime}", workflowInstanceId, _stopwatch.Elapsed);
}
public async Task RemoveTriggerAsync(ITrigger trigger, CancellationToken cancellationToken = default)
@ -113,7 +131,7 @@ namespace Elsa.Triggers
foreach (var workflowBlueprint in workflowBlueprints)
{
var blueprintDescriptors = await BuildDescriptorsForAsync(workflowBlueprint, cancellationToken);
var blueprintDescriptors = await BuildDescriptorsForAsync(workflowBlueprint, null, cancellationToken);
descriptors.AddRange(blueprintDescriptors);
}
@ -122,7 +140,7 @@ namespace Elsa.Triggers
.ToDictionary(x => x.Key, x => (ICollection<TriggerDescriptor>) x.ToList());
}
private async Task<IEnumerable<TriggerDescriptor>> BuildDescriptorsForAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken)
private async Task<IEnumerable<TriggerDescriptor>> BuildDescriptorsForAsync(IWorkflowBlueprint workflowBlueprint, string? workflowInstanceId, CancellationToken cancellationToken)
{
var descriptors = new List<TriggerDescriptor>();
@ -131,10 +149,12 @@ namespace Elsa.Triggers
descriptors.AddRange(startTriggers);
// Build triggers for workflow instances.
var specification = new WorkflowInstanceDefinitionIdSpecification(workflowBlueprint.Id)
.WithTenant(workflowBlueprint.TenantId)
.WithStatus(WorkflowStatus.Suspended);
var specification = workflowInstanceId == null
? new WorkflowInstanceDefinitionIdSpecification(workflowBlueprint.Id)
.WithTenant(workflowBlueprint.TenantId)
.WithStatus(WorkflowStatus.Suspended)
: new WorkflowInstanceIdSpecification(workflowInstanceId);
var workflowInstances = await _workflowInstanceStore
.FindManyAsync(specification, cancellationToken: cancellationToken)
.ToList();
@ -153,7 +173,7 @@ namespace Elsa.Triggers
{
var startActivities = workflowBlueprint.GetStartActivities();
var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, cancellationToken: cancellationToken);
// This is a transient workflow instance; setting EntityId to null ensures trigger providers don't try and load the workflow instance.
workflowInstance.Id = null!;
return await BuildDescriptorsAsync(workflowBlueprint, startActivities, workflowInstance, cancellationToken);
@ -170,12 +190,12 @@ namespace Elsa.Triggers
var scope = _serviceProvider.CreateScope();
var workflowExecutionContext = new WorkflowExecutionContext(scope, workflowBlueprint, workflowInstance);
var isTransientWorkflowInstance = workflowInstance.Id == null!;
workflowExecutionContext.WorkflowContext =
workflowBlueprint.ContextOptions != null &&
!isTransientWorkflowInstance &&
!string.IsNullOrWhiteSpace(workflowInstance.ContextId)
? await _workflowContextManager.LoadContext(new LoadWorkflowContext(workflowExecutionContext), cancellationToken)
workflowExecutionContext.WorkflowContext =
workflowBlueprint.ContextOptions != null &&
!isTransientWorkflowInstance &&
!string.IsNullOrWhiteSpace(workflowInstance.ContextId)
? await _workflowContextManager.LoadContext(new LoadWorkflowContext(workflowExecutionContext), cancellationToken)
: default;
foreach (var blockingActivity in blockingActivities)

View file

@ -11,10 +11,10 @@ namespace Elsa.Rebus.AzureServiceBus
{
public static ElsaOptions UseAzureServiceBus(this ElsaOptions elsaOptions, string connectionString, ITokenProvider? tokenProvider = default)
{
return elsaOptions.UseServiceBus(context => ConfigureAzureServiceBusEndpoint(context, connectionString, tokenProvider));
return elsaOptions.UseServiceBus(context => ConfigureAzureServiceBusEndpoint(elsaOptions, context, connectionString, tokenProvider));
}
private static void ConfigureAzureServiceBusEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString, ITokenProvider? tokenProvider)
private static void ConfigureAzureServiceBusEndpoint(ElsaOptions elsaOptions, ServiceBusEndpointConfigurationContext context, string connectionString, ITokenProvider? tokenProvider)
{
var queueName = context.QueueName;
var loggerFactory = context.ServiceProvider.GetRequiredService<ILoggerFactory>();
@ -22,7 +22,8 @@ namespace Elsa.Rebus.AzureServiceBus
context.Configurer
.Logging(l => l.MicrosoftExtensionsLogging(loggerFactory))
.Transport(t => t.UseAzureServiceBus(connectionString, queueName, tokenProvider))
.Routing(r => r.TypeBased().Map(context.MessageTypeMap));
.Routing(r => r.TypeBased().Map(context.MessageTypeMap))
.Options(o => o.Apply(elsaOptions.ServiceBusOptions));
}
}
}