KafkaDomainEventBus

class KafkaDomainEventBus(topicConverter: EventStreamTopicConverter = DefaultEventStreamTopicConverter(), senderOptions: SenderOptions<String, String>, receiverOptions: ReceiverOptions<String, String>, receiverOptionsCustomizer: ReceiverOptionsCustomizer = NoOpReceiverOptionsCustomizer, receiverPolicy: KafkaReceiverPolicy = KafkaReceiverPolicy(), recordDecodeFailureHandler: KafkaRecordDecodeFailureHandler = FailKafkaRecordDecodeFailureHandler) : DistributedDomainEventBus, AbstractKafkaBus<DomainEventStream, EventStreamExchange>

Constructors

Link copied to clipboard
constructor(topicConverter: EventStreamTopicConverter = DefaultEventStreamTopicConverter(), senderOptions: SenderOptions<String, String>, receiverOptions: ReceiverOptions<String, String>, receiverOptionsCustomizer: ReceiverOptionsCustomizer = NoOpReceiverOptionsCustomizer, receiverPolicy: KafkaReceiverPolicy = KafkaReceiverPolicy(), recordDecodeFailureHandler: KafkaRecordDecodeFailureHandler = FailKafkaRecordDecodeFailureHandler)

Properties

Link copied to clipboard
open override val messageType: Class<DomainEventStream>

Functions

Link copied to clipboard
open override fun DomainEventStream.toExchange(receiverOffset: ReceiverOffset): EventStreamExchange