Browse Source

Use AsyncEventingBasicConsumer for RabbitMQ's JobQueue.

pull/6413/head
maliming 6 years ago
parent
commit
453ac6ab9d
  1. 6
      framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs

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

@ -135,7 +135,7 @@ namespace Volo.Abp.BackgroundJobs.RabbitMQ
if (AbpBackgroundJobOptions.IsJobExecutionEnabled) if (AbpBackgroundJobOptions.IsJobExecutionEnabled)
{ {
Consumer = new EventingBasicConsumer(ChannelAccessor.Channel); var Consumer = new AsyncEventingBasicConsumer(ChannelAccessor.Channel);
Consumer.Received += MessageReceived; Consumer.Received += MessageReceived;
//TODO: What BasicConsume returns? //TODO: What BasicConsume returns?
@ -173,7 +173,7 @@ namespace Volo.Abp.BackgroundJobs.RabbitMQ
return properties; return properties;
} }
protected virtual void MessageReceived(object sender, BasicDeliverEventArgs ea) protected virtual async Task MessageReceived(object sender, BasicDeliverEventArgs ea)
{ {
using (var scope = ServiceScopeFactory.CreateScope()) using (var scope = ServiceScopeFactory.CreateScope())
{ {
@ -185,7 +185,7 @@ namespace Volo.Abp.BackgroundJobs.RabbitMQ
try try
{ {
AsyncHelper.RunSync(() => JobExecuter.ExecuteAsync(context)); await JobExecuter.ExecuteAsync(context);
ChannelAccessor.Channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false); ChannelAccessor.Channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
} }
catch (BackgroundJobExecutionException) catch (BackgroundJobExecutionException)

Loading…
Cancel
Save