Browse Source

Set DispatchConsumersAsync to true

pull/7071/head
liangshiwei 6 years ago
parent
commit
0adc4afb5c
  1. 5
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ConnectionPool.cs
  2. 6
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConnections.cs
  3. 23
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs

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

@ -22,8 +22,7 @@ namespace Volo.Abp.RabbitMQ
public virtual IConnection Get(string connectionName = null) public virtual IConnection Get(string connectionName = null)
{ {
connectionName = connectionName connectionName ??= RabbitMqConnections.DefaultConnectionName;
?? RabbitMqConnections.DefaultConnectionName;
return Connections.GetOrAdd( return Connections.GetOrAdd(
connectionName, connectionName,
@ -58,4 +57,4 @@ namespace Volo.Abp.RabbitMQ
Connections.Clear(); Connections.Clear();
} }
} }
} }

6
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConnections.cs

@ -9,7 +9,7 @@ namespace Volo.Abp.RabbitMQ
public class RabbitMqConnections : Dictionary<string, ConnectionFactory> public class RabbitMqConnections : Dictionary<string, ConnectionFactory>
{ {
public const string DefaultConnectionName = "Default"; public const string DefaultConnectionName = "Default";
[NotNull] [NotNull]
public ConnectionFactory Default public ConnectionFactory Default
{ {
@ -19,7 +19,7 @@ namespace Volo.Abp.RabbitMQ
public RabbitMqConnections() public RabbitMqConnections()
{ {
Default = new ConnectionFactory(); Default = new ConnectionFactory() { DispatchConsumersAsync = true };
} }
public ConnectionFactory GetOrDefault(string connectionName) public ConnectionFactory GetOrDefault(string connectionName)
@ -32,4 +32,4 @@ namespace Volo.Abp.RabbitMQ
return Default; return Default;
} }
} }
} }

23
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs

@ -143,10 +143,10 @@ namespace Volo.Abp.RabbitMQ
try try
{ {
var channel = ConnectionPool Channel = ConnectionPool
.Get(ConnectionName) .Get(ConnectionName)
.CreateModel(); .CreateModel();
channel.ExchangeDeclare( Channel.ExchangeDeclare(
exchange: Exchange.ExchangeName, exchange: Exchange.ExchangeName,
type: Exchange.Type, type: Exchange.Type,
durable: Exchange.Durable, durable: Exchange.Durable,
@ -154,7 +154,7 @@ namespace Volo.Abp.RabbitMQ
arguments: Exchange.Arguments arguments: Exchange.Arguments
); );
channel.QueueDeclare( Channel.QueueDeclare(
queue: Queue.QueueName, queue: Queue.QueueName,
durable: Queue.Durable, durable: Queue.Durable,
exclusive: Queue.Exclusive, exclusive: Queue.Exclusive,
@ -162,19 +162,14 @@ namespace Volo.Abp.RabbitMQ
arguments: Queue.Arguments arguments: Queue.Arguments
); );
var consumer = new EventingBasicConsumer(channel); var consumer = new AsyncEventingBasicConsumer(Channel);
consumer.Received += async (model, basicDeliverEventArgs) => consumer.Received += HandleIncomingMessageAsync;
{
await HandleIncomingMessageAsync(channel, basicDeliverEventArgs);
};
channel.BasicConsume( Channel.BasicConsume(
queue: Queue.QueueName, queue: Queue.QueueName,
autoAck: false, autoAck: false,
consumer: consumer consumer: consumer
); );
Channel = channel;
} }
catch (Exception ex) catch (Exception ex)
{ {
@ -183,16 +178,16 @@ namespace Volo.Abp.RabbitMQ
} }
} }
protected virtual async Task HandleIncomingMessageAsync(IModel channel, BasicDeliverEventArgs basicDeliverEventArgs) protected virtual async Task HandleIncomingMessageAsync(object sender, BasicDeliverEventArgs basicDeliverEventArgs)
{ {
try try
{ {
foreach (var callback in Callbacks) foreach (var callback in Callbacks)
{ {
await callback(channel, basicDeliverEventArgs); await callback(Channel, basicDeliverEventArgs);
} }
channel.BasicAck(basicDeliverEventArgs.DeliveryTag, multiple: false); Channel.BasicAck(basicDeliverEventArgs.DeliveryTag, multiple: false);
} }
catch (Exception ex) catch (Exception ex)
{ {

Loading…
Cancel
Save