Browse Source

Fix duplicate consumption of events when connection is broken

pull/22119/head
liangshiwei 2 years ago
parent
commit
de10dcc2c4
  1. 2
      framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs
  2. 1
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs
  3. 38
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs

2
framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs

@ -122,7 +122,7 @@ public class JobQueue<TArgs> : IJobQueue<TArgs>
protected virtual Task EnsureInitializedAsync() protected virtual Task EnsureInitializedAsync()
{ {
if (ChannelAccessor != null) if (ChannelAccessor != null && ChannelAccessor.Channel.IsOpen)
{ {
return Task.CompletedTask; return Task.CompletedTask;
} }

1
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/AbpRabbitMqModule.cs

@ -20,6 +20,7 @@ public class AbpRabbitMqModule : AbpModule
foreach (var connectionFactory in options.Connections.Values) foreach (var connectionFactory in options.Connections.Values)
{ {
connectionFactory.DispatchConsumersAsync = true; connectionFactory.DispatchConsumersAsync = true;
connectionFactory.AutomaticRecoveryEnabled = false;
} }
}); });
} }

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

@ -24,21 +24,19 @@ public class ConnectionPool : IConnectionPool, ISingletonDependency
public virtual IConnection Get(string? connectionName = null) public virtual IConnection Get(string? connectionName = null)
{ {
connectionName ??= RabbitMqConnections.DefaultConnectionName; connectionName ??= RabbitMqConnections.DefaultConnectionName;
var connectionFactory = Options.Connections.GetOrDefault(connectionName);
try try
{ {
var lazyConnection = Connections.GetOrAdd( var connection = GetConnection(connectionName, connectionFactory);
connectionName, () => new Lazy<IConnection>(() =>
{ if (connection.IsOpen)
var connection = Options.Connections.GetOrDefault(connectionName); {
var hostnames = connection.HostName.TrimEnd(';').Split(';'); return connection;
// Handle Rabbit MQ Cluster. }
return hostnames.Length == 1 ? connection.CreateConnection() : connection.CreateConnection(hostnames);
connection.Dispose();
}) Connections.TryRemove(connectionName, out _);
); return GetConnection(connectionName, connectionFactory);
return lazyConnection.Value;
} }
catch (Exception) catch (Exception)
{ {
@ -47,6 +45,20 @@ public class ConnectionPool : IConnectionPool, ISingletonDependency
} }
} }
protected virtual IConnection GetConnection(string connectionName, ConnectionFactory connectionFactory)
{
return Connections.GetOrAdd(
connectionName, () => new Lazy<IConnection>(() =>
{
var hostnames = connectionFactory.HostName.TrimEnd(';').Split(';');
// Handle Rabbit MQ Cluster.
return hostnames.Length == 1
? connectionFactory.CreateConnection()
: connectionFactory.CreateConnection(hostnames);
})
).Value;
}
public void Dispose() public void Dispose()
{ {
if (_isDisposed) if (_isDisposed)

Loading…
Cancel
Save