Optimize synchronized workflow instance execution in parallel

This commit is contained in:
Sipke Schoorstra 2021-01-02 12:57:27 +01:00
parent 717ced918f
commit 34983144c3
16 changed files with 109 additions and 62 deletions

View file

@ -33,7 +33,7 @@ namespace Elsa.Activities.AzureServiceBus.Services
private async Task OnMessageReceived(Message message, CancellationToken cancellationToken)
{
await _backgroundWorker.ScheduleTask(async () => await TriggerWorkflowsAsync(message, cancellationToken), cancellationToken);
await _backgroundWorker.ScheduleTask(GetType().FullName, async () => await TriggerWorkflowsAsync(message, cancellationToken), cancellationToken);
await _messageReceiver.CompleteAsync(message.SystemProperties.LockToken);
}

View file

@ -3,6 +3,7 @@ using Elsa.Activities.Timers.Options;
using Elsa.Activities.Timers.Quartz.Jobs;
using Elsa.Activities.Timers.Quartz.Services;
using Elsa.Activities.Timers.Services;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;
using Quartz;
@ -22,7 +23,6 @@ namespace Elsa
.AddQuartzHostedService(ConfigureQuartzHostedService)
.AddSingleton<IWorkflowScheduler, QuartzWorkflowScheduler>()
.AddSingleton<ICrontabParser, QuartzCrontabParser>()
.AddSingleton<WorkflowRunnerQueue>()
.AddTransient<RunQuartzWorkflowJob>();
}

View file

@ -13,14 +13,14 @@ namespace Elsa.Activities.Timers.Quartz.Jobs
private readonly IWorkflowRunner _workflowRunner;
private readonly IWorkflowRegistry _workflowRegistry;
private readonly IWorkflowInstanceStore _workflowInstanceStore;
private readonly WorkflowRunnerQueue _workflowRunnerQueue;
private readonly IWorkflowQueue _workflowQueue;
public RunQuartzWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, WorkflowRunnerQueue workflowRunnerQueue)
public RunQuartzWorkflowJob(IWorkflowRunner workflowRunner, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, IWorkflowQueue workflowQueue)
{
_workflowRunner = workflowRunner;
_workflowRegistry = workflowRegistry;
_workflowInstanceStore = workflowInstanceStore;
_workflowRunnerQueue = workflowRunnerQueue;
_workflowQueue = workflowQueue;
}
public async Task Execute(IJobExecutionContext context)
@ -47,7 +47,7 @@ namespace Elsa.Activities.Timers.Quartz.Jobs
}
else
{
await _workflowRunnerQueue.Enqueue(workflowInstanceId, activityId, cancellationToken);
await _workflowQueue.Enqueue(workflowInstanceId, activityId, cancellationToken);
}
}
}

View file

@ -1,13 +1,13 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using NodaTime;
namespace Elsa.Services
{
public interface IBackgroundWorker
{
Task ScheduleTask(Func<ValueTask> task, CancellationToken cancellationToken = default);
Task ScheduleTask(Action task, CancellationToken cancellationToken = default);
Task RunAsync(CancellationToken cancellationToken = default);
Task ScheduleTask(string channelName, Func<ValueTask> task, CancellationToken cancellationToken = default, Duration? channelTtl = default);
Task ScheduleTask(string channelName, Action task, CancellationToken cancellationToken = default, Duration? channelTtl = default);
}
}

View file

@ -0,0 +1,13 @@
using System.Threading;
using System.Threading.Tasks;
namespace Elsa.Services
{
public interface IWorkflowQueue
{
/// <summary>
/// Enqueues the specified workflow instance and activity for execution in the background.
/// </summary>
Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default);
}
}

View file

@ -33,7 +33,7 @@ namespace Elsa.Activities.Signaling
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context)
{
var triggeredSignal = context.GetInput<Signal>();
var triggeredSignal = context.GetInput<Signal>()!;
return Done(triggeredSignal.Input);
}
}

View file

