Interface Stream<V>

  • Type Parameters:
    V - Type of data emitted
    All Superinterfaces:
    org.reactivestreams.Publisher<V>

    public interface Stream<V>
    extends org.reactivestreams.Publisher<V>
    Stream is a reactive-streams compliant way to chain operations and transformation on a stream of data. It extends Publisher by providing additional reactive methods allowing to act on it. The goal of this interface is to decouple ourself from reactive framework (RxJava/Reactor).
    • Method Summary

      All Methods Instance Methods Abstract Methods 
      Modifier and Type Method Description
      <O> Stream<O> flatMap​(org.forgerock.util.Function<? super V,​? extends org.reactivestreams.Publisher<? extends O>,​Exception> function, int maxConcurrency)
      Transforms each data emitted by this stream into a new stream of data.
      <O> Stream<O> map​(org.forgerock.util.Function<V,​O,​Exception> function)
      Transforms the data emitted by this stream.
      Stream<V> onComplete​(Action onComplete)
      Invokes the on complete Action when this stream is completed.
      Stream<V> onError​(Consumer<Throwable> onError)
      Invokes the on error Consumer when an error occurs on this stream.
      Stream<V> onErrorResumeWith​(org.forgerock.util.Function<Throwable,​org.reactivestreams.Publisher<V>,​Exception> function)
      When an error occurs in this stream, continue the processing with the new Stream provided by the function.
      Stream<V> onNext​(Consumer<V> onNext)
      Invokes the on next Consumer when this stream emits a value.
      void subscribe()
      Subscribes to this stream and drop all data produced by it.
      • Methods inherited from interface org.reactivestreams.Publisher

        subscribe
    • Method Detail

      • map

        <O> Stream<O> map​(org.forgerock.util.Function<V,​O,​Exception> function)
        Transforms the data emitted by this stream.
        Type Parameters:
        O - Type of data emitted after transformation
        Parameters:
        function - The function to apply to each data of this stream
        Returns:
        a new Stream emitting transformed data
      • flatMap

        <O> Stream<O> flatMap​(org.forgerock.util.Function<? super V,​? extends org.reactivestreams.Publisher<? extends O>,​Exception> function,
                              int maxConcurrency)
        Transforms each data emitted by this stream into a new stream of data. All these streams are then merged together.
        Type Parameters:
        O - Type of data emitted after transformation
        Parameters:
        function - The function to transform each data into a new stream
        maxConcurrency - Maximum number of output stream which can be merged. Once this number is reached, the Stream will stop requesting data from this stream.
        Returns:
        A new Stream performing the transformation
      • onErrorResumeWith

        Stream<V> onErrorResumeWith​(org.forgerock.util.Function<Throwable,​org.reactivestreams.Publisher<V>,​Exception> function)
        When an error occurs in this stream, continue the processing with the new Stream provided by the function.
        Parameters:
        function - Generates the stream which must will used to resume operation when this Stream failed.
        Returns:
        A new Stream
      • onNext

        Stream<V> onNext​(Consumer<V> onNext)
        Invokes the on next Consumer when this stream emits a value.
        Parameters:
        onNext - The Consumer to invoke when a value is emitted by this stream
        Returns:
        a new Stream
      • onComplete

        Stream<V> onComplete​(Action onComplete)
        Invokes the on complete Action when this stream is completed.
        Parameters:
        onComplete - The Action to invoke on stream completion
        Returns:
        a new Stream
      • subscribe

        void subscribe()
        Subscribes to this stream and drop all data produced by it.