diff --git a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs index fd6d773952..46eb1c0827 100644 --- a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs @@ -6,8 +6,10 @@ using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; +using Volo.Abp.Data; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.EventBus.Local; using Volo.Abp.Kafka; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; @@ -34,8 +36,9 @@ namespace Volo.Abp.EventBus.Kafka IKafkaMessageConsumerFactory messageConsumerFactory, IOptions abpDistributedEventBusOptions, IKafkaSerializer serializer, - IProducerPool producerPool) - : base(serviceScopeFactory, currentTenant) + IProducerPool producerPool, + IEventErrorHandler errorHandler) + : base(serviceScopeFactory, currentTenant, errorHandler) { AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; @@ -54,6 +57,7 @@ namespace Volo.Abp.EventBus.Kafka AbpKafkaEventBusOptions.GroupId, AbpKafkaEventBusOptions.ConnectionName); + Consumer.Consume(); Consumer.OnMessageReceived(ProcessEventAsync); SubscribeHandlers(AbpDistributedEventBusOptions.Handlers); @@ -68,9 +72,10 @@ namespace Volo.Abp.EventBus.Kafka return; } - var eventData = Serializer.Deserialize(message.Value, eventType); + var eventMessage = Serializer.Deserialize(message.Value); - await TriggerHandlersAsync(eventType, eventData); + await TriggerHandlersAsync(eventType, eventMessage, + context => { context.SetProperty(KafkaEventErrorHandler.HeadersKey, message.Headers); }); } public IDisposable Subscribe(IDistributedEventHandler handler) where TEvent : class @@ -147,6 +152,11 @@ namespace Volo.Abp.EventBus.Kafka } public override async Task PublishAsync(Type eventType, object eventData) + { + await PublishAsync(eventType, eventData, null); + } + + public virtual async Task PublishAsync(Type eventType, object eventData, Headers headers) { var eventName = EventNameAttribute.GetNameOrDefault(eventType); var body = Serializer.Serialize(eventData); @@ -157,7 +167,7 @@ namespace Volo.Abp.EventBus.Kafka AbpKafkaEventBusOptions.TopicName, new Message { - Key = eventName, Value = body + Key = eventName, Value = body, Headers = headers }); } diff --git a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs new file mode 100644 index 0000000000..65c8561149 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs @@ -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 options, + IKafkaSerializer serializer, + KafkaDistributedEventBus eventBus, + IKafkaMessageConsumerFactory consumerFactory, + IProducerPool producerPool, + IOptions 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(HeadersKey) ?? new Headers(); + var index = Serializer.Deserialize(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 + { + 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(HeadersKey); + var index = 1; + + if (headers == null) + { + return true; + } + + index = Serializer.Deserialize(headers.GetLastBytes(RetryIndexKey)); + + return Options.RetryStrategyOptions.Count < index; + } + } +} diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs index e8c7fc6c5b..5baeda6d9f 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs @@ -41,8 +41,9 @@ namespace Volo.Abp.EventBus.RabbitMq IServiceScopeFactory serviceScopeFactory, IOptions distributedEventBusOptions, IRabbitMqMessageConsumerFactory messageConsumerFactory, - ICurrentTenant currentTenant) - : base(serviceScopeFactory, currentTenant) + ICurrentTenant currentTenant, + IEventErrorHandler errorHandler) + : base(serviceScopeFactory, currentTenant, errorHandler) { ConnectionPool = connectionPool; Serializer = serializer; diff --git a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs index ce3549646e..97bedaf8fb 100644 --- a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs @@ -30,8 +30,9 @@ namespace Volo.Abp.EventBus.Rebus ICurrentTenant currentTenant, IBus rebus, IOptions abpDistributedEventBusOptions, - IOptions abpEventBusRebusOptions) : - base(serviceScopeFactory, currentTenant) + IOptions abpEventBusRebusOptions, + IEventErrorHandler errorHandler) : + base(serviceScopeFactory, currentTenant, errorHandler) { Rebus = rebus; AbpRebusEventBusOptions = abpEventBusRebusOptions.Value; diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs new file mode 100644 index 0000000000..18d3290c25 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs @@ -0,0 +1,21 @@ +using System; + +namespace Volo.Abp.EventBus +{ + public class AbpEventBusOptions + { + public bool EnabledErrorHandle { get; set; } + + public Func ErrorHandleSelector { get; set; } + + public string ErrorQueue { get; set; } + + public AbpEventBusRetryStrategyOptions RetryStrategyOptions { get; set; } + + public void UseRetryStrategy(Action action = null) + { + RetryStrategyOptions = new AbpEventBusRetryStrategyOptions(); + action?.Invoke(RetryStrategyOptions); + } + } +} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs new file mode 100644 index 0000000000..5762e8d7a0 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs @@ -0,0 +1,10 @@ +namespace Volo.Abp.EventBus +{ + public class AbpEventBusRetryStrategyOptions + { + + public int IntervalMillisecond { get; set; } = 3000; + + public int Count { get; set; } = 3; + } +} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs index 80d4134db2..e505f30ccb 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -8,6 +8,7 @@ using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Volo.Abp.Collections; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.EventBus.Local; using Volo.Abp.MultiTenancy; using Volo.Abp.Reflection; @@ -19,10 +20,16 @@ namespace Volo.Abp.EventBus protected ICurrentTenant CurrentTenant { get; } - protected EventBusBase(IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant) + protected IEventErrorHandler ErrorHandler { get; } + + protected EventBusBase( + IServiceScopeFactory serviceScopeFactory, + ICurrentTenant currentTenant, + IEventErrorHandler errorHandler) { ServiceScopeFactory = serviceScopeFactory; CurrentTenant = currentTenant; + ErrorHandler = errorHandler; } /// @@ -89,7 +96,7 @@ namespace Volo.Abp.EventBus /// public abstract Task PublishAsync(Type eventType, object eventData); - public virtual async Task TriggerHandlersAsync(Type eventType, object eventData) + public virtual async Task TriggerHandlersAsync(Type eventType, object eventData, Action onErrorAction = null) { var exceptions = new List(); @@ -97,16 +104,13 @@ namespace Volo.Abp.EventBus if (exceptions.Any()) { - if (exceptions.Count == 1) - { - exceptions[0].ReThrow(); - } - - throw new AggregateException("More than one error has occurred while triggering the event: " + eventType, exceptions); + var context = new EventExecutionErrorContext(exceptions, eventData, eventType); + onErrorAction?.Invoke(context); + await ErrorHandler.Handle(context); } } - protected virtual async Task TriggerHandlersAsync(Type eventType, object eventData, List exceptions) + protected virtual async Task TriggerHandlersAsync(Type eventType, object eventData , List exceptions) { await new SynchronizationContextRemover(); @@ -217,6 +221,11 @@ namespace Volo.Abp.EventBus }; } + protected virtual void OnErrorHandle(EventExecutionErrorContext context) + { + + } + protected class EventTypeWithEventHandlerFactories { public Type EventType { get; } diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs new file mode 100644 index 0000000000..14a27c4221 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs @@ -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 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; + } + } +} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs new file mode 100644 index 0000000000..cf40102579 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs @@ -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 Exceptions { get; } + + public object EventData { get; } + + public Type EventType { get; } + + public EventExecutionErrorContext(List exceptions, object eventData, Type eventType) + { + Exceptions = exceptions; + EventData = eventData; + EventType = eventType; + } + } +} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventErrorHandler.cs new file mode 100644 index 0000000000..27d4951ead --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventErrorHandler.cs @@ -0,0 +1,9 @@ +using System.Threading.Tasks; + +namespace Volo.Abp.EventBus +{ + public interface IEventErrorHandler + { + Task Handle(EventExecutionErrorContext context); + } +} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs index 8c7bef6f2d..0a0ae4f378 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs @@ -7,6 +7,7 @@ using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; +using Volo.Abp.Data; using Volo.Abp.DependencyInjection; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; @@ -31,8 +32,9 @@ namespace Volo.Abp.EventBus.Local public LocalEventBus( IOptions options, IServiceScopeFactory serviceScopeFactory, - ICurrentTenant currentTenant) - : base(serviceScopeFactory, currentTenant) + ICurrentTenant currentTenant, + IEventErrorHandler errorHandler) + : base(serviceScopeFactory, currentTenant, errorHandler) { Options = options.Value; Logger = NullLogger.Instance; @@ -119,19 +121,15 @@ namespace Volo.Abp.EventBus.Local public override async Task PublishAsync(Type eventType, object eventData) { - var exceptions = new List(); - - await TriggerHandlersAsync(eventType, eventData, exceptions); + await PublishAsync(new LocalEventMessage(Guid.NewGuid(), eventData, eventType)); + } - if (exceptions.Any()) + public virtual async Task PublishAsync(LocalEventMessage localEventMessage) + { + await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData, errorContext => { - if (exceptions.Count == 1) - { - exceptions[0].ReThrow(); - } - - throw new AggregateException("More than one error has occurred while triggering the event: " + eventType, exceptions); - } + errorContext.SetProperty("messageId", localEventMessage.MessageId); + }); } protected override IEnumerable GetHandlerFactories(Type eventType) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs new file mode 100644 index 0000000000..db8964d380 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs @@ -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 RetryTracking { get; } + + public LocalEventErrorHandler( + IOptions options, + ILocalEventBus localEventBus) + : base(options) + { + LocalEventBus = localEventBus; + RetryTracking = new Dictionary(); + } + + protected override async Task Retry(EventExecutionErrorContext context) + { + if (Options.RetryStrategyOptions.IntervalMillisecond > 0) + { + await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond); + } + + var messageId = context.GetProperty("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("messageId"); + + var index = RetryTracking.GetOrDefault(messageId); + + if (Options.RetryStrategyOptions.Count >= index) + { + RetryTracking.Remove(messageId); + return false; + } + + RetryTracking[messageId] = ++index; + + return true; + } + } +} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventMessage.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventMessage.cs new file mode 100644 index 0000000000..550de897bd --- /dev/null +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventMessage.cs @@ -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; + } + } +} diff --git a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaMessageConsumer.cs b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaMessageConsumer.cs index 87872b31a2..721f0b5a9b 100644 --- a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaMessageConsumer.cs +++ b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaMessageConsumer.cs @@ -7,5 +7,7 @@ namespace Volo.Abp.Kafka public interface IKafkaMessageConsumer { void OnMessageReceived(Func, Task> callback); + + void Consume(); } } diff --git a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaSerializer.cs b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaSerializer.cs index 58e718831c..a283eb7e50 100644 --- a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaSerializer.cs +++ b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/IKafkaSerializer.cs @@ -7,5 +7,7 @@ namespace Volo.Abp.Kafka byte[] Serialize(object obj); object Deserialize(byte[] value, Type type); + + T Deserialize(byte[] value); } } diff --git a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs index 3b8f022012..f5c726063e 100644 --- a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs +++ b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs @@ -63,7 +63,6 @@ namespace Volo.Abp.Kafka GroupId = groupId; AsyncHelper.RunSync(CreateTopicAsync); - Consume(); } public virtual void OnMessageReceived(Func, Task> callback) @@ -98,7 +97,7 @@ namespace Volo.Abp.Kafka } } - protected virtual void Consume() + public virtual void Consume() { Consumer = ConsumerPool.Get(GroupId, ConnectionName); diff --git a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/Utf8JsonKafkaSerializer.cs b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/Utf8JsonKafkaSerializer.cs index a04125f8a6..a8a199c140 100644 --- a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/Utf8JsonKafkaSerializer.cs +++ b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/Utf8JsonKafkaSerializer.cs @@ -23,5 +23,10 @@ namespace Volo.Abp.Kafka { return _jsonSerializer.Deserialize(type, Encoding.UTF8.GetString(value)); } + + public T Deserialize(byte[] value) + { + return _jsonSerializer.Deserialize(Encoding.UTF8.GetString(value)); + } } } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs index 771d1a05a6..2de5cd2191 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/IRabbitMqSerializer.cs @@ -7,5 +7,7 @@ namespace Volo.Abp.RabbitMQ byte[] Serialize(object obj); object Deserialize(byte[] value, Type type); + + T Deserialize(byte[] value); } } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs index c179193fdd..cc815f686a 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/Utf8JsonRabbitMqSerializer.cs @@ -23,5 +23,10 @@ namespace Volo.Abp.RabbitMQ { return _jsonSerializer.Deserialize(type, Encoding.UTF8.GetString(value)); } + + public T Deserialize(byte[] value) + { + return _jsonSerializer.Deserialize(Encoding.UTF8.GetString(value)); + } } -} \ No newline at end of file +}