Package com.forgerock.reactive
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 extendsPublisherby 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 completeActionwhen this stream is completed.Stream<V>onError(Consumer<Throwable> onError)Invokes the on errorConsumerwhen 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 newStreamprovided by the function.Stream<V>onNext(Consumer<V> onNext)Invokes the on nextConsumerwhen this stream emits a value.voidsubscribe()Subscribes to this stream and drop all data produced by it.
-
-
-
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
Streamemitting 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 streammaxConcurrency- Maximum number of output stream which can be merged. Once this number is reached, theStreamwill stop requesting data from this stream.- Returns:
- A new
Streamperforming 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 newStreamprovided by the function.
-
onNext
Stream<V> onNext(Consumer<V> onNext)
Invokes the on nextConsumerwhen this stream emits a value.
-
onError
Stream<V> onError(Consumer<Throwable> onError)
Invokes the on errorConsumerwhen an error occurs on this stream.
-
onComplete
Stream<V> onComplete(Action onComplete)
Invokes the on completeActionwhen this stream is completed.
-
subscribe
void subscribe()
Subscribes to this stream and drop all data produced by it.
-
-