KafkaStateEventBus

class KafkaStateEventBus(topicConverter: StateEventTopicConverter = DefaultStateEventTopicConverter(), senderOptions: SenderOptions<String, String>, receiverOptions: ReceiverOptions<String, String>, receiverOptionsCustomizer: ReceiverOptionsCustomizer = NoOpReceiverOptionsCustomizer, receiverPolicy: KafkaReceiverPolicy = KafkaReceiverPolicy(), recordDecodeFailureHandler: KafkaRecordDecodeFailureHandler = FailKafkaRecordDecodeFailureHandler) : DistributedStateEventBus, AbstractKafkaBus<StateEvent<*>, StateEventExchange<*>>

Constructors

Link copied to clipboard
constructor(topicConverter: StateEventTopicConverter = DefaultStateEventTopicConverter(), 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<StateEvent<*>>

Functions

Link copied to clipboard
open override fun StateEvent<*>.toExchange(receiverOffset: ReceiverOffset): StateEventExchange<*>