mirror of https://github.com/abpframework/abp.git
19 changed files with 369 additions and 34 deletions
@ -0,0 +1,90 @@ |
|||||
|
using System.Threading.Tasks; |
||||
|
using Confluent.Kafka; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.Data; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.Kafka; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus.Kafka |
||||
|
{ |
||||
|
public class KafkaEventErrorHandler : EventErrorHandlerBase, ISingletonDependency |
||||
|
{ |
||||
|
public const string HeadersKey = "headers"; |
||||
|
public const string RetryIndexKey = "retryIndex"; |
||||
|
|
||||
|
protected IKafkaSerializer Serializer { get; } |
||||
|
protected KafkaDistributedEventBus EventBus { get; } |
||||
|
protected IProducerPool ProducerPool { get; } |
||||
|
protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; } |
||||
|
|
||||
|
protected string ErrorTopicName { get; } |
||||
|
|
||||
|
public KafkaEventErrorHandler( |
||||
|
IOptions<AbpEventBusOptions> options, |
||||
|
IKafkaSerializer serializer, |
||||
|
KafkaDistributedEventBus eventBus, |
||||
|
IKafkaMessageConsumerFactory consumerFactory, |
||||
|
IProducerPool producerPool, |
||||
|
IOptions<AbpKafkaEventBusOptions> abpKafkaEventBusOptions) : base(options) |
||||
|
{ |
||||
|
Serializer = serializer; |
||||
|
EventBus = eventBus; |
||||
|
ProducerPool = producerPool; |
||||
|
AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; |
||||
|
|
||||
|
ErrorTopicName = options.Value.ErrorQueue ?? abpKafkaEventBusOptions.Value.TopicName + "_error"; |
||||
|
consumerFactory.Create(ErrorTopicName, string.Empty, abpKafkaEventBusOptions.Value.ConnectionName); |
||||
|
} |
||||
|
|
||||
|
protected override async Task Retry(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
if (Options.RetryStrategyOptions.IntervalMillisecond > 0) |
||||
|
{ |
||||
|
await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond); |
||||
|
} |
||||
|
|
||||
|
var headers = context.GetProperty<Headers>(HeadersKey) ?? new Headers(); |
||||
|
var index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey)); |
||||
|
|
||||
|
headers.Remove(RetryIndexKey); |
||||
|
headers.Add(RetryIndexKey, Serializer.Serialize(++index)); |
||||
|
|
||||
|
await EventBus.PublishAsync(context.EventType, context.EventData, headers); |
||||
|
} |
||||
|
|
||||
|
protected override async Task MoveToErrorQueue(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); |
||||
|
var eventName = EventNameAttribute.GetNameOrDefault(context.EventType); |
||||
|
var body = Serializer.Serialize(context.EventData); |
||||
|
|
||||
|
await producer.ProduceAsync( |
||||
|
AbpKafkaEventBusOptions.TopicName, |
||||
|
new Message<string, byte[]> |
||||
|
{ |
||||
|
Key = eventName, Value = body, |
||||
|
Headers = new Headers {{"exceptions", Serializer.Serialize(context.Exceptions)}} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
protected override bool ShouldRetry(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
if (!base.ShouldRetry(context)) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
var headers = context.GetProperty<Headers>(HeadersKey); |
||||
|
var index = 1; |
||||
|
|
||||
|
if (headers == null) |
||||
|
{ |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey)); |
||||
|
|
||||
|
return Options.RetryStrategyOptions.Count < index; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,21 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus |
||||
|
{ |
||||
|
public class AbpEventBusOptions |
||||
|
{ |
||||
|
public bool EnabledErrorHandle { get; set; } |
||||
|
|
||||
|
public Func<Type, bool> ErrorHandleSelector { get; set; } |
||||
|
|
||||
|
public string ErrorQueue { get; set; } |
||||
|
|
||||
|
public AbpEventBusRetryStrategyOptions RetryStrategyOptions { get; set; } |
||||
|
|
||||
|
public void UseRetryStrategy(Action<AbpEventBusRetryStrategyOptions> action = null) |
||||
|
{ |
||||
|
RetryStrategyOptions = new AbpEventBusRetryStrategyOptions(); |
||||
|
action?.Invoke(RetryStrategyOptions); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,10 @@ |
|||||
|
namespace Volo.Abp.EventBus |
||||
|
{ |
||||
|
public class AbpEventBusRetryStrategyOptions |
||||
|
{ |
||||
|
|
||||
|
public int IntervalMillisecond { get; set; } = 3000; |
||||
|
|
||||
|
public int Count { get; set; } = 3; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,57 @@ |
|||||
|
using System.Collections.Generic; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.Data; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus |
||||
|
{ |
||||
|
public abstract class EventErrorHandlerBase : IEventErrorHandler |
||||
|
{ |
||||
|
protected AbpEventBusOptions Options { get; } |
||||
|
|
||||
|
public EventErrorHandlerBase(IOptions<AbpEventBusOptions> options) |
||||
|
{ |
||||
|
Options = options.Value; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task Handle(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
if (!ShouldHandle(context)) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
if (ShouldRetry(context)) |
||||
|
{ |
||||
|
await Retry(context); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
await MoveToErrorQueue(context); |
||||
|
} |
||||
|
|
||||
|
protected abstract Task Retry(EventExecutionErrorContext context); |
||||
|
|
||||
|
protected abstract Task MoveToErrorQueue(EventExecutionErrorContext context); |
||||
|
|
||||
|
protected virtual bool ShouldHandle(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
if (!Options.EnabledErrorHandle) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
if (Options.ErrorHandleSelector != null) |
||||
|
{ |
||||
|
return Options.ErrorHandleSelector.Invoke(context.EventType); |
||||
|
} |
||||
|
|
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
protected virtual bool ShouldRetry(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
return Options.RetryStrategyOptions == null && false; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,23 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using Volo.Abp.EventBus.Local; |
||||
|
using Volo.Abp.ObjectExtending; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus |
||||
|
{ |
||||
|
public class EventExecutionErrorContext : ExtensibleObject |
||||
|
{ |
||||
|
public IReadOnlyList<Exception> Exceptions { get; } |
||||
|
|
||||
|
public object EventData { get; } |
||||
|
|
||||
|
public Type EventType { get; } |
||||
|
|
||||
|
public EventExecutionErrorContext(List<Exception> exceptions, object eventData, Type eventType) |
||||
|
{ |
||||
|
Exceptions = exceptions; |
||||
|
EventData = eventData; |
||||
|
EventType = eventType; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,9 @@ |
|||||
|
using System.Threading.Tasks; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus |
||||
|
{ |
||||
|
public interface IEventErrorHandler |
||||
|
{ |
||||
|
Task Handle(EventExecutionErrorContext context); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,71 @@ |
|||||
|
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.Local |
||||
|
{ |
||||
|
public class LocalEventErrorHandler : EventErrorHandlerBase, ISingletonDependency |
||||
|
{ |
||||
|
protected ILocalEventBus LocalEventBus { get; } |
||||
|
protected Dictionary<Guid, int> RetryTracking { get; } |
||||
|
|
||||
|
public LocalEventErrorHandler( |
||||
|
IOptions<AbpEventBusOptions> options, |
||||
|
ILocalEventBus localEventBus) |
||||
|
: base(options) |
||||
|
{ |
||||
|
LocalEventBus = localEventBus; |
||||
|
RetryTracking = new Dictionary<Guid, int>(); |
||||
|
} |
||||
|
|
||||
|
protected override async Task Retry(EventExecutionErrorContext context) |
||||
|
{ |
||||
|
if (Options.RetryStrategyOptions.IntervalMillisecond > 0) |
||||
|
{ |
||||
|
await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond); |
||||
|
} |
||||
|
|
||||
|
var messageId = context.GetProperty<Guid>("messageId"); |
||||
|
|
||||
|
await LocalEventBus.PublishAsync(context.EventType, |
||||
|
new LocalEventMessage(messageId, context.EventData, context.EventType)); |
||||
|
} |
||||
|
|
||||
|
protected override Task MoveToErrorQueue(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 messageId = context.GetProperty<Guid>("messageId"); |
||||
|
|
||||
|
var index = RetryTracking.GetOrDefault(messageId); |
||||
|
|
||||
|
if (Options.RetryStrategyOptions.Count >= index) |
||||
|
{ |
||||
|
RetryTracking.Remove(messageId); |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
RetryTracking[messageId] = ++index; |
||||
|
|
||||
|
return true; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,20 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus.Local |
||||
|
{ |
||||
|
public class LocalEventMessage |
||||
|
{ |
||||
|
public Guid MessageId { get; } |
||||
|
|
||||
|
public object EventData { get; } |
||||
|
|
||||
|
public Type EventType { get; } |
||||
|
|
||||
|
public LocalEventMessage(Guid messageId, object eventData, Type eventType) |
||||
|
{ |
||||
|
MessageId = messageId; |
||||
|
EventData = eventData; |
||||
|
EventType = eventType; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
Loading…
Reference in new issue