ConcurrentManySink

class ConcurrentManySink<T : Any>(val delegate: Sinks.Many<T>) : Sinks.Many<T> , Decorator<Sinks.Many<T>>

Concurrent-emission Sinks.Many decorator.

Reactor sinks reject concurrent emissions with Sinks.EmitResult.FAIL_NON_SERIALIZED. A ReentrantLock serializes every emission call on its originating thread, preserving synchronous results, exception identity, thread-local context and downstream reentrancy. Subscription, scan and subscriber-count operations are delegated without acquiring this lock.

Parameters

delegate

sink whose emissions are serialized

Type Parameters

T

non-null element type

Constructors

Link copied to clipboard
constructor(delegate: Sinks.Many<T>)

Properties

Link copied to clipboard
val Scannable.cancelled: Boolean
Link copied to clipboard
open override val delegate: Sinks.Many<T>
Link copied to clipboard
Link copied to clipboard
val Scannable.terminated: Boolean

Functions

Link copied to clipboard
open fun actuals(): Stream<out Scannable?>?
Link copied to clipboard
open override fun asFlux(): Flux<T>
Link copied to clipboard
fun <T : Any> Sinks.Many<T>.concurrent(): ConcurrentManySink<T>

Makes every emission method of this sink safe for concurrent producers.

Link copied to clipboard
open override fun currentSubscriberCount(): Int
Link copied to clipboard
open override fun emitComplete(failureHandler: Sinks.EmitFailureHandler)
Link copied to clipboard
open override fun emitError(error: Throwable, failureHandler: Sinks.EmitFailureHandler)
Link copied to clipboard
open override fun emitNext(t: T, failureHandler: Sinks.EmitFailureHandler)
Link copied to clipboard
open fun inners(): Stream<out Scannable?>?
Link copied to clipboard
open fun name(): String?
Link copied to clipboard
open fun parents(): Stream<out Scannable?>?
Link copied to clipboard
open fun <T : Any?> scan(key: Scannable.Attr<T?>?): @Nullable T?
Link copied to clipboard
open fun <T : Any?> scanOrDefault(key: Scannable.Attr<T?>?, defaultValue: T?): T?
Link copied to clipboard
open override fun scanUnsafe(key: Scannable.Attr<*>): Any?
Link copied to clipboard
open fun stepName(): String?
Link copied to clipboard
open fun steps(): Stream<String?>?
Link copied to clipboard
open fun tags(): Stream<Tuple2<String?, String?>?>?
Link copied to clipboard
Link copied to clipboard
open override fun tryEmitComplete(): Sinks.EmitResult
Link copied to clipboard
open override fun tryEmitError(error: Throwable): Sinks.EmitResult
Link copied to clipboard
open override fun tryEmitNext(t: T): Sinks.EmitResult