Class SafeSubscriber<T>

  • Type Parameters:
    T - the value type
    All Implemented Interfaces:
    FlowableSubscriber<T>, org.reactivestreams.Subscriber<T>, org.reactivestreams.Subscription

    public final class SafeSubscriber<@NonNull T>
    extends java.lang.Object
    implements FlowableSubscriber<T>, org.reactivestreams.Subscription
    Wraps another Subscriber and ensures all onXXX methods conform the protocol (except the requirement for serialized access).
    • Field Summary

      Fields 
      Modifier and Type Field Description
      (package private) boolean done
      Indicates a terminal state.
      (package private) org.reactivestreams.Subscriber<? super T> downstream
      The actual Subscriber.
      (package private) org.reactivestreams.Subscription upstream
      The subscription.
    • Constructor Summary

      Constructors 
      Constructor Description
      SafeSubscriber​(@NonNull org.reactivestreams.Subscriber<? super @NonNull T> downstream)
      Constructs a SafeSubscriber by wrapping the given actual Subscriber.
    • Field Detail

      • downstream

        final org.reactivestreams.Subscriber<? super T> downstream
        The actual Subscriber.
      • upstream

        org.reactivestreams.Subscription upstream
        The subscription.
      • done

        boolean done
        Indicates a terminal state.
    • Constructor Detail

      • SafeSubscriber

        public SafeSubscriber​(@NonNull
                              @NonNull org.reactivestreams.Subscriber<? super @NonNull T> downstream)
        Constructs a SafeSubscriber by wrapping the given actual Subscriber.
        Parameters:
        downstream - the actual Subscriber to wrap, not null (not validated)
    • Method Detail

      • onSubscribe

        public void onSubscribe​(@NonNull
                                @NonNull org.reactivestreams.Subscription s)
        Description copied from interface: FlowableSubscriber
        Implementors of this method should make sure everything that needs to be visible in Subscriber.onNext(Object) is established before calling Subscription.request(long). In practice this means no initialization should happen after the request() call and additional behavior is thread safe in respect to onNext.
        Specified by:
        onSubscribe in interface FlowableSubscriber<T>
        Specified by:
        onSubscribe in interface org.reactivestreams.Subscriber<T>
      • onNext

        public void onNext​(@NonNull
                           @NonNull T t)
        Specified by:
        onNext in interface org.reactivestreams.Subscriber<T>
      • onNextNoSubscription

        void onNextNoSubscription()
      • onError

        public void onError​(@NonNull
                            @NonNull java.lang.Throwable t)
        Specified by:
        onError in interface org.reactivestreams.Subscriber<T>
      • onComplete

        public void onComplete()
        Specified by:
        onComplete in interface org.reactivestreams.Subscriber<T>
      • onCompleteNoSubscription

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