Class FlowableCreate.BaseEmitter<T>
java.lang.Object
java.lang.Number
java.util.concurrent.atomic.AtomicLong
io.reactivex.rxjava3.internal.operators.flowable.FlowableCreate.BaseEmitter<T>
- All Implemented Interfaces:
Emitter<T>
,FlowableEmitter<T>
,Serializable
,org.reactivestreams.Subscription
- Direct Known Subclasses:
FlowableCreate.BufferAsyncEmitter
,FlowableCreate.LatestAsyncEmitter
,FlowableCreate.MissingEmitter
,FlowableCreate.NoOverflowBaseAsyncEmitter
- Enclosing class:
FlowableCreate<T>
abstract static class FlowableCreate.BaseEmitter<T>
extends AtomicLong
implements FlowableEmitter<T>, org.reactivestreams.Subscription
-
Field Summary
FieldsModifier and TypeFieldDescription(package private) final org.reactivestreams.Subscriber
<? super T> (package private) final SequentialDisposable
private static final long
-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionfinal void
cancel()
protected void
protected boolean
final boolean
Returns true if the downstream cancelled the sequence or the emitter was terminated viaEmitter.onError(Throwable)
,Emitter.onComplete()
or a successfulFlowableEmitter.tryOnError(Throwable)
.void
Signal a completion.final void
Signal aThrowable
exception.(package private) void
(package private) void
final void
request
(long n) final long
The current outstanding request amount.final FlowableEmitter
<T> Ensures that calls toonNext
,onError
andonComplete
are properly serialized.final void
Sets aCancellable
on this emitter; any previousDisposable
orCancellable
will be disposed/cancelled.final void
Sets a Disposable on this emitter; any previousDisposable
orCancellable
will be disposed/cancelled.boolean
toString()
final boolean
Attempts to emit the specifiedThrowable
error if the downstream hasn't cancelled the sequence or is otherwise terminated, returning false if the emission is not allowed to happen due to lifecycle restrictions.Methods inherited from class java.util.concurrent.atomic.AtomicLong
accumulateAndGet, addAndGet, compareAndExchange, compareAndExchangeAcquire, compareAndExchangeRelease, compareAndSet, decrementAndGet, doubleValue, floatValue, get, getAcquire, getAndAccumulate, getAndAdd, getAndDecrement, getAndIncrement, getAndSet, getAndUpdate, getOpaque, getPlain, incrementAndGet, intValue, lazySet, longValue, set, setOpaque, setPlain, setRelease, updateAndGet, weakCompareAndSet, weakCompareAndSetAcquire, weakCompareAndSetPlain, weakCompareAndSetRelease, weakCompareAndSetVolatile
Methods inherited from class java.lang.Number
byteValue, shortValue
-
Field Details
-
serialVersionUID
private static final long serialVersionUID- See Also:
-
downstream
-
serial
-
-
Constructor Details
-
BaseEmitter
BaseEmitter(org.reactivestreams.Subscriber<? super T> downstream)
-
-
Method Details
-
onComplete
public void onComplete()Description copied from interface:Emitter
Signal a completion.- Specified by:
onComplete
in interfaceEmitter<T>
-
completeDownstream
protected void completeDownstream() -
onError
Description copied from interface:Emitter
Signal aThrowable
exception. -
tryOnError
Description copied from interface:FlowableEmitter
Attempts to emit the specifiedThrowable
error if the downstream hasn't cancelled the sequence or is otherwise terminated, returning false if the emission is not allowed to happen due to lifecycle restrictions.Unlike
Emitter.onError(Throwable)
, theRxjavaPlugins.onError
is not called if the error could not be delivered.History: 2.1.1 - experimental
- Specified by:
tryOnError
in interfaceFlowableEmitter<T>
- Parameters:
e
- the throwable error to signal if possible- Returns:
- true if successful, false if the downstream is not able to accept further events
-
signalError
-
errorDownstream
-
cancel
public final void cancel()- Specified by:
cancel
in interfaceorg.reactivestreams.Subscription
-
onUnsubscribed
void onUnsubscribed() -
isCancelled
public final boolean isCancelled()Description copied from interface:FlowableEmitter
Returns true if the downstream cancelled the sequence or the emitter was terminated viaEmitter.onError(Throwable)
,Emitter.onComplete()
or a successfulFlowableEmitter.tryOnError(Throwable)
.This method is thread-safe.
- Specified by:
isCancelled
in interfaceFlowableEmitter<T>
- Returns:
- true if the downstream cancelled the sequence or the emitter was terminated
-
request
public final void request(long n) - Specified by:
request
in interfaceorg.reactivestreams.Subscription
-
onRequested
void onRequested() -
setDisposable
Description copied from interface:FlowableEmitter
Sets a Disposable on this emitter; any previousDisposable
orCancellable
will be disposed/cancelled.This method is thread-safe.
- Specified by:
setDisposable
in interfaceFlowableEmitter<T>
- Parameters:
d
- the disposable,null
is allowed
-
setCancellable
Description copied from interface:FlowableEmitter
Sets aCancellable
on this emitter; any previousDisposable
orCancellable
will be disposed/cancelled.This method is thread-safe.
- Specified by:
setCancellable
in interfaceFlowableEmitter<T>
- Parameters:
c
- theCancellable
resource,null
is allowed
-
requested
public final long requested()Description copied from interface:FlowableEmitter
The current outstanding request amount.This method is thread-safe.
- Specified by:
requested
in interfaceFlowableEmitter<T>
- Returns:
- the current outstanding request amount
-
serialize
Description copied from interface:FlowableEmitter
Ensures that calls toonNext
,onError
andonComplete
are properly serialized.- Specified by:
serialize
in interfaceFlowableEmitter<T>
- Returns:
- the serialized
FlowableEmitter
-
toString
- Overrides:
toString
in classAtomicLong
-