TransportReceiver

The consumer Transport.open returns.

records must be subscribed before readiness. readiness completes once records published from then on can no longer be missed (Kafka: partitions assigned and their positions anchored; Redis: the consumer groups exist), or fails when that setup fails. openProcessing starts consumption on a transport that holds it back until the runtime is ready; close revokes that admission before the subscription is cancelled. Both are prompt and idempotent.

Properties

Link copied to clipboard
open val readiness: Mono<Void>
Link copied to clipboard
abstract val records: Flux<TransportRecord>

Functions

Link copied to clipboard
open fun close()
Link copied to clipboard
open fun openProcessing()