KafkaCommandBus
class KafkaCommandBus(topicConverter: CommandTopicConverter = DefaultCommandTopicConverter(), senderOptions: SenderOptions<String, String>, receiverOptions: ReceiverOptions<String, String>, receiverOptionsCustomizer: ReceiverOptionsCustomizer = NoOpReceiverOptionsCustomizer, receiverPolicy: KafkaReceiverPolicy = KafkaReceiverPolicy(), recordDecodeFailureHandler: KafkaRecordDecodeFailureHandler = FailKafkaRecordDecodeFailureHandler) : DistributedCommandBus, AbstractKafkaBus<CommandMessage<*>, ServerCommandExchange<*>>
Constructors
Link copied to clipboard
constructor(topicConverter: CommandTopicConverter = DefaultCommandTopicConverter(), senderOptions: SenderOptions<String, String>, receiverOptions: ReceiverOptions<String, String>, receiverOptionsCustomizer: ReceiverOptionsCustomizer = NoOpReceiverOptionsCustomizer, receiverPolicy: KafkaReceiverPolicy = KafkaReceiverPolicy(), recordDecodeFailureHandler: KafkaRecordDecodeFailureHandler = FailKafkaRecordDecodeFailureHandler)
Functions
Link copied to clipboard
open override fun CommandMessage<*>.toExchange(receiverOffset: ReceiverOffset): ServerCommandExchange<*>