Package rx.internal.operators
Class OperatorEagerConcatMap.EagerOuterSubscriber<T,R>
- java.lang.Object
-
- rx.Subscriber<T>
-
- rx.internal.operators.OperatorEagerConcatMap.EagerOuterSubscriber<T,R>
-
- All Implemented Interfaces:
Observer<T>
,Subscription
- Enclosing class:
- OperatorEagerConcatMap<T,R>
static final class OperatorEagerConcatMap.EagerOuterSubscriber<T,R> extends Subscriber<T>
-
-
Field Summary
Fields Modifier and Type Field Description (package private) Subscriber<? super R>
actual
(package private) int
bufferSize
(package private) boolean
cancelled
(package private) boolean
done
(package private) java.lang.Throwable
error
(package private) Func1<? super T,? extends Observable<? extends R>>
mapper
private OperatorEagerConcatMap.EagerOuterProducer
sharedProducer
(package private) java.util.Queue<OperatorEagerConcatMap.EagerInnerSubscriber<R>>
subscribers
(package private) java.util.concurrent.atomic.AtomicInteger
wip
-
Constructor Summary
Constructors Constructor Description EagerOuterSubscriber(Func1<? super T,? extends Observable<? extends R>> mapper, int bufferSize, int maxConcurrent, Subscriber<? super R> actual)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description (package private) void
cleanup()
(package private) void
drain()
(package private) void
init()
void
onCompleted()
Notifies the Observer that theObservable
has finished sending push-based notifications.void
onError(java.lang.Throwable e)
Notifies the Observer that theObservable
has experienced an error condition.void
onNext(T t)
Provides the Observer with a new item to observe.-
Methods inherited from class rx.Subscriber
add, isUnsubscribed, onStart, request, setProducer, unsubscribe
-
-
-
-
Field Detail
-
mapper
final Func1<? super T,? extends Observable<? extends R>> mapper
-
bufferSize
final int bufferSize
-
actual
final Subscriber<? super R> actual
-
subscribers
final java.util.Queue<OperatorEagerConcatMap.EagerInnerSubscriber<R>> subscribers
-
done
volatile boolean done
-
error
java.lang.Throwable error
-
cancelled
volatile boolean cancelled
-
wip
final java.util.concurrent.atomic.AtomicInteger wip
-
sharedProducer
private OperatorEagerConcatMap.EagerOuterProducer sharedProducer
-
-
Constructor Detail
-
EagerOuterSubscriber
public EagerOuterSubscriber(Func1<? super T,? extends Observable<? extends R>> mapper, int bufferSize, int maxConcurrent, Subscriber<? super R> actual)
-
-
Method Detail
-
init
void init()
-
cleanup
void cleanup()
-
onNext
public void onNext(T t)
Description copied from interface:Observer
Provides the Observer with a new item to observe.The
Observable
may call this method 0 or more times.The
Observable
will not call this method again after it calls eitherObserver.onCompleted()
orObserver.onError(java.lang.Throwable)
.- Parameters:
t
- the item emitted by the Observable
-
onError
public void onError(java.lang.Throwable e)
Description copied from interface:Observer
Notifies the Observer that theObservable
has experienced an error condition.If the
Observable
calls this method, it will not thereafter callObserver.onNext(T)
orObserver.onCompleted()
.- Parameters:
e
- the exception encountered by the Observable
-
onCompleted
public void onCompleted()
Description copied from interface:Observer
Notifies the Observer that theObservable
has finished sending push-based notifications.The
Observable
will not call this method if it callsObserver.onError(java.lang.Throwable)
.
-
drain
void drain()
-
-