|
|
|
@ -8,6 +8,7 @@ using Microsoft.Extensions.DependencyInjection; |
|
|
|
using Microsoft.Extensions.Options; |
|
|
|
using Volo.Abp.DependencyInjection; |
|
|
|
using Volo.Abp.EventBus.Distributed; |
|
|
|
using Volo.Abp.EventBus.Local; |
|
|
|
using Volo.Abp.Guids; |
|
|
|
using Volo.Abp.Kafka; |
|
|
|
using Volo.Abp.MultiTenancy; |
|
|
|
@ -40,7 +41,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
IProducerPool producerPool, |
|
|
|
IGuidGenerator guidGenerator, |
|
|
|
IClock clock, |
|
|
|
IEventHandlerInvoker eventHandlerInvoker) |
|
|
|
IEventHandlerInvoker eventHandlerInvoker, |
|
|
|
ILocalEventBus localEventBus) |
|
|
|
: base( |
|
|
|
serviceScopeFactory, |
|
|
|
currentTenant, |
|
|
|
@ -48,7 +50,8 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
abpDistributedEventBusOptions, |
|
|
|
guidGenerator, |
|
|
|
clock, |
|
|
|
eventHandlerInvoker) |
|
|
|
eventHandlerInvoker, |
|
|
|
localEventBus) |
|
|
|
{ |
|
|
|
AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; |
|
|
|
MessageConsumerFactory = messageConsumerFactory; |
|
|
|
@ -88,7 +91,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
|
|
|
|
var eventData = Serializer.Deserialize(message.Value, eventType); |
|
|
|
|
|
|
|
await TriggerHandlersAsync(eventType, eventData); |
|
|
|
await TriggerHandlersDirectAsync(eventType, eventData); |
|
|
|
} |
|
|
|
|
|
|
|
public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) |
|
|
|
@ -177,11 +180,11 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
unitOfWork.AddOrReplaceDistributedEvent(eventRecord); |
|
|
|
} |
|
|
|
|
|
|
|
public override Task PublishFromOutboxAsync( |
|
|
|
public override async Task PublishFromOutboxAsync( |
|
|
|
OutgoingEventInfo outgoingEvent, |
|
|
|
OutboxConfig outboxConfig) |
|
|
|
{ |
|
|
|
return PublishAsync( |
|
|
|
await PublishAsync( |
|
|
|
AbpKafkaEventBusOptions.TopicName, |
|
|
|
outgoingEvent.EventName, |
|
|
|
outgoingEvent.EventData, |
|
|
|
@ -190,21 +193,28 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
{ "messageId", System.Text.Encoding.UTF8.GetBytes(outgoingEvent.Id.ToString("N")) } |
|
|
|
} |
|
|
|
); |
|
|
|
|
|
|
|
await TriggerDistributedEventSentAsync(new DistributedEventSent() |
|
|
|
{ |
|
|
|
Source = DistributedEventSource.Outbox, |
|
|
|
EventName = outgoingEvent.EventName, |
|
|
|
EventData = outgoingEvent.EventData |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
public override Task PublishManyFromOutboxAsync(IEnumerable<OutgoingEventInfo> outgoingEvents, OutboxConfig outboxConfig) |
|
|
|
public async override Task PublishManyFromOutboxAsync(IEnumerable<OutgoingEventInfo> outgoingEvents, OutboxConfig outboxConfig) |
|
|
|
{ |
|
|
|
var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); |
|
|
|
var outgoingEventArray = outgoingEvents.ToArray(); |
|
|
|
|
|
|
|
|
|
|
|
foreach (var outgoingEvent in outgoingEventArray) |
|
|
|
{ |
|
|
|
var messageId = outgoingEvent.Id.ToString("N"); |
|
|
|
var headers = new Headers |
|
|
|
var headers = new Headers |
|
|
|
{ |
|
|
|
{ "messageId", System.Text.Encoding.UTF8.GetBytes(messageId)} |
|
|
|
}; |
|
|
|
|
|
|
|
|
|
|
|
producer.Produce( |
|
|
|
AbpKafkaEventBusOptions.TopicName, |
|
|
|
new Message<string, byte[]> |
|
|
|
@ -213,9 +223,14 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
Value = outgoingEvent.EventData, |
|
|
|
Headers = headers |
|
|
|
}); |
|
|
|
|
|
|
|
await TriggerDistributedEventSentAsync(new DistributedEventSent() |
|
|
|
{ |
|
|
|
Source = DistributedEventSource.Outbox, |
|
|
|
EventName = outgoingEvent.EventName, |
|
|
|
EventData = outgoingEvent.EventData |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
return Task.CompletedTask; |
|
|
|
} |
|
|
|
|
|
|
|
public async override Task ProcessFromInboxAsync( |
|
|
|
@ -230,7 +245,7 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
|
|
|
|
var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); |
|
|
|
var exceptions = new List<Exception>(); |
|
|
|
await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); |
|
|
|
await TriggerHandlersFromInboxAsync(eventType, eventData, exceptions, inboxConfig); |
|
|
|
if (exceptions.Any()) |
|
|
|
{ |
|
|
|
ThrowOriginalExceptions(eventType, exceptions); |
|
|
|
@ -251,9 +266,9 @@ public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDepen |
|
|
|
} |
|
|
|
|
|
|
|
private Task<DeliveryResult<string, byte[]>> PublishAsync( |
|
|
|
string topicName, |
|
|
|
string topicName, |
|
|
|
string eventName, |
|
|
|
byte[] body, |
|
|
|
byte[] body, |
|
|
|
Headers headers) |
|
|
|
{ |
|
|
|
var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); |
|
|
|
|