Class FlowableCombineLatest.CombineLatestCoordinator<T,​R>

    • Constructor Summary

      Constructors 
      Constructor Description
      CombineLatestCoordinator​(org.reactivestreams.Subscriber<? super R> actual, Function<? super java.lang.Object[],​? extends R> combiner, int n, int bufferSize, boolean delayErrors)  
    • Method Summary

      All Methods Instance Methods Concrete Methods 
      Modifier and Type Method Description
      void cancel()  
      (package private) void cancelAll()  
      (package private) boolean checkTerminated​(boolean d, boolean empty, org.reactivestreams.Subscriber<?> a, SpscLinkedArrayQueue<?> q)  
      void clear()
      Removes all enqueued items from this queue.
      (package private) void drain()  
      (package private) void drainAsync()  
      (package private) void drainOutput()  
      (package private) void innerComplete​(int index)  
      (package private) void innerError​(int index, java.lang.Throwable e)  
      (package private) void innerValue​(int index, T value)  
      boolean isEmpty()
      Returns true if the queue is empty.
      R poll()
      Tries to dequeue a value (non-null) or returns null if the queue is empty.
      void request​(long n)  
      int requestFusion​(int requestedMode)
      Request a fusion mode from the upstream.
      (package private) void subscribe​(org.reactivestreams.Publisher<? extends T>[] sources, int n)  
      • 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
      • combiner

        final Function<? super java.lang.Object[],​? extends R> combiner
      • latest

        final java.lang.Object[] latest
      • delayErrors

        final boolean delayErrors
      • outputFused

        boolean outputFused
      • nonEmptySources

        int nonEmptySources
      • completedSources

        int completedSources
      • cancelled

        volatile boolean cancelled
      • requested

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

        volatile boolean done
    • Constructor Detail

      • CombineLatestCoordinator

        CombineLatestCoordinator​(org.reactivestreams.Subscriber<? super R> actual,
                                 Function<? super java.lang.Object[],​? extends R> combiner,
                                 int n,
                                 int bufferSize,
                                 boolean delayErrors)
    • Method Detail

      • request

        public void request​(long n)
      • cancel

        public void cancel()
      • subscribe

        void subscribe​(org.reactivestreams.Publisher<? extends T>[] sources,
                       int n)
      • innerValue

        void innerValue​(int index,
                        T value)
      • innerComplete

        void innerComplete​(int index)
      • innerError

        void innerError​(int index,
                        java.lang.Throwable e)
      • drainOutput

        void drainOutput()
      • drainAsync

        void drainAsync()
      • drain

        void drain()
      • checkTerminated

        boolean checkTerminated​(boolean d,
                                boolean empty,
                                org.reactivestreams.Subscriber<?> a,
                                SpscLinkedArrayQueue<?> q)
      • cancelAll

        void cancelAll()
      • 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.

        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.
      • clear

        public void clear()
        Description copied from interface: SimpleQueue
        Removes all enqueued items from this queue.
      • 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).

        Returns:
        true if the queue is empty