Browse Source

Implement EventErrorHandler for rabbitmq & rebus

pull/8829/head
liangshiwei 6 years ago
parent
commit
29e0cee1aa
  1. 2
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs
  2. 2
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs
  3. 43
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  4. 82
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqEventErrorHandler.cs
  5. 26
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusModule.cs
  6. 12
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusOptionsSetup.cs
  7. 8
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpRebusEventBusOptions.cs
  8. 44
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusEventErrorHandler.cs
  9. 8
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs
  10. 2
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs
  11. 3
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs
  12. 2
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs
  13. 2
      framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/AbpKafkaOptions.cs
  14. 36
      framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs
  15. 16
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs
  16. 26
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs

2
framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs

@ -49,7 +49,7 @@ namespace Volo.Abp.EventBus.Kafka
MessageConsumerFactory = messageConsumerFactory;
Serializer = serializer;
ProducerPool = producerPool;
DeadLetterTopicName = AbpEventBusOptions.DeadLetterQueue ?? AbpKafkaEventBusOptions.TopicName + "_error";
DeadLetterTopicName = AbpEventBusOptions.DeadLetterName ?? AbpKafkaEventBusOptions.TopicName + "_dead_letter";
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>();

2
framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs

@ -64,7 +64,7 @@ namespace Volo.Abp.EventBus.Kafka
var index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey));
return Options.RetryStrategyOptions.Count < index;
return Options.RetryStrategyOptions.MaxRetryAttempts < index;
}
}
}

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