@ -36,7 +36,6 @@ namespace Microsoft.Extensions.DependencyInjection
this IServiceCollection services,
Action<ElsaOptions>? configure = default)
{
services.AddHostedService<StartBackgroundWorker>();
services.AddStartupRunner();
var options = new ElsaOptions(services);
@ -97,6 +96,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton<IWorkflowBlueprintMaterializer, WorkflowBlueprintMaterializer>()
.AddSingleton<IWorkflowBlueprintReflector, WorkflowBlueprintReflector>()
.AddSingleton<IBackgroundWorker, BackgroundWorker>()
.AddSingleton<IWorkflowQueue, WorkflowQueue>()
.AddScoped<IWorkflowSelector, WorkflowSelector>()
.AddScoped<IWorkflowPublisher, WorkflowPublisher>()
.AddScoped<IWorkflowContextManager, WorkflowContextManager>()

View file

@ -1,14 +0,0 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Services;
using Microsoft.Extensions.Hosting;
namespace Elsa.HostedServices
{
public class StartBackgroundWorker : BackgroundService
{
private readonly IBackgroundWorker _backgroundWorker;
public StartBackgroundWorker(IBackgroundWorker backgroundWorker) => _backgroundWorker = backgroundWorker;
protected override async Task ExecuteAsync(CancellationToken stoppingToken) => await _backgroundWorker.RunAsync(stoppingToken);
}
}

View file

@ -1,29 +1,88 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Channels;
using System.Threading.Tasks;
using NodaTime;
namespace Elsa.Services
{
public class BackgroundWorker : IBackgroundWorker
{
private readonly Channel<Func<ValueTask>> _channel;
public BackgroundWorker() => _channel = Channel.CreateBounded<Func<ValueTask>>(10);
private readonly IDictionary<string, Channel<Func<ValueTask>>> _channels = new Dictionary<string, Channel<Func<ValueTask>>>();
public async Task ScheduleTask(Func<ValueTask> task, CancellationToken cancellationToken) => await _channel.Writer.WriteAsync(task, cancellationToken);
public async Task ScheduleTask(string channelName, Func<ValueTask> task, CancellationToken cancellationToken, Duration? channelTtl)
{
var channel = GetOrCreateChannel(channelName, channelTtl, cancellationToken);
await channel.Writer.WriteAsync(task, cancellationToken);
}
public async Task ScheduleTask(Action task, CancellationToken cancellationToken) =>
await _channel.Writer.WriteAsync(() =>
public async Task ScheduleTask(string channelName, Action task, CancellationToken cancellationToken, Duration? channelTtl)
{
await ScheduleTask(channelName, () =>
{
task();
return new ValueTask();
}, cancellationToken);
}, cancellationToken, channelTtl);
}
public async Task RunAsync(CancellationToken cancellationToken)
private Channel<Func<ValueTask>> GetOrCreateChannel(string channelName, Duration? channelTtl, CancellationToken cancellationToken)
{
while (await _channel.Reader.WaitToReadAsync(cancellationToken))
while (_channel.Reader.TryRead(out var task))
await task();
var channel = _channels.ContainsKey(channelName) ? _channels[channelName] : default;
if (channel != null)
return channel;
channel = Channel.CreateBounded<Func<ValueTask>>(10);
_channels[channelName] = channel;
var reader = new BackgroundChannelReader(channel, () => CompleteChannel(channelName));
reader.Start(channelTtl, cancellationToken);
return channel;
}
private void CompleteChannel(string name)
{
if (_channels.ContainsKey(name))
_channels.Remove(name);
}
private class BackgroundChannelReader
{
private readonly Channel<Func<ValueTask>> _channel;
private readonly Action _onComplete;
public BackgroundChannelReader(Channel<Func<ValueTask>> channel, Action onComplete)
{
_channel = channel;
_onComplete = onComplete;
}
public void Start(Duration? channelTtl, CancellationToken cancellationToken) => Task.Factory.StartNew(() => ProcessAsync(channelTtl, cancellationToken), cancellationToken);
private async Task ProcessAsync(Duration? channelTtl, CancellationToken cancellationToken)
{
var timeOutTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
if (channelTtl != null)
timeOutTokenSource.CancelAfter(channelTtl.Value.ToTimeSpan());
try
{
while (await _channel.Reader.WaitToReadAsync(timeOutTokenSource.Token))
{
var task = await _channel.Reader.ReadAsync(CancellationToken.None);
await task();
if (channelTtl != null)
timeOutTokenSource.CancelAfter(channelTtl.Value.ToTimeSpan()); // Reset timeout.
}
}
catch(TaskCanceledException)
{
_onComplete();
}
}
}
}
}

View file

