Class FutureMultiObserver<T>
- java.lang.Object
-
- java.util.concurrent.CountDownLatch
-
- io.reactivex.rxjava3.internal.observers.FutureMultiObserver<T>
-
- Type Parameters:
T- the value type
- All Implemented Interfaces:
CompletableObserver,MaybeObserver<T>,SingleObserver<T>,Disposable,java.util.concurrent.Future<T>
public final class FutureMultiObserver<T> extends java.util.concurrent.CountDownLatch implements MaybeObserver<T>, SingleObserver<T>, CompletableObserver, java.util.concurrent.Future<T>, Disposable
An Observer + Future that expects exactly one upstream value and provides it via the (blocking) Future API.
-
-
Field Summary
Fields Modifier and Type Field Description (package private) java.lang.Throwableerror(package private) java.util.concurrent.atomic.AtomicReference<Disposable>upstream(package private) Tvalue
-
Constructor Summary
Constructors Constructor Description FutureMultiObserver()
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description booleancancel(boolean mayInterruptIfRunning)voiddispose()Dispose the resource, the operation should be idempotent.Tget()Tget(long timeout, @NonNull java.util.concurrent.TimeUnit unit)booleanisCancelled()booleanisDisposed()Returns true if this resource has been disposed.booleanisDone()voidonComplete()Called once the deferred computation completes normally.voidonError(java.lang.Throwable t)Notifies theMaybeObserverthat theMaybehas experienced an error condition.voidonSubscribe(Disposable d)Provides theMaybeObserverwith the means of cancelling (disposing) the connection (channel) with theMaybein both synchronous (from withinonSubscribe(Disposable)itself) and asynchronous manner.voidonSuccess(T t)Notifies theMaybeObserverwith one item and that theMaybehas finished sending push-based notifications.
-
-
-
Field Detail
-
value
T value
-
error
java.lang.Throwable error
-
upstream
final java.util.concurrent.atomic.AtomicReference<Disposable> upstream
-
-
Method Detail
-
cancel
public boolean cancel(boolean mayInterruptIfRunning)
- Specified by:
cancelin interfacejava.util.concurrent.Future<T>
-
isCancelled
public boolean isCancelled()
- Specified by:
isCancelledin interfacejava.util.concurrent.Future<T>
-
isDone
public boolean isDone()
- Specified by:
isDonein interfacejava.util.concurrent.Future<T>
-
get
public T get() throws java.lang.InterruptedException, java.util.concurrent.ExecutionException
- Specified by:
getin interfacejava.util.concurrent.Future<T>- Throws:
java.lang.InterruptedExceptionjava.util.concurrent.ExecutionException
-
get
public T get(long timeout, @NonNull @NonNull java.util.concurrent.TimeUnit unit) throws java.lang.InterruptedException, java.util.concurrent.ExecutionException, java.util.concurrent.TimeoutException
- Specified by:
getin interfacejava.util.concurrent.Future<T>- Throws:
java.lang.InterruptedExceptionjava.util.concurrent.ExecutionExceptionjava.util.concurrent.TimeoutException
-
onSubscribe
public void onSubscribe(Disposable d)
Description copied from interface:MaybeObserverProvides theMaybeObserverwith the means of cancelling (disposing) the connection (channel) with theMaybein both synchronous (from withinonSubscribe(Disposable)itself) and asynchronous manner.- Specified by:
onSubscribein interfaceCompletableObserver- Specified by:
onSubscribein interfaceMaybeObserver<T>- Specified by:
onSubscribein interfaceSingleObserver<T>- Parameters:
d- theDisposableinstance whoseDisposable.dispose()can be called anytime to cancel the connection
-
onSuccess
public void onSuccess(T t)
Description copied from interface:MaybeObserverNotifies theMaybeObserverwith one item and that theMaybehas finished sending push-based notifications.The
Maybewill not call this method if it callsMaybeObserver.onError(java.lang.Throwable).- Specified by:
onSuccessin interfaceMaybeObserver<T>- Specified by:
onSuccessin interfaceSingleObserver<T>- Parameters:
t- the item emitted by theMaybe
-
onError
public void onError(java.lang.Throwable t)
Description copied from interface:MaybeObserverNotifies theMaybeObserverthat theMaybehas experienced an error condition.If the
Maybecalls this method, it will not thereafter callMaybeObserver.onSuccess(T).- Specified by:
onErrorin interfaceCompletableObserver- Specified by:
onErrorin interfaceMaybeObserver<T>- Specified by:
onErrorin interfaceSingleObserver<T>- Parameters:
t- the exception encountered by theMaybe
-
onComplete
public void onComplete()
Description copied from interface:MaybeObserverCalled once the deferred computation completes normally.- Specified by:
onCompletein interfaceCompletableObserver- Specified by:
onCompletein interfaceMaybeObserver<T>
-
dispose
public void dispose()
Description copied from interface:DisposableDispose the resource, the operation should be idempotent.- Specified by:
disposein interfaceDisposable
-
isDisposed
public boolean isDisposed()
Description copied from interface:DisposableReturns true if this resource has been disposed.- Specified by:
isDisposedin interfaceDisposable- Returns:
- true if this resource has been disposed
-
-