@ -25,6 +25,7 @@ namespace Volo.Abp.EventBus.RabbitMq
{
protected AbpRabbitMqEventBusOptions AbpRabbitMqEventBusOptions { get; }
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; }
protected AbpEventBusOptions AbpEventBusOptions { get; }
protected IConnectionPool ConnectionPool { get; }
protected IRabbitMqSerializer Serializer { get; }
@ -42,12 +43,14 @@ namespace Volo.Abp.EventBus.RabbitMq
IOptions<AbpDistributedEventBusOptions> distributedEventBusOptions,
IRabbitMqMessageConsumerFactory messageConsumerFactory,
ICurrentTenant currentTenant,
IEventErrorHandler errorHandler)
IEventErrorHandler errorHandler,
IOptions<AbpEventBusOptions> abpEventBusOptions)
: base(serviceScopeFactory, currentTenant, errorHandler)
{
ConnectionPool = connectionPool;
Serializer = serializer;
MessageConsumerFactory = messageConsumerFactory;
AbpEventBusOptions = abpEventBusOptions.Value;
AbpDistributedEventBusOptions = distributedEventBusOptions.Value;
AbpRabbitMqEventBusOptions = options.Value;
@ -57,6 +60,8 @@ namespace Volo.Abp.EventBus.RabbitMq
public void Initialize()
{
const string suffix = "_dead_letter";
Consumer = MessageConsumerFactory.Create(
new ExchangeDeclareConfiguration(
AbpRabbitMqEventBusOptions.ExchangeName,
@ -67,7 +72,12 @@ namespace Volo.Abp.EventBus.RabbitMq
AbpRabbitMqEventBusOptions.ClientName,
durable: true,
exclusive: false,
autoDelete: false
autoDelete: false,
arguments: new Dictionary<string, object>
{
{"x-dead-letter-exchange", AbpRabbitMqEventBusOptions.ExchangeName + suffix},
{"x-dead-letter-routing-key", AbpEventBusOptions.DeadLetterName ?? AbpRabbitMqEventBusOptions.ClientName + suffix}
}
),
AbpRabbitMqEventBusOptions.ConnectionName
);
@ -197,6 +207,35 @@ namespace Volo.Abp.EventBus.RabbitMq
return Task.CompletedTask;
}
public Task PublishAsync(Type eventType, object eventData, Dictionary<string, object> headers)
{
var eventName = EventNameAttribute.GetNameOrDefault(eventType);
var body = Serializer.Serialize(eventData);
using (var channel = ConnectionPool.Get(AbpRabbitMqEventBusOptions.ConnectionName).CreateModel())
{
channel.ExchangeDeclare(
AbpRabbitMqEventBusOptions.ExchangeName,
"direct",
durable: true
);
var properties = channel.CreateBasicProperties();
properties.DeliveryMode = RabbitMqConsts.DeliveryModes.Persistent;
properties.Headers = headers;
channel.BasicPublish(
exchange: AbpRabbitMqEventBusOptions.ExchangeName,
routingKey: eventName,
mandatory: true,
basicProperties: properties,
body: body
);
}
return Task.CompletedTask;
}
private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType)
{
return HandlerFactories.GetOrAdd(

82
framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqEventErrorHandler.cs

@ -0,0 +1,82 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Microsoft.Extensions.Options;
using Volo.Abp.Data;
using Volo.Abp.DependencyInjection;
namespace Volo.Abp.EventBus.RabbitMq
{
public class RabbitMqEventErrorHandler : EventErrorHandlerBase, ISingletonDependency
{
public const string HeadersKey = "headers";
public const string RetryIndexKey = "retryIndex";
protected RabbitMqDistributedEventBus EventBus { get; }
public RabbitMqEventErrorHandler(
IOptions<AbpEventBusOptions> options,
RabbitMqDistributedEventBus eventBus)
: base(options)
{
EventBus = eventBus;
}
protected override async Task Retry(EventExecutionErrorContext context)
{
if (Options.RetryStrategyOptions.IntervalMillisecond > 0)
{
await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond);
}
var headers = context.GetProperty<Dictionary<string, object>>(HeadersKey) ??
new Dictionary<string, object>();
var index = 1;
if (headers.ContainsKey(RetryIndexKey))
{
index = (int) headers[RetryIndexKey];
headers[RetryIndexKey] = ++index;
}
else
{
headers[RetryIndexKey] = index;
}
headers["exceptions"] = context.Exceptions;
await EventBus.PublishAsync(context.EventType, context.EventData, headers);
}
protected override Task MoveToDeadLetter(EventExecutionErrorContext context)
{
if (context.Exceptions.Count == 1)
{
context.Exceptions[0].ReThrow();
}
throw new AggregateException(
"More than one error has occurred while triggering the event: " + context.EventType,
context.Exceptions);
}
protected override bool ShouldRetry(EventExecutionErrorContext context)
{
if (!base.ShouldRetry(context))
{
return false;
}
var headers = context.GetProperty<Dictionary<string, object>>(HeadersKey);
if (headers == null || !headers.ContainsKey(RetryIndexKey))
{
return true;
}
var index = (int) headers[RetryIndexKey];
return Options.RetryStrategyOptions.MaxRetryAttempts < index;
}
}
}

26
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusModule.cs

@ -1,5 +1,6 @@
using Microsoft.Extensions.DependencyInjection;
using Rebus.Handlers;
using Rebus.Retry.Simple;
using Rebus.ServiceProvider;
using Volo.Abp.Modularity;
@ -11,23 +12,28 @@ namespace Volo.Abp.EventBus.Rebus
{
public override void ConfigureServices(ServiceConfigurationContext context)
{
var options = context.Services.ExecutePreConfiguredActions<AbpRebusEventBusOptions>();
var abpEventBusOptions = context.Services.ExecutePreConfiguredActions<AbpEventBusOptions>();
context.Services.AddTransient(typeof(IHandleMessages<>), typeof(RebusDistributedEventHandlerAdapter<>));
Configure<AbpRebusEventBusOptions>(rebusOptions =>
{
rebusOptions.Configurer = options.Configurer;
rebusOptions.Publish = options.Publish;
rebusOptions.InputQueueName = options.InputQueueName;
});
context.Services.ExecutePreConfiguredActions(rebusOptions);
context.Services.AddRebus(configurer =>
{
options.Configurer?.Invoke(configurer);
return configurer;
});
context.Services.AddRebus(configure =>
{
if (abpEventBusOptions.RetryStrategyOptions != null)
{
configure.Options(b =>
b.SimpleRetryStrategy(
errorQueueAddress: abpEventBusOptions.DeadLetterName ?? rebusOptions.InputQueueName + "_dead_letter",
maxDeliveryAttempts: abpEventBusOptions.RetryStrategyOptions.MaxRetryAttempts));
}
rebusOptions.Configurer?.Invoke(configure);
return configure;
});
});
}
public override void OnApplicationInitialization(ApplicationInitializationContext context)

12
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusOptionsSetup.cs

@ -0,0 +1,12 @@
using Microsoft.Extensions.Options;
namespace Volo.Abp.EventBus.Rebus
{
public class AbpEventBusRebusOptionsSetup : IConfigureOptions<AbpEventBusOptions>
{
public void Configure(AbpEventBusOptions options)
{
throw new System.NotImplementedException();
}
}
}

8
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpRebusEventBusOptions.cs