@ -3,32 +3,29 @@ using System.Threading;
using System.Threading.Tasks;
using Elsa.Models;
using Elsa.Persistence;
using Elsa.Services;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using NodaTime;
namespace Elsa.Activities.Timers.Quartz.Services
namespace Elsa.Services
{
// TODO: Consider turning this into a global service to allow background, sequential execution of a given workflow instance (but allow for parallel execution of workflow instances with different definitions). Executing the same workflow instances sequentially prevents loss of data during update concurrency conflicts.
public class WorkflowRunnerQueue
public class WorkflowQueue : IWorkflowQueue
{
private readonly IBackgroundWorker _backgroundWorker;
private readonly IServiceProvider _serviceProvider;
private readonly ILogger _logger;
public WorkflowRunnerQueue(IBackgroundWorker backgroundWorker, IServiceProvider serviceProvider, ILogger<WorkflowRunnerQueue> logger)
public WorkflowQueue(IBackgroundWorker backgroundWorker, IServiceProvider serviceProvider, ILogger<WorkflowQueue> logger)
{
_backgroundWorker = backgroundWorker;
_serviceProvider = serviceProvider;
_logger = logger;
}
public async Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default)
{
await _backgroundWorker.ScheduleTask(async () => await RunWorkflowAsync(workflowInstanceId, activityId, cancellationToken), cancellationToken);
}
public async Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default) =>
await _backgroundWorker.ScheduleTask(workflowInstanceId, async () => await RunWorkflowAsync(workflowInstanceId, activityId, cancellationToken), cancellationToken, Duration.FromMinutes(1));
private async Task RunWorkflowAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default)
private async ValueTask RunWorkflowAsync(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default)
{
using var scope = _serviceProvider.CreateScope();
var store = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();

View file

@ -19,7 +19,7 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows
// Demonstrating that we can create activities and connect to them later on by using the activity builder reference.
var join = workflow.Add<Join>(x => x.WithMode(Join.JoinMode.WaitAny));
join.Finish();
workflow
// The workflow context type of this workflow.
.WithContextType<Document>()
@ -28,7 +28,7 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows
.ReceiveHttpPostRequest<Document>("/documents")
// Store the document as the workflow context. It will be saved automatically when the workflow gets suspended.
.Then(context => context.SetWorkflowContext(context.GetInput<HttpRequestModel>().GetBody<Document>())).LoadWorkflowContext()
.Then(context => context.SetWorkflowContext(context.GetInput<HttpRequestModel>()!.GetBody<Document>())).LoadWorkflowContext()
// Write an HTTP response.
.WriteHttpResponse(
@ -62,8 +62,8 @@ namespace Elsa.Samples.ContextualWorkflowHttp.Workflows
private static void StoreComment(ActivityExecutionContext context)
{
var document = (Document)context.WorkflowExecutionContext.WorkflowContext!;
var comment = (Comment)((HttpRequestModel)context.Input!).Body!;
var document = (Document) context.WorkflowExecutionContext.WorkflowContext!;
var comment = (Comment) ((HttpRequestModel) context.Input!).Body!;
document.Comments.Add(comment);
}

View file

@ -17,7 +17,7 @@ namespace ElsaDashboard.Samples.Monolith.Workflows
.Then(() => Console.WriteLine("What is your age?")).WithDisplayName("Write").WithDescription("What is your age?")
.ReadLine()
.Timer(Duration.FromMinutes(5))
.SetVariable("Age", context => int.Parse(context.GetInput<string>()))
.SetVariable("Age", context => int.Parse(context.GetInput<string>()!))
.IfElse(
context => context.GetVariable<int>("Age") < 18,
whenTrue => whenTrue.WriteLine("You are not allowed to drink beer."),

View file

@ -19,7 +19,7 @@ namespace ElsaDashboard.Samples.Monolith.Workflows
.Correlate(() => Guid.NewGuid().ToString("N"))
.WriteHttpResponse(x => x.WithStatusCode(HttpStatusCode.OK).WithContent(context => $"Tell me your name please. Use correlation ID {context.WorkflowExecutionContext.CorrelationId}").WithContentType("text/plain"))
.HttpRequestReceived(x => x.WithPath("/signup").WithMethod(HttpMethod.Post.ToString()).WithReadContent())
.SetVariable("Name", context => (string?) context.GetInput<HttpRequestModel>().Body)
.SetVariable("Name", context => (string?) context.GetInput<HttpRequestModel>()!.Body)
.SetName(context => context.GetVariable<string>("Name"))
.WriteHttpResponse(x => x

View file

@ -23,7 +23,7 @@ namespace Elsa.Samples.CustomAttributesChildWorker.Workflows
.Timer(Duration.FromSeconds(5))
.SetVariable("CustomerId", SelectRandomCustomerId)
.WriteLine(context => $"Creating a new order for customer {context.GetVariable<string>("CustomerId")}.")
.Then<SendRebusMessage>(message => message.Set(x => x.Message, context => new OrderReceived { CustomerId = context.GetVariable<string>("CustomerId") }));
.Then<SendRebusMessage>(message => message.Set(x => x.Message, context => new OrderReceived { CustomerId = context.GetVariable<string>("CustomerId")! }));
}
private string SelectRandomCustomerId()

View file

@ -18,8 +18,8 @@ namespace Elsa.Samples.CustomAttributesChildWorker.Workflows
workflow
.StartWith<RebusMessageReceived>(activity => activity.Set(x => x.MessageType, typeof(OrderReceived)))
.SetVariable(context => context.GetInput<OrderReceived>())
.WriteLine(context => $"Received a new order for {context.GetVariable<OrderReceived>().CustomerId}.")
.RunWorkflow(activity => activity.WithCustomAttributes(context => new Variables().Set("Customer", context.GetVariable<OrderReceived>().CustomerId)))
.WriteLine(context => $"Received a new order for {context.GetVariable<OrderReceived>()!.CustomerId}.")
.RunWorkflow(activity => activity.WithCustomAttributes(context => new Variables().Set("Customer", context.GetVariable<OrderReceived>()!.CustomerId)))
.WriteLine("Returned back from child workflow.");
}
}

View file

@ -1,19 +1,11 @@
using Elsa.Activities.Console;
using Elsa.Builders;
using Elsa.Samples.Timers.Activities;
using NodaTime;
namespace Elsa.Samples.Timers.Workflows
{
public class RecurringTaskWorkflow : IWorkflow
{
private readonly IClock _clock;
public RecurringTaskWorkflow(IClock clock)
{
_clock = clock;
}
public void Build(IWorkflowBuilder workflow)
{
workflow