From c5d712ce51f692d1153863fce50574588b5bd2b6 Mon Sep 17 00:00:00 2001 From: maliming Date: Mon, 31 Mar 2025 13:39:19 +0800 Subject: [PATCH] Refactor RabbitMQ channel and connection handling for improved efficiency and reliability --- .../RabbitMq/RabbitMqDistributedEventBus.cs | 11 ++++------ .../Volo/Abp/RabbitMQ/ChannelPool.cs | 19 +++++++++++------ .../Volo/Abp/RabbitMQ/ConnectionPool.cs | 21 ++++++++++++------- 3 files changed, 31 insertions(+), 20 deletions(-) diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs index b4fc4f5768..3a647388b5 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs @@ -225,7 +225,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis OutboxConfig outboxConfig) { using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) - .CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true))) + .CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true, new ThrottlingRateLimiter(256)))) { var outgoingEventArray = outgoingEvents.ToArray(); @@ -298,8 +298,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis Guid? eventId = null, string? correlationId = null) { - using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) - .CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true))) + using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)).CreateChannelAsync()) { await PublishAsync(channel, eventName, body, headersArguments, eventId, correlationId); } @@ -335,7 +334,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis await channel.BasicPublishAsync( exchange: AbpRabbitMqEventBusOptions.ExchangeName, routingKey: eventName, - mandatory: true, + mandatory: false, basicProperties: properties, body: body ); @@ -350,9 +349,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis try { - using (var temporaryChannel = await(await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) - .CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true, - publisherConfirmationTrackingEnabled: true))) + using (var temporaryChannel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)).CreateChannelAsync()) { await temporaryChannel.ExchangeDeclarePassiveAsync(AbpRabbitMqEventBusOptions.ExchangeName); } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs index 02711189b1..bbb594a5f7 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs @@ -49,14 +49,14 @@ public class ChannelPool : IChannelPool, ISingletonDependency { using (await Semaphore.LockAsync()) { - if (!Channels.TryGetValue(channelName, out var channel)) + if (Channels.TryGetValue(channelName, out var existingChannelPoolItem2)) { - poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); - Channels.TryAdd(channelName, poolItem); + poolItem = existingChannelPoolItem2; } else { - poolItem = channel; + poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); + Channels.TryAdd(channelName, poolItem); } } } @@ -70,8 +70,15 @@ public class ChannelPool : IChannelPool, ISingletonDependency using (await Semaphore.LockAsync()) { - poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); - Channels.TryAdd(channelName, poolItem); + if (Channels.TryGetValue(channelName, out var existingChannelPoolItem3)) + { + poolItem = existingChannelPoolItem3; + } + else + { + poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); + Channels.TryAdd(channelName, poolItem); + } } poolItem.Acquire(); diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs index a4425810b8..938b0e4f8c 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs @@ -39,18 +39,25 @@ public class ConnectionPool : IConnectionPool, ISingletonDependency { using (await Semaphore.LockAsync()) { - var connectionFactory = Options.Connections.GetOrDefault(connectionName); try { - connection = await GetConnectionAsync(connectionName, connectionFactory); - Connections.TryAdd(connectionName, connection); - - if (!connection.IsOpen) + var connectionFactory = Options.Connections.GetOrDefault(connectionName); + if (Connections.TryGetValue(connectionName, out var existingConnection2)) + { + connection = existingConnection2; + } + else { - connection.Dispose(); - Connections.TryRemove(connectionName, out _); connection = await GetConnectionAsync(connectionName, connectionFactory); Connections.TryAdd(connectionName, connection); + + if (!connection.IsOpen) + { + connection.Dispose(); + Connections.TryRemove(connectionName, out _); + connection = await GetConnectionAsync(connectionName, connectionFactory); + Connections.TryAdd(connectionName, connection); + } } } catch (Exception)