@ -32,7 +32,7 @@ namespace Volo.Abp.EventBus.Rebus
public AbpRebusEventBusOptions()
{
_publish = DefaultPublish;
_configurer = DefaultConfigurer;
_configurer = DefaultConfigure;
}
private async Task DefaultPublish(IBus bus, Type eventType, object eventData)
@ -40,10 +40,10 @@ namespace Volo.Abp.EventBus.Rebus
await bus.Advanced.Routing.Send(InputQueueName, eventData);
}
private void DefaultConfigurer(RebusConfigurer configurer)
private void DefaultConfigure(RebusConfigurer configure)
{
configurer.Subscriptions(s => s.StoreInMemory());
configurer.Transport(t => t.UseInMemoryTransport(new InMemNetwork(), InputQueueName));
configure.Subscriptions(s => s.StoreInMemory());
configure.Transport(t => t.UseInMemoryTransport(new InMemNetwork(), InputQueueName));
}
}
}

44
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusEventErrorHandler.cs

@ -0,0 +1,44 @@
using System;
using System.Threading.Tasks;
using Microsoft.Extensions.Options;
using Volo.Abp.DependencyInjection;
namespace Volo.Abp.EventBus.Rebus
{
public class RebusEventErrorHandler : EventErrorHandlerBase, ISingletonDependency
{
public RebusEventErrorHandler(
IOptions<AbpEventBusOptions> options)
: base(options)
{
}
protected override Task Retry(EventExecutionErrorContext context)
{
Throw(context);
return Task.CompletedTask;
}
protected override Task MoveToDeadLetter(EventExecutionErrorContext context)
{
Throw(context);
return Task.CompletedTask;
}
private void Throw(EventExecutionErrorContext context)
{
// Rebus will automatic retries and error handling: https://github.com/rebus-org/Rebus/wiki/Automatic-retries-and-error-handling
if (context.Exceptions.Count == 1)
{
context.Exceptions[0].ReThrow();
}
throw new AggregateException(
"More than one error has occurred while triggering the event: " + context.EventType,
context.Exceptions);
}
}
}

8
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs

