diff --git a/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj b/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj index 6cecd5b0d8..c4b590412d 100644 --- a/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj +++ b/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj @@ -18,4 +18,8 @@ + + + + diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj b/framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj index 19b7f7032e..6cbaf62aed 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj @@ -16,6 +16,7 @@ + diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs index 43c51c0b73..176f281caa 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs @@ -1,5 +1,6 @@ using System; using System.Threading.Tasks; +using Medallion.Threading; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; @@ -15,17 +16,22 @@ namespace Volo.Abp.EventBus.Boxes protected IServiceProvider ServiceProvider { get; } protected AbpTimer Timer { get; } protected IDistributedEventBus DistributedEventBus { get; } + protected IDistributedLockProvider DistributedLockProvider { get; } protected IEventOutbox Outbox { get; private set; } + protected OutboxConfig OutboxConfig { get; private set; } + protected string DistributedLockName => "Outbox_" + OutboxConfig.Name; public ILogger Logger { get; set; } public OutboxSender( IServiceProvider serviceProvider, AbpTimer timer, - IDistributedEventBus distributedEventBus) + IDistributedEventBus distributedEventBus, + IDistributedLockProvider distributedLockProvider) { ServiceProvider = serviceProvider; Timer = timer; DistributedEventBus = distributedEventBus; + DistributedLockProvider = distributedLockProvider; Timer.Period = 2000; //TODO: Config? Timer.Elapsed += TimerOnElapsed; Logger = NullLogger.Instance; @@ -33,6 +39,7 @@ namespace Volo.Abp.EventBus.Boxes public virtual Task StartAsync(OutboxConfig outboxConfig) { + OutboxConfig = outboxConfig; Outbox = (IEventOutbox)ServiceProvider.GetRequiredService(outboxConfig.ImplementationType); Timer.Start(); return Task.CompletedTask; @@ -51,25 +58,39 @@ namespace Volo.Abp.EventBus.Boxes protected virtual async Task RunAsync() { - while (true) + await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName)) { - var waitingEvents = await Outbox.GetWaitingEventsAsync(1000); //TODO: Config? - if (waitingEvents.Count <= 0) + if (handle != null) { - break; - } + Logger.LogDebug("Obtained the distributed lock: " + DistributedLockName); + + while (true) + { + var waitingEvents = await Outbox.GetWaitingEventsAsync(1000); //TODO: Config? + if (waitingEvents.Count <= 0) + { + break; + } - Logger.LogInformation($"Found {waitingEvents.Count} events in the outbox."); + Logger.LogInformation($"Found {waitingEvents.Count} events in the outbox."); - foreach (var waitingEvent in waitingEvents) - { - await DistributedEventBus - .AsRawEventPublisher() - .PublishRawAsync(waitingEvent.Id, waitingEvent.EventName, waitingEvent.EventData); + foreach (var waitingEvent in waitingEvents) + { + await DistributedEventBus + .AsRawEventPublisher() + .PublishRawAsync(waitingEvent.Id, waitingEvent.EventName, waitingEvent.EventData); - await Outbox.DeleteAsync(waitingEvent.Id); + await Outbox.DeleteAsync(waitingEvent.Id); - Logger.LogInformation($"Sent the event to the message broker with id = {waitingEvent.Id:N}"); + Logger.LogInformation($"Sent the event to the message broker with id = {waitingEvent.Id:N}"); + } + } + + await Task.Delay(30000); + } + else + { + Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName); } } } diff --git a/test/DistEvents/DistDemoApp/DistDemoApp.csproj b/test/DistEvents/DistDemoApp/DistDemoApp.csproj index adfe26ebea..66623bfde4 100644 --- a/test/DistEvents/DistDemoApp/DistDemoApp.csproj +++ b/test/DistEvents/DistDemoApp/DistDemoApp.csproj @@ -6,6 +6,7 @@ + diff --git a/test/DistEvents/DistDemoApp/DistDemoAppModule.cs b/test/DistEvents/DistDemoApp/DistDemoAppModule.cs index 42e82f5e46..c4eab2937f 100644 --- a/test/DistEvents/DistDemoApp/DistDemoAppModule.cs +++ b/test/DistEvents/DistDemoApp/DistDemoAppModule.cs @@ -1,4 +1,7 @@ +using Medallion.Threading; +using Medallion.Threading.Redis; using Microsoft.Extensions.DependencyInjection; +using StackExchange.Redis; using Volo.Abp.Autofac; using Volo.Abp.Domain.Entities.Events.Distributed; using Volo.Abp.EntityFrameworkCore; @@ -21,6 +24,8 @@ namespace DistDemoApp { public override void ConfigureServices(ServiceConfigurationContext context) { + var configuration = context.Services.GetConfiguration(); + context.Services.AddHostedService(); context.Services.AddAbpDbContext(options => @@ -46,6 +51,12 @@ namespace DistDemoApp config.UseDbContext(); }); }); + + context.Services.AddSingleton(sp => + { + var connection = ConnectionMultiplexer.Connect(configuration["Redis:Configuration"]); + return new RedisDistributedSynchronizationProvider(connection.GetDatabase()); + }); } } } \ No newline at end of file diff --git a/test/DistEvents/DistDemoApp/appsettings.json b/test/DistEvents/DistDemoApp/appsettings.json index 8f4773ae6a..1ed30c80a9 100644 --- a/test/DistEvents/DistDemoApp/appsettings.json +++ b/test/DistEvents/DistDemoApp/appsettings.json @@ -12,5 +12,8 @@ "ClientName": "DistDemoApp", "ExchangeName": "DistDemo" } + }, + "Redis": { + "Configuration": "127.0.0.1" } } \ No newline at end of file