Class IntegrationReactiveUtils

java.lang.Object
org.springframework.integration.util.IntegrationReactiveUtils

public final class IntegrationReactiveUtils extends Object
Utilities for adapting integration components to/from reactive types.
Since:
5.3
Author:
Artem Bilan, Fardan An
  • Field Details

    • DELAY_WHEN_EMPTY_KEY

      public static final String DELAY_WHEN_EMPTY_KEY
      The subscriber context entry for Flux.delayElements(Duration) from the Mono.repeatWhenEmpty(java.util.function.Function).
      See Also:
    • DEFAULT_DELAY_WHEN_EMPTY

      public static final Duration DEFAULT_DELAY_WHEN_EMPTY
      A default delay before repeating an empty source Mono as 1 second Duration.
    • isContextPropagationPresent

      public static final boolean isContextPropagationPresent
      The indicator that io.micrometer:context-propagation library is on classpath.
      Since:
      6.2.5
  • Method Details

    • captureReactorContext

      public static reactor.util.context.ContextView captureReactorContext()
      Capture a Reactor ContextView from the current thread local state according to the ContextSnapshotFactory logic. If no io.micrometer:context-propagation library is on classpath, the Context.empty() is returned.
      Returns:
      the Reactor ContextView from the current thread local state or Context.empty().
      Since:
      6.2.5
    • setThreadLocalsFromReactorContext

      public static @Nullable AutoCloseable setThreadLocalsFromReactorContext(reactor.util.context.ContextView context)
      Populate thread local variables from the provided Reactor ContextView according to the ContextSnapshotFactory logic.
      Parameters:
      context - the Reactor ContextView to populate from.
      Returns:
      the ContextSnapshot.Scope as a AutoCloseable to not pollute the target classpath. Can be cast if necessary. Or null if there is no io.micrometer:context-propagation library is on classpath.
      Since:
      6.2.5
    • messageSourceToFlux

      public static <T> reactor.core.publisher.Flux<Message<T>> messageSourceToFlux(MessageSource<T> messageSource)
      Wrap a provided MessageSource into a Flux for pulling messages on demand. When MessageSource.receive() returns null, the source Mono goes to the Mono.repeatWhenEmpty(Function) state and performs a delay based on the DELAY_WHEN_EMPTY_KEY Duration entry in the subscriber context or falls back to 1-second duration. If a produced message has an IntegrationMessageHeaderAccessor.ACKNOWLEDGMENT_CALLBACK header it is ack'ed in Flux.doOnNext(Consumer) when this flux emits to its subscriber, and nack'ed in the Mono.doOnError(Consumer).

      Cancellation while MessageSource.receive() is in flight (the sink is already CANCELLED) is handled by Reactor's discard hook on this flux. Only a PollableChannel adapted through messageChannelToFlux(MessageChannel) is best-effort re-queued via non-blocking MessageChannel.send(Message, long) (interceptor chain is re-entered; FIFO is not preserved if the queue is not empty; a full bounded channel is logged and the message is not nack'd). A direct messageSourceToFlux(() -> channel.receive(0)) lambda is not re-queued. Other sources are left unacknowledged so the source's own redelivery (broker requeue, uncommitted offset) can recover them. The discard hook only sees drops at or upstream of this flux; operators a caller adds afterward (ReactiveStreamsConsumer.setReactiveCustomizer, a flatMap on the ReactiveMessageHandler path) use a different context and are not rescued here.

      Type Parameters:
      T - the expected payload type.
      Parameters:
      messageSource - the MessageSource to adapt.
      Returns:
      a Flux which pulls messages from the MessageSource on demand.
    • messageChannelToFlux

      public static <T> reactor.core.publisher.Flux<Message<T>> messageChannelToFlux(MessageChannel messageChannel)
      Adapt a provided MessageChannel into a Flux source: - a FluxMessageChannel is returned as is because it is already a Publisher; - a SubscribableChannel is subscribed with a MessageHandler for the Sinks.Many.tryEmitNext(Object) which is returned from this method; - a PollableChannel is wrapped into a IntegrationReactiveUtils.PollableChannelMessageSource so cancelled in-flight receives can be re-queued, and reuses messageSourceToFlux(MessageSource).
      Type Parameters:
      T - the expected payload type.
      Parameters:
      messageChannel - the MessageChannel to adapt.
      Returns:
      a Flux which uses a provided MessageChannel as a source for events to publish.