Browse Source

Refactor RabbitMQ channel and connection handling for improved efficiency and reliability

pull/22510/head
maliming 2 years ago
parent
commit
c5d712ce51
No known key found for this signature in database GPG Key ID: A646B9CB645ECEA4
  1. 11
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  2. 19
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs
  3. 21
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs

11
framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs

@ -225,7 +225,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis
OutboxConfig outboxConfig) OutboxConfig outboxConfig)
{ {
using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) 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(); var outgoingEventArray = outgoingEvents.ToArray();
@ -298,8 +298,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis
Guid? eventId = null, Guid? eventId = null,
string? correlationId = null) string? correlationId = null)
{ {
using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) using (var channel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)).CreateChannelAsync())
.CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true, publisherConfirmationTrackingEnabled: true)))
{ {
await PublishAsync(channel, eventName, body, headersArguments, eventId, correlationId); await PublishAsync(channel, eventName, body, headersArguments, eventId, correlationId);
} }
@ -335,7 +334,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis
await channel.BasicPublishAsync( await channel.BasicPublishAsync(
exchange: AbpRabbitMqEventBusOptions.ExchangeName, exchange: AbpRabbitMqEventBusOptions.ExchangeName,
routingKey: eventName, routingKey: eventName,
mandatory: true, mandatory: false,
basicProperties: properties, basicProperties: properties,
body: body body: body
); );
@ -350,9 +349,7 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, IRabbitMqDis
try try
{ {
using (var temporaryChannel = await(await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)) using (var temporaryChannel = await (await ConnectionPool.GetAsync(AbpRabbitMqEventBusOptions.ConnectionName)).CreateChannelAsync())
.CreateChannelAsync(new CreateChannelOptions(publisherConfirmationsEnabled: true,
publisherConfirmationTrackingEnabled: true)))
{ {
await temporaryChannel.ExchangeDeclarePassiveAsync(AbpRabbitMqEventBusOptions.ExchangeName); await temporaryChannel.ExchangeDeclarePassiveAsync(AbpRabbitMqEventBusOptions.ExchangeName);
} }

19
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ChannelPool.cs

@ -49,14 +49,14 @@ public class ChannelPool : IChannelPool, ISingletonDependency
{ {
using (await Semaphore.LockAsync()) 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)); poolItem = existingChannelPoolItem2;
Channels.TryAdd(channelName, poolItem);
} }
else 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()) using (await Semaphore.LockAsync())
{ {
poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName)); if (Channels.TryGetValue(channelName, out var existingChannelPoolItem3))
Channels.TryAdd(channelName, poolItem); {
poolItem = existingChannelPoolItem3;
}
else
{
poolItem = new ChannelPoolItem(await CreateChannelAsync(channelName, connectionName));
Channels.TryAdd(channelName, poolItem);
}
} }
poolItem.Acquire(); poolItem.Acquire();

21
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs

@ -39,18 +39,25 @@ public class ConnectionPool : IConnectionPool, ISingletonDependency
{ {
using (await Semaphore.LockAsync()) using (await Semaphore.LockAsync())
{ {
var connectionFactory = Options.Connections.GetOrDefault(connectionName);
try try
{ {
connection = await GetConnectionAsync(connectionName, connectionFactory); var connectionFactory = Options.Connections.GetOrDefault(connectionName);
Connections.TryAdd(connectionName, connection); if (Connections.TryGetValue(connectionName, out var existingConnection2))
{
if (!connection.IsOpen) connection = existingConnection2;
}
else
{ {
connection.Dispose();
Connections.TryRemove(connectionName, out _);
connection = await GetConnectionAsync(connectionName, connectionFactory); connection = await GetConnectionAsync(connectionName, connectionFactory);
Connections.TryAdd(connectionName, connection); Connections.TryAdd(connectionName, connection);
if (!connection.IsOpen)
{
connection.Dispose();
Connections.TryRemove(connectionName, out _);
connection = await GetConnectionAsync(connectionName, connectionFactory);
Connections.TryAdd(connectionName, connection);
}
} }
} }
catch (Exception) catch (Exception)

Loading…
Cancel
Save