Add ParallelForEach activity

This commit is contained in:
Sipke Schoorstra 2020-10-21 19:54:42 +02:00
parent 400b4f827a
commit ed7265bb61
5 changed files with 129 additions and 0 deletions

View file

@ -0,0 +1,31 @@
using System.Collections.Generic;
using System.Collections.ObjectModel;
using System.Linq;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
[ActivityDefinition(
Category = "Control Flow",
Description = "Iterate over a collection in parallel.",
Icon = "far fa-circle",
Outcomes = new[] { OutcomeNames.Iterate, OutcomeNames.Done }
)]
public class ParallelForEach : Activity
{
[ActivityProperty(Hint = "A collection of items to iterate over.")]
public ICollection<object> Items { get; set; } = new Collection<object>();
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
{
var items = Items.Reverse().ToList();
var results = new List<IActivityExecutionResult> { Done() };
results.AddRange(items.Select(x => Outcome(OutcomeNames.Iterate, x)));
return Combine(results);
}
}
}

View file

@ -0,0 +1,37 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Elsa.Builders;
using Elsa.Services;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public static class ParallelForEachBuilderExtensions
{
public static IActivityBuilder ParallelForEach(
this IBuilder builder,
Action<ISetupActivity<ParallelForEach>> setup,
Action<IOutcomeBuilder> iterate) =>
builder.Then(setup, branch => iterate(branch.When(OutcomeNames.Iterate)));
public static IActivityBuilder ParallelForEach(
this IBuilder builder,
Func<ActivityExecutionContext, ICollection<object>> items,
Action<IOutcomeBuilder> iterate) =>
builder.ParallelForEach(activity => activity.Set(x => x.Items, items), iterate);
public static IActivityBuilder ParallelForEach(
this IBuilder builder,
Func<ICollection<object>> items,
Action<IOutcomeBuilder> iterate) =>
builder.ParallelForEach(activity => activity.Set(x => x.Items, items), iterate);
public static IActivityBuilder ParallelForEach(
this IBuilder builder,
ICollection<object> items,
Action<IOutcomeBuilder> iterate) =>
builder.ParallelForEach(activity => activity.Set(x => x.Items, items), iterate);
}
}

View file

@ -130,6 +130,7 @@ namespace Microsoft.Extensions.DependencyInjection
.AddActivity<Complete>()
.AddActivity<For>()
.AddActivity<ForEach>()
.AddActivity<ParallelForEach>()
.AddActivity<Fork>()
.AddActivity<IfElse>()
.AddActivity<Join>()

View file

@ -0,0 +1,29 @@
using System.Linq;
using System.Threading.Tasks;
using Elsa.Core.IntegrationTests.Workflows;
using Elsa.Models;
using Elsa.Testing.Shared.Helpers;
using Xunit;
using Xunit.Abstractions;
namespace Elsa.Core.IntegrationTests
{
public class ParallelForEachWorkflowTests : WorkflowsUnitTestBase
{
public ParallelForEachWorkflowTests(ITestOutputHelper testOutputHelper) : base(testOutputHelper)
{
}
[Fact(DisplayName = "Runs Parallel For Each workflow.")]
public async Task Test01()
{
var items = Enumerable.Range(1, 10).Select(x => $"Item {x}").ToList();
var workflow = new ParallelForEachWorkflow(items);
var workflowInstance = await WorkflowHost.RunWorkflowAsync(workflow);
var iterationLogs = workflowInstance.ExecutionLog.Where(x => x.ActivityId == "WriteLine").ToList();
Assert.Equal(WorkflowStatus.Suspended, workflowInstance.Status);
Assert.Equal(items.Count, iterationLogs.Count);
}
}
}

View file

@ -0,0 +1,31 @@
using System.Collections.Generic;
using System.Linq;
using Elsa.Activities.Console;
using Elsa.Activities.ControlFlow;
using Elsa.Activities.Signaling;
using Elsa.Builders;
namespace Elsa.Core.IntegrationTests.Workflows
{
public class ParallelForEachWorkflow : IWorkflow
{
private readonly List<object> _items;
public ParallelForEachWorkflow(IEnumerable<string> items)
{
_items = items.Cast<object>().ToList();
}
public void Build(IWorkflowBuilder workflow)
{
workflow
.ParallelForEach(
_items,
iterate => iterate
.Then<WriteLine>(activity => activity.Set(x => x.Text, context => $"{context.Input}")).WithId("WriteLine")
.Then<Signaled>() /* Block workflow.*/
.WriteLine("Resumed"))
.WriteLine("All iterations executing in parallel");
}
}
}