@ -20,6 +20,14 @@ namespace Volo.Abp.EventBus
AddEventHandlers(context.Services);
}
public override void ConfigureServices(ServiceConfigurationContext context)
{
Configure<AbpEventBusOptions>(options =>
{
context.Services.ExecutePreConfiguredActions(options);
});
}
private static void AddEventHandlers(IServiceCollection services)
{
var localHandlers = new List<Type>();

2
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs

@ -8,7 +8,7 @@ namespace Volo.Abp.EventBus
public Func<Type, bool> ErrorHandleSelector { get; set; }
public string DeadLetterQueue { get; set; }
public string DeadLetterName { get; set; }
public AbpEventBusRetryStrategyOptions RetryStrategyOptions { get; set; }

3
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs

@ -2,9 +2,8 @@
{
public class AbpEventBusRetryStrategyOptions
{
public int IntervalMillisecond { get; set; } = 3000;
public int Count { get; set; } = 3;
public int MaxRetryAttempts { get; set; } = 3;
}
}

2
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs

@ -57,7 +57,7 @@ namespace Volo.Abp.EventBus.Local
var index = RetryTracking.GetOrDefault(messageId);
if (Options.RetryStrategyOptions.Count >= index)
if (Options.RetryStrategyOptions.MaxRetryAttempts >= index)
{
RetryTracking.Remove(messageId);
return false;

2
framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/AbpKafkaOptions.cs

@ -14,8 +14,6 @@ namespace Volo.Abp.Kafka
public Action<TopicSpecification> ConfigureTopic { get; set; }
public bool ReQueue { get; set; } = true;
public AbpKafkaOptions()
{
Connections = new KafkaConnections();

36
framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs

@ -27,6 +27,8 @@ namespace Volo.Abp.Kafka
protected AbpKafkaOptions Options { get; }
protected AbpAsyncTimer Timer { get; }
protected ConcurrentBag<Func<Message<string, byte[]>, Task>> Callbacks { get; }
protected IConsumer<string, byte[]> Consumer { get; private set; }
@ -37,22 +39,27 @@ namespace Volo.Abp.Kafka
protected string TopicName { get; private set; }
protected string DeadLetterTopicName { get; private set; }
public KafkaMessageConsumer(
IConsumerPool consumerPool,
IExceptionNotifier exceptionNotifier,
IOptions<AbpKafkaOptions> options,
IProducerPool producerPool)
IProducerPool producerPool,
AbpAsyncTimer timer)
{
ConsumerPool = consumerPool;
ExceptionNotifier = exceptionNotifier;
ProducerPool = producerPool;
Timer = timer;
Options = options.Value;
Logger = NullLogger<KafkaMessageConsumer>.Instance;
Callbacks = new ConcurrentBag<Func<Message<string, byte[]>, Task>>();
Timer.Period = 5000; //5 sec.
Timer.Elapsed = Timer_Elapsed;
Timer.RunOnStart = true;
}
public virtual void Initialize(
@ -68,9 +75,7 @@ namespace Volo.Abp.Kafka
DeadLetterTopicName = deadLetterTopicName;
ConnectionName = connectionName ?? KafkaConnections.DefaultConnectionName;
GroupId = groupId;
AsyncHelper.RunSync(CreateTopicAsync);
Consume();
Timer.Start();
}
public virtual void OnMessageReceived(Func<Message<string, byte[]>, Task> callback)
@ -78,6 +83,14 @@ namespace Volo.Abp.Kafka
Callbacks.Add(callback);
}
protected virtual async Task Timer_Elapsed(AbpAsyncTimer timer)
{
await CreateTopicAsync();
Consume();
Timer.Stop();
}
protected virtual async Task CreateTopicAsync()
{
using (var adminClient = new AdminClientBuilder(Options.Connections.GetOrDefault(ConnectionName)).Build())
@ -158,8 +171,6 @@ namespace Volo.Abp.Kafka
}
catch (Exception ex)
{
await RequeueAsync(consumeResult);
Logger.LogException(ex);
await ExceptionNotifier.NotifyAsync(ex);
}
@ -169,17 +180,6 @@ namespace Volo.Abp.Kafka
}
}
protected virtual async Task RequeueAsync(ConsumeResult<string, byte[]> consumeResult)
{
if (!Options.ReQueue)
{
return;
}
var producer = ProducerPool.Get(ConnectionName);
await producer.ProduceAsync(consumeResult.Topic, consumeResult.Message);
}
public virtual void Dispose()
{
if (Consumer == null)

16
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs

@ -6,8 +6,7 @@ namespace Volo.Abp.RabbitMQ
{
public class QueueDeclareConfiguration
{
[NotNull]
public string QueueName { get; }
[NotNull] public string QueueName { get; }
public bool Durable { get; set; }
@ -18,16 +17,17 @@ namespace Volo.Abp.RabbitMQ
public IDictionary<string, object> Arguments { get; }
public QueueDeclareConfiguration(
[NotNull] string queueName,
bool durable = true,
bool exclusive = false,
bool autoDelete = false)
[NotNull] string queueName,
bool durable = true,
bool exclusive = false,
bool autoDelete = false,
Dictionary<string, object> arguments = null)
{
QueueName = queueName;
Durable = durable;
Exclusive = exclusive;
AutoDelete = autoDelete;
Arguments = new Dictionary<string, object>();
Arguments = arguments ?? new Dictionary<string, object>();
}
public virtual QueueDeclareOk Declare(IModel channel)
@ -41,4 +41,4 @@ namespace Volo.Abp.RabbitMQ
);
}
}
}
}

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

@ -147,6 +147,7 @@ namespace Volo.Abp.RabbitMQ
Channel = ConnectionPool
.Get(ConnectionName)
.CreateModel();
Channel.ExchangeDeclare(
exchange: Exchange.ExchangeName,
type: Exchange.Type,
@ -155,6 +156,23 @@ namespace Volo.Abp.RabbitMQ
arguments: Exchange.Arguments
);
if (Queue.Arguments.ContainsKey("x-dead-letter-exchange") &&
Queue.Arguments.ContainsKey("x-dead-letter-routing-key"))
{
Channel.ExchangeDeclare(
Exchange.Arguments["x-dead-letter-exchange"].ToString(),
Exchange.Type,
Exchange.Durable,
Exchange.AutoDelete
);
Channel.QueueDeclare(
Queue.Arguments["x-dead-letter-routing-key"].ToString(),
Queue.Durable,
Queue.Exclusive,
Queue.AutoDelete);
}
Channel.QueueDeclare(
queue: Queue.QueueName,
durable: Queue.Durable,
@ -194,14 +212,10 @@ namespace Volo.Abp.RabbitMQ
{
try
{
Channel.BasicNack(
basicDeliverEventArgs.DeliveryTag,
multiple: false,
requeue: true
);
Channel.BasicReject(basicDeliverEventArgs.DeliveryTag, false);
}
catch { }
Logger.LogException(ex);
await ExceptionNotifier.NotifyAsync(ex);
}

Loading…
Cancel
Save