T - the type of elements held in the nested Fluxespublic interface FluxT<T>
extends com.aol.cyclops.control.monads.transformers.values.FoldableTransformerSeq<T>
| Modifier and Type | Method and Description |
|---|---|
default <B> FluxT<B> |
bind(java.util.function.Function<? super T,FluxT<? extends B>> f)
Flat Map the wrapped Flux
|
default <U> FluxT<U> |
cast(java.lang.Class<? extends U> type) |
default FluxT<T> |
combine(java.util.function.BiPredicate<? super T,? super T> predicate,
java.util.function.BinaryOperator<T> op) |
default FluxT<T> |
cycle(int times) |
default FluxT<T> |
cycle(com.aol.cyclops.Monoid<T> m,
int times) |
default FluxT<T> |
cycleUntil(java.util.function.Predicate<? super T> predicate) |
default FluxT<T> |
cycleWhile(java.util.function.Predicate<? super T> predicate) |
default FluxT<T> |
distinct() |
default FluxT<T> |
dropRight(int num) |
default FluxT<T> |
dropUntil(java.util.function.Predicate<? super T> p) |
default FluxT<T> |
dropWhile(java.util.function.Predicate<? super T> p) |
<R> FluxT<R> |
empty() |
static <T> FluxTSeq<T> |
emptyFlux() |
static <T> FluxTValue<T> |
emptyOptional() |
FluxT<T> |
filter(java.util.function.Predicate<? super T> test)
Filter the wrapped Flux
|
default FluxT<T> |
filterNot(java.util.function.Predicate<? super T> fn) |
<B> FluxT<B> |
flatMap(java.util.function.Function<? super T,? extends reactor.core.publisher.Flux<? extends B>> f)
Perform a flatMap operation on each nested Flux
|
reactor.core.publisher.Flux<T> |
flux() |
default reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
fluxOfFlux()
Convert this FluxTransformer into a Flux of Fluxes
|
static <A> FluxT<A> |
fromAnyM(com.aol.cyclops.control.AnyM<A> anyM)
|
static <A> FluxTSeq<A> |
fromAnyMSeq(com.aol.cyclops.types.anyM.AnyMSeq<A> anyM)
Create a FluxT from an AnyMSeq by wrapping the elements stored in the AnyMSeq in a Flux
|
static <A> FluxTValue<A> |
fromAnyMValue(com.aol.cyclops.types.anyM.AnyMValue<A> anyM)
Create a FluxT from an AnyMValue by wrapping the element stored in the AnyMValue in a Flux
|
static <A> FluxTValue<A> |
fromFuture(java.util.concurrent.CompletableFuture<reactor.core.publisher.Flux<A>> future)
Create a FluxTValue from a JDK CompletableFuture that contains nested Fluxes
|
static <A> FluxTSeq<A> |
fromIterable(java.lang.Iterable<reactor.core.publisher.Flux<A>> iterableOfFluxs)
Create a FluxT from an Iterable that contains nested Fluxes
|
static <A> FluxTValue<A> |
fromIterableValue(java.lang.Iterable<reactor.core.publisher.Flux<A>> iterableOfFluxs)
Create a FluxTValue from an Iterable that contains nested Fluxes
|
static <A> FluxTValue<A> |
fromMono(reactor.core.publisher.Mono<reactor.core.publisher.Flux<A>> mono)
Create a FluxTValue from a Reactor Mono type that contains nested Fluxes
|
static <A> FluxTValue<A> |
fromOptional(java.util.Optional<reactor.core.publisher.Flux<A>> optional)
Create a FluxTValue from a JDK Optional that contains nested Fluxes
|
static <A> FluxTSeq<A> |
fromPublisher(org.reactivestreams.Publisher<reactor.core.publisher.Flux<A>> publisherOfFluxs)
Create a FluxTSeq from a Publisher that contains nested Fluxes
|
static <A,V extends com.aol.cyclops.types.MonadicValue<? extends reactor.core.publisher.Flux<A>>> |
fromValue(V monadicValue)
Create a FluxTValue from a cyclops-react Value that contains nested Fluxes
|
default <K> FluxT<org.jooq.lambda.tuple.Tuple2<K,org.jooq.lambda.Seq<T>>> |
grouped(java.util.function.Function<? super T,? extends K> classifier) |
default <K,A,D> FluxT<org.jooq.lambda.tuple.Tuple2<K,D>> |
grouped(java.util.function.Function<? super T,? extends K> classifier,
java.util.stream.Collector<? super T,A,D> downFlux) |
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> |
grouped(int groupSize) |
default <C extends java.util.Collection<? super T>> |
grouped(int size,
java.util.function.Supplier<C> supplier) |
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> |
groupedStatefullyUntil(java.util.function.BiPredicate<com.aol.cyclops.data.collections.extensions.standard.ListX<? super T>,? super T> predicate) |
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> |
groupedUntil(java.util.function.Predicate<? super T> predicate) |
default <C extends java.util.Collection<? super T>> |
groupedUntil(java.util.function.Predicate<? super T> predicate,
java.util.function.Supplier<C> factory) |
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> |
groupedWhile(java.util.function.Predicate<? super T> predicate) |
default <C extends java.util.Collection<? super T>> |
groupedWhile(java.util.function.Predicate<? super T> predicate,
java.util.function.Supplier<C> factory) |
default FluxT<T> |
intersperse(T value) |
static <U,R> java.util.function.Function<FluxT<U>,FluxT<R>> |
lift(java.util.function.Function<? super U,? extends R> fn)
Lift a function into one that accepts and returns an FluxT
This allows multiple monad types to add functionality to existing functions and methods
e.g.
|
default FluxT<T> |
limit(long num) |
default FluxT<T> |
limitLast(int num) |
default FluxT<T> |
limitUntil(java.util.function.Predicate<? super T> p) |
default FluxT<T> |
limitWhile(java.util.function.Predicate<? super T> p) |
<B> FluxT<B> |
map(java.util.function.Function<? super T,? extends B> f)
Map the wrapped Flux
|
default FluxT<T> |
notNull() |
static <A> FluxT<A> |
of(com.aol.cyclops.control.AnyM<? extends reactor.core.publisher.Flux<A>> monads)
Create a FluxT from an AnyM that wraps a monad containing a Flux
|
default <U> FluxT<U> |
ofType(java.lang.Class<? extends U> type) |
default FluxT<T> |
onEmpty(T value) |
default FluxT<T> |
onEmptyGet(java.util.function.Supplier<? extends T> supplier) |
default <X extends java.lang.Throwable> |
onEmptyThrow(java.util.function.Supplier<? extends X> supplier) |
default <R> FluxT<R> |
patternMatch(java.util.function.Function<com.aol.cyclops.control.Matchable.CheckValue1<T,R>,com.aol.cyclops.control.Matchable.CheckValue1<T,R>> case1,
java.util.function.Supplier<? extends R> otherwise) |
FluxT<T> |
peek(java.util.function.Consumer<? super T> peek)
Peek at the current value of the Flux
|
default FluxT<T> |
reverse() |
default FluxT<T> |
scanLeft(com.aol.cyclops.Monoid<T> monoid) |
default <U> FluxT<U> |
scanLeft(U seed,
java.util.function.BiFunction<? super U,? super T,? extends U> function) |
default FluxT<T> |
scanRight(com.aol.cyclops.Monoid<T> monoid) |
default <U> FluxT<U> |
scanRight(U identity,
java.util.function.BiFunction<? super T,? super U,? extends U> combiner) |
default FluxT<T> |
shuffle() |
default FluxT<T> |
shuffle(java.util.Random random) |
default FluxT<T> |
skip(long num) |
default FluxT<T> |
skipLast(int num) |
default FluxT<T> |
skipUntil(java.util.function.Predicate<? super T> p) |
default FluxT<T> |
skipWhile(java.util.function.Predicate<? super T> p) |
default FluxT<T> |
slice(long from,
long to) |
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> |
sliding(int windowSize) |
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> |
sliding(int windowSize,
int increment) |
default FluxT<T> |
sorted() |
default FluxT<T> |
sorted(java.util.Comparator<? super T> c) |
default <U extends java.lang.Comparable<? super U>> |
sorted(java.util.function.Function<? super T,? extends U> function) |
default FluxT<T> |
takeRight(int num) |
default FluxT<T> |
takeUntil(java.util.function.Predicate<? super T> p) |
default FluxT<T> |
takeWhile(java.util.function.Predicate<? super T> p) |
default <R> FluxT<R> |
trampoline(java.util.function.Function<? super T,? extends com.aol.cyclops.control.Trampoline<? extends R>> mapper) |
<R> FluxT<R> |
unit(R t)
Create an instance of the same FluxTransformer that contains a Flux with just the value provided
|
<R> FluxT<R> |
unitIterator(java.util.Iterator<R> it)
Create an instance of the same FluxTransformer type from the provided Iterator over raw values.
|
com.aol.cyclops.control.AnyM<reactor.core.publisher.Flux<T>> |
unwrap() |
default <U> FluxT<org.jooq.lambda.tuple.Tuple2<T,U>> |
zip(java.lang.Iterable<? extends U> other) |
default <U,R> FluxT<R> |
zip(java.lang.Iterable<? extends U> other,
java.util.function.BiFunction<? super T,? super U,? extends R> zipper) |
default <U> FluxT<org.jooq.lambda.tuple.Tuple2<T,U>> |
zip(org.jooq.lambda.Seq<? extends U> other) |
default <U,R> FluxT<R> |
zip(org.jooq.lambda.Seq<? extends U> other,
java.util.function.BiFunction<? super T,? super U,? extends R> zipper) |
default <U> FluxT<org.jooq.lambda.tuple.Tuple2<T,U>> |
zip(java.util.stream.Stream<? extends U> other) |
default <U,R> FluxT<R> |
zip(java.util.stream.Stream<? extends U> other,
java.util.function.BiFunction<? super T,? super U,? extends R> zipper) |
default <S,U> FluxT<org.jooq.lambda.tuple.Tuple3<T,S,U>> |
zip3(java.util.stream.Stream<? extends S> second,
java.util.stream.Stream<? extends U> third) |
default <T2,T3,T4> FluxT<org.jooq.lambda.tuple.Tuple4<T,T2,T3,T4>> |
zip4(java.util.stream.Stream<? extends T2> second,
java.util.stream.Stream<? extends T3> third,
java.util.stream.Stream<? extends T4> fourth) |
default FluxT<org.jooq.lambda.tuple.Tuple2<T,java.lang.Long>> |
zipWithIndex() |
streamfutureOperations, isSeqPresent, lazyOperations, subscribe, transformerStream, unitAnyMseq, toCompletableFuture, toDequeX, toEvalAlways, toEvalLater, toEvalNow, toFutureStream, toFutureStream, toFutureW, toIor, toIorSecondary, toListX, toMapX, toMaybe, toOptional, toPBagX, toPMapX, toPOrderedSetX, toPQueueX, toPSetX, toPStackX, toPVectorX, toQueueX, toSetX, toSimpleReact, toSimpleReact, toSortedSetX, toStreamable, toTry, toValue, toValueMap, toValueSet, toXor, toXorSecondaryfutureStream, getStreamable, isEmpty, iterator, jdkStream, reactiveSeq, reveresedJDKStream, reveresedStreamendsWith, endsWithIterable, findFirst, firstValue, foldRight, foldRight, foldRight, foldRightMapToType, get, groupBy, headAndTail, join, join, join, mapReduce, mapReduce, nestedFoldables, print, print, printErr, printOut, reduce, reduce, reduce, reduce, reduce, reduce, schedule, scheduleFixedDelay, scheduleFixedRate, single, single, singleOptional, startsWith, startsWithIterable, toConcurrentLazyCollection, toConcurrentLazyStreamable, toLazyCollection, validate, visit<R> FluxT<R> unitIterator(java.util.Iterator<R> it)
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.just(1,2,3),Flux.just(4,5,6));
FluxT<String> fluxTStrings = fluxT.unitIterator(Arrays.asList("hello","world").iterator());
//List[Flux["hello","world"]]
it - Iterator over raw values to add to a new Flux Transformer<R> FluxT<R> unit(R t)
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.just(1,2,3),Flux.just(4,5,6));
FluxT<String> fluxTStrings = fluxT.unit("hello");
//List[Flux["hello"]]
t - Value to embed in new Flux Transformer<R> FluxT<R> empty()
<B> FluxT<B> flatMap(java.util.function.Function<? super T,? extends reactor.core.publisher.Flux<? extends B>> f)
f - FlatMapping functiondefault reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> fluxOfFlux()
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.just(1,2,3),Flux.just(4,5,6));
Flux<Flux<Integer>> fluxes = fluxT.fluxOfFlux();
//Flux[Flux[1,2,3],Flux[4,5,6]]
com.aol.cyclops.control.AnyM<reactor.core.publisher.Flux<T>> unwrap()
FluxT<T> peek(java.util.function.Consumer<? super T> peek)
FluxT.of(AnyM.fromFlux(Arrays.asFlux(10))
.peek(System.out::println);
//prints 10
peek in interface com.aol.cyclops.types.Functor<T>peek - Consumer to accept current value of FluxFluxT<T> filter(java.util.function.Predicate<? super T> test)
FluxT.of(AnyM.fromFlux(Arrays.asFlux(10,11))
.filter(t->t!=10);
//FluxT<AnyM<Flux<Flux[11]>>>
<B> FluxT<B> map(java.util.function.Function<? super T,? extends B> f)
FluxT.of(AnyM.fromFlux(Arrays.asFlux(10))
.map(t->t=t+1);
//FluxT<AnyM<Flux<Flux[11]>>>
default <B> FluxT<B> bind(java.util.function.Function<? super T,FluxT<? extends B>> f)
FluxT.of(AnyM.fromFlux(Flux.just(10))
.flatMap(t->Flux.empty());
//FluxT<AnyM<Flux<Flux.empty>>>
f - FlatMap functionstatic <U,R> java.util.function.Function<FluxT<U>,FluxT<R>> lift(java.util.function.Function<? super U,? extends R> fn)
fn - Function to enhance with functionality from Flux and another monad typestatic <A> FluxT<A> fromAnyM(com.aol.cyclops.control.AnyM<A> anyM)
anyM - AnyM that doesn't contain a monad wrapping an Fluxstatic <A> FluxT<A> of(com.aol.cyclops.control.AnyM<? extends reactor.core.publisher.Flux<A>> monads)
FluxT<Integer> fluxT = FluxT.of(AnyM.fromOptional(Optional.of(Flux.just(1,2,3)));
monads - AnyM containing nested Fluxesstatic <A> FluxTValue<A> fromAnyMValue(com.aol.cyclops.types.anyM.AnyMValue<A> anyM)
anyM - Monad to embed a Flux inside (wrapping it's current value)static <A> FluxTSeq<A> fromAnyMSeq(com.aol.cyclops.types.anyM.AnyMSeq<A> anyM)
anyM - Monad to embed a Flux inside (wrapping it's current values individually in Fluxes)static <A> FluxTSeq<A> fromIterable(java.lang.Iterable<reactor.core.publisher.Flux<A>> iterableOfFluxs)
FluxTSeq<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.just(1,2,3));
iterableOfFluxs - An Iterable containing nested Fluxesstatic <A> FluxTSeq<A> fromPublisher(org.reactivestreams.Publisher<reactor.core.publisher.Flux<A>> publisherOfFluxs)
FluxTSeq<Integer> fluxT = FluxT.fromPublisher(Flux.just(Flux.just(1,2,3));
publisherOfFluxs - A Publisher containing nested Fluxesstatic <A,V extends com.aol.cyclops.types.MonadicValue<? extends reactor.core.publisher.Flux<A>>> FluxTValue<A> fromValue(V monadicValue)
FluxTValue<Integer> fluxT = FluxT.fromValue(Maybe.just(Flux.just(1,2,3));
monadicValue - A Value containing nested Fluxesstatic <A> FluxTValue<A> fromOptional(java.util.Optional<reactor.core.publisher.Flux<A>> optional)
FluxTValue<Integer> fluxT = FluxT.fromOptional(Optional.of(Flux.just(1,2,3));
optional - An Optional containing nested Fluxesstatic <A> FluxTValue<A> fromFuture(java.util.concurrent.CompletableFuture<reactor.core.publisher.Flux<A>> future)
FluxTValue<Integer> fluxT = FluxT.fromFuture(CompletableFuture.completedFuture(Flux.just(1,2,3));
future - A CompletableFuture containing nested Fluxesstatic <A> FluxTValue<A> fromMono(reactor.core.publisher.Mono<reactor.core.publisher.Flux<A>> mono)
FluxTValue<Integer> fluxT = FluxT.fromMono(Mono.just(Flux.just(1,2,3));
future - A Mono containing nested Fluxesstatic <A> FluxTValue<A> fromIterableValue(java.lang.Iterable<reactor.core.publisher.Flux<A>> iterableOfFluxs)
FluxTValue<Integer> fluxT = FluxT.fromIterableValue(Arrays.asList(Flux.just(1,2,3));
iterableOfFluxs - An Iterable containing nested Fluxesstatic <T> FluxTValue<T> emptyOptional()
static <T> FluxTSeq<T> emptyFlux()
reactor.core.publisher.Flux<T> flux()
default <U> FluxT<U> cast(java.lang.Class<? extends U> type)
cast in interface com.aol.cyclops.types.Functor<T>default <R> FluxT<R> trampoline(java.util.function.Function<? super T,? extends com.aol.cyclops.control.Trampoline<? extends R>> mapper)
trampoline in interface com.aol.cyclops.types.Functor<T>default <R> FluxT<R> patternMatch(java.util.function.Function<com.aol.cyclops.control.Matchable.CheckValue1<T,R>,com.aol.cyclops.control.Matchable.CheckValue1<T,R>> case1, java.util.function.Supplier<? extends R> otherwise)
patternMatch in interface com.aol.cyclops.types.Functor<T>default <U> FluxT<U> ofType(java.lang.Class<? extends U> type)
ofType in interface com.aol.cyclops.types.Filterable<T>default FluxT<T> filterNot(java.util.function.Predicate<? super T> fn)
filterNot in interface com.aol.cyclops.types.Filterable<T>default FluxT<T> notNull()
notNull in interface com.aol.cyclops.types.Filterable<T>default FluxT<T> combine(java.util.function.BiPredicate<? super T,? super T> predicate, java.util.function.BinaryOperator<T> op)
default <U,R> FluxT<R> zip(java.lang.Iterable<? extends U> other, java.util.function.BiFunction<? super T,? super U,? extends R> zipper)
default <U,R> FluxT<R> zip(org.jooq.lambda.Seq<? extends U> other, java.util.function.BiFunction<? super T,? super U,? extends R> zipper)
default <U,R> FluxT<R> zip(java.util.stream.Stream<? extends U> other, java.util.function.BiFunction<? super T,? super U,? extends R> zipper)
default <U> FluxT<org.jooq.lambda.tuple.Tuple2<T,U>> zip(java.util.stream.Stream<? extends U> other)
default <U> FluxT<org.jooq.lambda.tuple.Tuple2<T,U>> zip(org.jooq.lambda.Seq<? extends U> other)
default <S,U> FluxT<org.jooq.lambda.tuple.Tuple3<T,S,U>> zip3(java.util.stream.Stream<? extends S> second, java.util.stream.Stream<? extends U> third)
default <T2,T3,T4> FluxT<org.jooq.lambda.tuple.Tuple4<T,T2,T3,T4>> zip4(java.util.stream.Stream<? extends T2> second, java.util.stream.Stream<? extends T3> third, java.util.stream.Stream<? extends T4> fourth)
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> sliding(int windowSize)
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> sliding(int windowSize, int increment)
default <C extends java.util.Collection<? super T>> FluxT<C> grouped(int size, java.util.function.Supplier<C> supplier)
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> groupedUntil(java.util.function.Predicate<? super T> predicate)
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> groupedStatefullyUntil(java.util.function.BiPredicate<com.aol.cyclops.data.collections.extensions.standard.ListX<? super T>,? super T> predicate)
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> groupedWhile(java.util.function.Predicate<? super T> predicate)
default <C extends java.util.Collection<? super T>> FluxT<C> groupedWhile(java.util.function.Predicate<? super T> predicate, java.util.function.Supplier<C> factory)
default <C extends java.util.Collection<? super T>> FluxT<C> groupedUntil(java.util.function.Predicate<? super T> predicate, java.util.function.Supplier<C> factory)
default FluxT<com.aol.cyclops.data.collections.extensions.standard.ListX<T>> grouped(int groupSize)
default <K,A,D> FluxT<org.jooq.lambda.tuple.Tuple2<K,D>> grouped(java.util.function.Function<? super T,? extends K> classifier, java.util.stream.Collector<? super T,A,D> downFlux)
default <K> FluxT<org.jooq.lambda.tuple.Tuple2<K,org.jooq.lambda.Seq<T>>> grouped(java.util.function.Function<? super T,? extends K> classifier)
default <U> FluxT<U> scanLeft(U seed, java.util.function.BiFunction<? super U,? super T,? extends U> function)
default <U> FluxT<U> scanRight(U identity, java.util.function.BiFunction<? super T,? super U,? extends U> combiner)
default <X extends java.lang.Throwable> FluxT<T> onEmptyThrow(java.util.function.Supplier<? extends X> supplier)