Merge pull request #6498 from truthz03/feature/AddJobRunnerWorkerCountSetting
Allow to configure JobRunner worker count
This commit is contained in:
commit
40a6d3ea8d
|
|
@ -14,7 +14,6 @@ using Elsa.Mediator.Options;
|
|||
using Elsa.Mediator.Services;
|
||||
using JetBrains.Annotations;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
// ReSharper disable once CheckNamespace
|
||||
namespace Microsoft.Extensions.DependencyInjection;
|
||||
|
|
@ -54,22 +53,10 @@ public static class DependencyInjectionExtensions
|
|||
.AddSingleton<INotificationsChannel, NotificationsChannel>()
|
||||
.AddSingleton<ICommandsChannel, CommandsChannel>()
|
||||
.AddSingleton<IJobsChannel, JobsChannel>()
|
||||
.AddSingleton<IJobQueue, JobQueue>()
|
||||
.AddSingleton<IJobQueue, JobQueue>()
|
||||
.AddHostedService<JobRunnerHostedService>()
|
||||
.AddHostedService(sp =>
|
||||
{
|
||||
using var scope = sp.CreateScope();
|
||||
|
||||
var options = scope.ServiceProvider.GetRequiredService<IOptions<MediatorOptions>>().Value;
|
||||
return ActivatorUtilities.CreateInstance<BackgroundCommandSenderHostedService>(scope.ServiceProvider, options.CommandWorkerCount);
|
||||
})
|
||||
.AddHostedService(sp =>
|
||||
{
|
||||
using var scope = sp.CreateScope();
|
||||
|
||||
var options = scope.ServiceProvider.GetRequiredService<IOptions<MediatorOptions>>().Value;
|
||||
return ActivatorUtilities.CreateInstance<BackgroundEventPublisherHostedService>(scope.ServiceProvider, options.NotificationWorkerCount);
|
||||
});
|
||||
.AddHostedService<BackgroundCommandSenderHostedService>()
|
||||
.AddHostedService<BackgroundEventPublisherHostedService>();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
|
|
|
|||
|
|
@ -1,8 +1,10 @@
|
|||
using System.Threading.Channels;
|
||||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Mediator.Options;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace Elsa.Mediator.HostedServices;
|
||||
|
||||
|
|
@ -18,13 +20,13 @@ public class BackgroundCommandSenderHostedService : BackgroundService
|
|||
private readonly ILogger _logger;
|
||||
|
||||
/// <inheritdoc />
|
||||
public BackgroundCommandSenderHostedService(int workerCount, ICommandsChannel commandsChannel, IServiceScopeFactory scopeFactory, ILogger<BackgroundCommandSenderHostedService> logger)
|
||||
public BackgroundCommandSenderHostedService(IOptions<MediatorOptions> options, ICommandsChannel commandsChannel, IServiceScopeFactory scopeFactory, ILogger<BackgroundCommandSenderHostedService> logger)
|
||||
{
|
||||
_workerCount = workerCount;
|
||||
_workerCount = options.Value.CommandWorkerCount;
|
||||
_commandsChannel = commandsChannel;
|
||||
_scopeFactory = scopeFactory;
|
||||
_logger = logger;
|
||||
_outputs = new List<Channel<ICommand>>(workerCount);
|
||||
_outputs = new List<Channel<ICommand>>(_workerCount);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
|
|
@ -1,8 +1,10 @@
|
|||
using System.Threading.Channels;
|
||||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Mediator.Options;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace Elsa.Mediator.HostedServices;
|
||||
|
||||
|
|
@ -18,13 +20,13 @@ public class BackgroundEventPublisherHostedService : BackgroundService
|
|||
private readonly ILogger _logger;
|
||||
|
||||
/// <inheritdoc />
|
||||
public BackgroundEventPublisherHostedService(int workerCount, INotificationsChannel notificationsChannel, IServiceScopeFactory scopeFactory, ILogger<BackgroundEventPublisherHostedService> logger)
|
||||
public BackgroundEventPublisherHostedService(IOptions<MediatorOptions> options, INotificationsChannel notificationsChannel, IServiceScopeFactory scopeFactory, ILogger<BackgroundEventPublisherHostedService> logger)
|
||||
{
|
||||
_workerCount = workerCount;
|
||||
_workerCount = options.Value.NotificationWorkerCount;
|
||||
_notificationsChannel = notificationsChannel;
|
||||
_scopeFactory = scopeFactory;
|
||||
_logger = logger;
|
||||
_outputs = new List<Channel<INotification>>(workerCount);
|
||||
_outputs = new List<Channel<INotification>>(_workerCount);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Mediator.Options;
|
||||
using Microsoft.Extensions.Hosting;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace Elsa.Mediator.HostedServices;
|
||||
|
||||
|
|
@ -9,13 +11,14 @@ namespace Elsa.Mediator.HostedServices;
|
|||
/// </summary>
|
||||
public class JobRunnerHostedService : BackgroundService
|
||||
{
|
||||
private const int WorkerCount = 4;
|
||||
private readonly int _workerCount;
|
||||
private readonly IJobsChannel _jobsChannel;
|
||||
private readonly ILogger<JobRunnerHostedService> _logger;
|
||||
|
||||
/// <inheritdoc />
|
||||
public JobRunnerHostedService(IJobsChannel jobsChannel, ILogger<JobRunnerHostedService> logger)
|
||||
{
|
||||
public JobRunnerHostedService(IOptions<MediatorOptions> options, IJobsChannel jobsChannel, ILogger<JobRunnerHostedService> logger)
|
||||
{
|
||||
_workerCount = options.Value.JobWorkerCount;
|
||||
_jobsChannel = jobsChannel;
|
||||
_logger = logger;
|
||||
}
|
||||
|
|
@ -23,9 +26,9 @@ public class JobRunnerHostedService : BackgroundService
|
|||
/// <inheritdoc />
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||
{
|
||||
var workers = new Task[WorkerCount];
|
||||
|
||||
for (var i = 0; i < WorkerCount; i++)
|
||||
var workers = new Task[_workerCount];
|
||||
|
||||
for (var i = 0; i < _workerCount; i++)
|
||||
workers[i] = ProcessJobsAsync(stoppingToken);
|
||||
|
||||
await Task.WhenAll(workers);
|
||||
|
|
|
|||
|
|
@ -10,12 +10,17 @@ public class MediatorOptions
|
|||
/// <summary>
|
||||
/// The max number of workers to process commands.
|
||||
/// </summary>
|
||||
public int CommandWorkerCount { get; set; } = 4;
|
||||
|
||||
public int CommandWorkerCount { get; set; } = 4;
|
||||
|
||||
/// <summary>
|
||||
/// The max number of workers to process notifications.
|
||||
/// </summary>
|
||||
public int NotificationWorkerCount { get; set; } = 4;
|
||||
public int NotificationWorkerCount { get; set; } = 4;
|
||||
|
||||
/// <summary>
|
||||
/// The max number of workers to process jobs.
|
||||
/// </summary>
|
||||
public int JobWorkerCount { get; set; } = 4;
|
||||
|
||||
/// <summary>
|
||||
/// The default publishing strategy to use.
|
||||
|
|
|
|||
Loading…
Reference in a new issue