Class MaybeFlattenStreamAsFlowable.FlattenStreamMultiObserver<T,​R>

    • Constructor Summary

      Constructors 
      Constructor Description
      FlattenStreamMultiObserver​(org.reactivestreams.Subscriber<? super R> downstream, Function<? super T,​? extends java.util.stream.Stream<? extends R>> mapper)  
    • Method Summary

      All Methods Instance Methods Concrete Methods 
      Modifier and Type Method Description
      void cancel()  
      void clear()
      Removes all enqueued items from this queue.
      (package private) void close​(java.lang.AutoCloseable c)  
      (package private) void drain()  
      boolean isEmpty()
      Returns true if the queue is empty.
      void onComplete()
      Called once the deferred computation completes normally.
      void onError​(@NonNull java.lang.Throwable e)
      Notifies the MaybeObserver that the Maybe has experienced an error condition.
      void onSubscribe​(@NonNull Disposable d)
      Provides the MaybeObserver with the means of cancelling (disposing) the connection (channel) with the Maybe in both synchronous (from within onSubscribe(Disposable) itself) and asynchronous manner.
      void onSuccess​(T t)
      Notifies the MaybeObserver with one item and that the Maybe has finished sending push-based notifications.
      R poll()
      Tries to dequeue a value (non-null) or returns null if the queue is empty.
      void request​(long n)  
      int requestFusion​(int mode)
      Request a fusion mode from the upstream.
      • Methods inherited from class java.util.concurrent.atomic.AtomicInteger

        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, toString, updateAndGet, weakCompareAndSet, weakCompareAndSetAcquire, weakCompareAndSetPlain, weakCompareAndSetRelease, weakCompareAndSetVolatile
      • Methods inherited from class java.lang.Number

        byteValue, shortValue
      • Methods inherited from class java.lang.Object

        clone, equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, wait
    • Field Detail

      • downstream

        final org.reactivestreams.Subscriber<? super R> downstream
      • mapper

        final Function<? super T,​? extends java.util.stream.Stream<? extends R>> mapper
      • requested

        final java.util.concurrent.atomic.AtomicLong requested
      • iterator

        volatile java.util.Iterator<? extends R> iterator
      • close

        java.lang.AutoCloseable close
      • once

        boolean once
      • cancelled

        volatile boolean cancelled
      • outputFused

        boolean outputFused
      • emitted

        long emitted
    • Constructor Detail

      • FlattenStreamMultiObserver

        FlattenStreamMultiObserver​(org.reactivestreams.Subscriber<? super R> downstream,
                                   Function<? super T,​? extends java.util.stream.Stream<? extends R>> mapper)
    • Method Detail

      • onComplete

        public void onComplete()
        Description copied from interface: MaybeObserver
        Called once the deferred computation completes normally.
        Specified by:
        onComplete in interface MaybeObserver<T>
      • request

        public void request​(long n)
        Specified by:
        request in interface org.reactivestreams.Subscription
      • cancel

        public void cancel()
        Specified by:
        cancel in interface org.reactivestreams.Subscription
      • poll

        @Nullable
        public R poll()
               throws java.lang.Throwable
        Description copied from interface: SimpleQueue
        Tries to dequeue a value (non-null) or returns null if the queue is empty.

        If the producer uses SimpleQueue.offer(Object, Object) and when polling in pairs, if the first poll() returns a non-null item, the second poll() is guaranteed to return a non-null item as well.

        Specified by:
        poll in interface SimpleQueue<T>
        Returns:
        the item or null to indicate an empty queue
        Throws:
        java.lang.Throwable - if some pre-processing of the dequeued item (usually through fused functions) throws.
      • isEmpty

        public boolean isEmpty()
        Description copied from interface: SimpleQueue
        Returns true if the queue is empty.

        Note however that due to potential fused functions in SimpleQueue.poll() it is possible this method returns false but then poll() returns null because the fused function swallowed the available item(s).

        Specified by:
        isEmpty in interface SimpleQueue<T>
        Returns:
        true if the queue is empty
      • clear

        public void clear()
        Description copied from interface: SimpleQueue
        Removes all enqueued items from this queue.
        Specified by:
        clear in interface SimpleQueue<T>
      • close

        void close​(java.lang.AutoCloseable c)
      • drain

        void drain()