AbstractRedisMessageBus

abstract class AbstractRedisMessageBus<M : Message<*, *>, AggregateIdCapable, NamedAggregate, E : MessageExchange<*, M>>(redisTemplate: ReactiveStringRedisTemplate, topicConverter: AggregateTopicConverter, pollTimeout: Duration = Duration.ofSeconds(2), recoveryOptions: RedisStreamRecoveryOptions = RedisStreamRecoveryOptions.DEFAULT, messageBusObserver: RedisMessageBusObserver = RedisMessageBusObserver.NOOP) : DistributedMessageBus<M, E>

Inheritors

Constructors

Link copied to clipboard
constructor(redisTemplate: ReactiveStringRedisTemplate, topicConverter: AggregateTopicConverter, pollTimeout: Duration = Duration.ofSeconds(2), recoveryOptions: RedisStreamRecoveryOptions = RedisStreamRecoveryOptions.DEFAULT, messageBusObserver: RedisMessageBusObserver = RedisMessageBusObserver.NOOP)

Types

Link copied to clipboard
object Companion

Properties

Link copied to clipboard
abstract val messageType: Class<M>

Functions

Link copied to clipboard
open override fun close()
Link copied to clipboard
open override fun receive(subscription: MessageSubscription): Flux<E>
Link copied to clipboard
open override fun receiver(subscription: MessageSubscription): MessageReceiver<E>
Link copied to clipboard
Link copied to clipboard
open override fun send(message: M): Mono<Void>
Link copied to clipboard
abstract fun M.toExchange(acknowledgePublisher: Mono<Void>): E