public class FluxTs
extends java.lang.Object
| Constructor and Description |
|---|
FluxTs() |
| Modifier and Type | Method and Description |
|---|---|
static <T> AnyMSeq<T> |
anyM(FluxT<T> flux)
Construct an AnyM type from a Flux.
|
static <T> FluxTSeq<T> |
fluxT(org.reactivestreams.Publisher<reactor.core.publisher.Flux<T>> nested)
Construct a FluxT from a Publisher containing Fluxes.
|
static <T,R1,R> FluxT<R> |
forEach(FluxT<? extends T> value1,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<R1>> value2,
java.util.function.BiFunction<? super T,? super R1,java.lang.Boolean> filterFunction,
java.util.function.BiFunction<? super T,? super R1,? extends R> yieldingFunction)
Perform a For Comprehension over a FluxT, accepting a generating function.
|
static <T,R1,R> FluxT<R> |
forEach(FluxT<? extends T> value1,
java.util.function.Function<? super T,org.reactivestreams.Publisher<R1>> value2,
java.util.function.BiFunction<? super T,? super R1,? extends R> yieldingFunction)
Perform a For Comprehension over a FluxT, accepting a generating function.
|
static <T1,T2,R1,R2,R> |
forEach3(FluxT<? extends T1> value1,
java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2,
java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3,
TriFunction<? super T1,? super R1,? super R2,? extends R> yieldingFunction)
Perform a For Comprehension over a FluxT, accepting 2 generating functions.
|
static <T1,T2,R1,R2,R> |
forEach3(FluxT<? extends T1> value1,
java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2,
java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3,
TriFunction<? super T1,? super R1,? super R2,java.lang.Boolean> filterFunction,
TriFunction<? super T1,? super R1,? super R2,? extends R> yieldingFunction)
Perform a For Comprehension over a FluxT, accepting 2 generating functions.
|
static <T1,T2,T3,R1,R2,R3,R> |
forEach4(FluxT<? extends T1> value1,
java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2,
java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3,
TriFunction<? super T1,? super R1,? super R2,? extends org.reactivestreams.Publisher<R3>> value4,
QuadFunction<? super T1,? super R1,? super R2,? super R3,? extends R> yieldingFunction)
Perform a For Comprehension over a FluxT, accepting 3 generating functions.
|
static <T1,T2,T3,R1,R2,R3,R> |
forEach4(FluxT<? extends T1> value1,
java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2,
java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3,
TriFunction<? super T1,? super R1,? super R2,? extends org.reactivestreams.Publisher<R3>> value4,
QuadFunction<? super T1,? super R1,? super R2,? super R3,java.lang.Boolean> filterFunction,
QuadFunction<? super T1,? super R1,? super R2,? super R3,? extends R> yieldingFunction)
Perform a For Comprehension over a FluxT, accepting 3 generating functions.
|
public static <T> AnyMSeq<T> anyM(FluxT<T> flux)
AnyMSeq<Integer> fluxT = Reactor.fromFluxT(fluxT);
AnyMSeq<Integer> transformedFluxT = myGenericOperation(flux);
public AnyMSeq<Integer> myGenericOperation(AnyMSeq<Integer> monad);
fluxT - To wrap inside an AnyMpublic static <T> FluxTSeq<T> fluxT(org.reactivestreams.Publisher<reactor.core.publisher.Flux<T>> nested)
nested - Publisher of Fluxespublic static <T1,T2,T3,R1,R2,R3,R> FluxT<R> forEach4(FluxT<? extends T1> value1, java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2, java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3, TriFunction<? super T1,? super R1,? super R2,? extends org.reactivestreams.Publisher<R3>> value4, QuadFunction<? super T1,? super R1,? super R2,? super R3,? extends R> yieldingFunction)
import static com.aol.cyclops.reactor.FluxTs.forEach4;
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.range(10,2),Flux.range(100,2)));
forEach4(fluxT,
a-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b)-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b,c)-> ReactiveSeq.iterate(a,i->i+1).limit(2),
Tuple::tuple)
.forEach(System.out::println);
//(10, 10, 10, 10)
(11, 11, 11, 11)
(100, 100, 100, 100)
(101, 101, 101, 101)
value1 - top level FluxTvalue2 - Nested publishervalue3 - Nested publishervalue4 - Nested publisheryieldingFunction - Generates a result per combinationpublic static <T1,T2,T3,R1,R2,R3,R> FluxT<R> forEach4(FluxT<? extends T1> value1, java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2, java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3, TriFunction<? super T1,? super R1,? super R2,? extends org.reactivestreams.Publisher<R3>> value4, QuadFunction<? super T1,? super R1,? super R2,? super R3,java.lang.Boolean> filterFunction, QuadFunction<? super T1,? super R1,? super R2,? super R3,? extends R> yieldingFunction)
import static com.aol.cyclops.reactor.FluxTs.forEach4;
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.range(10,2),Flux.range(100,2)));
forEach4(fluxT,
a-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b)-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b,c)-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b,c,d)->a+b+c+d<102,
Tuple::tuple)
.forEach(System.out::println);
//(10, 10, 10, 10)
(11, 11, 11, 11)
value1 - top level FluxTvalue2 - Nested publishervalue3 - Nested publishervalue4 - Nested publisherfilterFunction - A filtering function, keeps values where the predicate holdsyieldingFunction - Generates a result per combinationpublic static <T1,T2,R1,R2,R> FluxT<R> forEach3(FluxT<? extends T1> value1, java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2, java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3, TriFunction<? super T1,? super R1,? super R2,? extends R> yieldingFunction)
import static com.aol.cyclops.reactor.FluxTs.forEach3;
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.range(10,2),Flux.range(100,2)));
forEach3(fluxT,
a-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b)-> ReactiveSeq.iterate(a,i->i+1).limit(2),
Tuple::tuple)
.forEach(System.out::println);
value1 - top level FluxTvalue2 - Nested publishervalue3 - Nested publisheryieldingFunction - Generates a result per combinationpublic static <T1,T2,R1,R2,R> FluxT<R> forEach3(FluxT<? extends T1> value1, java.util.function.Function<? super T1,? extends org.reactivestreams.Publisher<R1>> value2, java.util.function.BiFunction<? super T1,? super R1,? extends org.reactivestreams.Publisher<R2>> value3, TriFunction<? super T1,? super R1,? super R2,java.lang.Boolean> filterFunction, TriFunction<? super T1,? super R1,? super R2,? extends R> yieldingFunction)
import static com.aol.cyclops.reactor.FluxTs.forEach3;
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.range(10,2),Flux.range(100,2)));
forEach3(fluxT,
a-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b)-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b,c)->a+b+c<102,
Tuple::tuple)
.forEach(System.out::println);
value1 - top level FluxTvalue2 - Nested publishervalue3 - Nested publisherfilterFunction - A filtering function, keeps values where the predicate holdsyieldingFunction - Generates a result per combinationpublic static <T,R1,R> FluxT<R> forEach(FluxT<? extends T> value1, java.util.function.Function<? super T,org.reactivestreams.Publisher<R1>> value2, java.util.function.BiFunction<? super T,? super R1,? extends R> yieldingFunction)
import static com.aol.cyclops.reactor.FluxTs.forEach;
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.range(10,2),Flux.range(100,2)));
forEach(fluxT,
a-> ReactiveSeq.iterate(a,i->i+1).limit(2),
Tuple::tuple)
.forEach(System.out::println);
value1 - top level FluxTvalue2 - Nested publishervalue3 - Nested publisheryieldingFunction - Generates a result per combinationpublic static <T,R1,R> FluxT<R> forEach(FluxT<? extends T> value1, java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<R1>> value2, java.util.function.BiFunction<? super T,? super R1,java.lang.Boolean> filterFunction, java.util.function.BiFunction<? super T,? super R1,? extends R> yieldingFunction)
import static com.aol.cyclops.reactor.FluxTs.forEach;
FluxT<Integer> fluxT = FluxT.fromIterable(Arrays.asList(Flux.range(10,2),Flux.range(100,2)));
forEach(fluxT,
a-> ReactiveSeq.iterate(a,i->i+1).limit(2),
(a,b)->a+b<102,
Tuple::tuple)
.forEach(System.out::println);
value1 - top level FluxTvalue2 - Nested publishervalue3 - Nested publisherfilterFunction - A filtering function, keeps values where the predicate holdsyieldingFunction - Generates a result per combination