T - Data type stored within the Fluxpublic final class FluxType<T> extends java.lang.Object implements com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T>, org.reactivestreams.Publisher<T>
| Modifier and Type | Class and Description |
|---|---|
static class |
FluxType.µ
Witness type
|
| Constructor and Description |
|---|
FluxType() |
| Modifier and Type | Method and Description |
|---|---|
reactor.core.publisher.Mono<java.lang.Boolean> |
all(java.util.function.Predicate<? super T> predicate) |
reactor.core.publisher.Mono<java.lang.Boolean> |
any(java.util.function.Predicate<? super T> predicate) |
<P> P |
as(java.util.function.Function<? super reactor.core.publisher.Flux<T>,P> transformer) |
reactor.core.publisher.Flux<T> |
awaitOnSubscribe() |
T |
blockFirst() |
T |
blockFirst(java.time.Duration d) |
T |
blockFirstMillis(long timeout) |
T |
blockLast() |
T |
blockLast(java.time.Duration d) |
T |
blockLastMillis(long timeout) |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer() |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer(java.time.Duration timespan) |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer(java.time.Duration timespan,
java.time.Duration timeshift) |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer(int maxSize) |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer(int maxSize,
java.time.Duration timespan) |
<C extends java.util.Collection<? super T>> |
buffer(int maxSize,
java.time.Duration timespan,
java.util.function.Supplier<C> bufferSupplier) |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer(int maxSize,
int skip) |
<C extends java.util.Collection<? super T>> |
buffer(int maxSize,
int skip,
java.util.function.Supplier<C> bufferSupplier) |
<C extends java.util.Collection<? super T>> |
buffer(int maxSize,
java.util.function.Supplier<C> bufferSupplier) |
reactor.core.publisher.Flux<java.util.List<T>> |
buffer(org.reactivestreams.Publisher<?> other) |
<C extends java.util.Collection<? super T>> |
buffer(org.reactivestreams.Publisher<?> other,
java.util.function.Supplier<C> bufferSupplier) |
<U,V> reactor.core.publisher.Flux<java.util.List<T>> |
buffer(org.reactivestreams.Publisher<U> bucketOpening,
java.util.function.Function<? super U,? extends org.reactivestreams.Publisher<V>> closeSelector) |
<U,V,C extends java.util.Collection<? super T>> |
buffer(org.reactivestreams.Publisher<U> bucketOpening,
java.util.function.Function<? super U,? extends org.reactivestreams.Publisher<V>> closeSelector,
java.util.function.Supplier<C> bufferSupplier) |
reactor.core.publisher.Flux<java.util.List<T>> |
bufferMillis(int maxSize,
long timespan) |
reactor.core.publisher.Flux<java.util.List<T>> |
bufferMillis(int maxSize,
long timespan,
reactor.core.scheduler.TimedScheduler timer) |
<C extends java.util.Collection<? super T>> |
bufferMillis(int maxSize,
long timespan,
reactor.core.scheduler.TimedScheduler timer,
java.util.function.Supplier<C> bufferSupplier) |
reactor.core.publisher.Flux<java.util.List<T>> |
bufferMillis(long timespan) |
reactor.core.publisher.Flux<java.util.List<T>> |
bufferMillis(long timespan,
long timeshift) |
reactor.core.publisher.Flux<java.util.List<T>> |
bufferMillis(long timespan,
long timeshift,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<java.util.List<T>> |
bufferMillis(long timespan,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<T> |
cache() |
reactor.core.publisher.Flux<T> |
cache(java.time.Duration ttl) |
reactor.core.publisher.Flux<T> |
cache(int history) |
reactor.core.publisher.Flux<T> |
cache(int history,
java.time.Duration ttl) |
reactor.core.publisher.Flux<T> |
cancelOn(reactor.core.scheduler.Scheduler scheduler) |
<E> reactor.core.publisher.Flux<E> |
cast(java.lang.Class<E> clazz) |
<R,A> reactor.core.publisher.Mono<R> |
collect(java.util.stream.Collector<T,A,R> collector) |
<E> reactor.core.publisher.Mono<E> |
collect(java.util.function.Supplier<E> containerSupplier,
java.util.function.BiConsumer<E,? super T> collector) |
reactor.core.publisher.Mono<java.util.List<T>> |
collectList() |
<K> reactor.core.publisher.Mono<java.util.Map<K,T>> |
collectMap(java.util.function.Function<? super T,? extends K> keyExtractor) |
<K,V> reactor.core.publisher.Mono<java.util.Map<K,V>> |
collectMap(java.util.function.Function<? super T,? extends K> keyExtractor,
java.util.function.Function<? super T,? extends V> valueExtractor) |
<K,V> reactor.core.publisher.Mono<java.util.Map<K,V>> |
collectMap(java.util.function.Function<? super T,? extends K> keyExtractor,
java.util.function.Function<? super T,? extends V> valueExtractor,
java.util.function.Supplier<java.util.Map<K,V>> mapSupplier) |
<K> reactor.core.publisher.Mono<java.util.Map<K,java.util.Collection<T>>> |
collectMultimap(java.util.function.Function<? super T,? extends K> keyExtractor) |
<K,V> reactor.core.publisher.Mono<java.util.Map<K,java.util.Collection<V>>> |
collectMultimap(java.util.function.Function<? super T,? extends K> keyExtractor,
java.util.function.Function<? super T,? extends V> valueExtractor) |
<K,V> reactor.core.publisher.Mono<java.util.Map<K,java.util.Collection<V>>> |
collectMultimap(java.util.function.Function<? super T,? extends K> keyExtractor,
java.util.function.Function<? super T,? extends V> valueExtractor,
java.util.function.Supplier<java.util.Map<K,java.util.Collection<V>>> mapSupplier) |
reactor.core.publisher.Mono<java.util.List<T>> |
collectSortedList() |
reactor.core.publisher.Mono<java.util.List<T>> |
collectSortedList(java.util.Comparator<? super T> comparator) |
<V> reactor.core.publisher.Flux<V> |
compose(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<V>> transformer) |
<V> reactor.core.publisher.Flux<V> |
concatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper) |
<V> reactor.core.publisher.Flux<V> |
concatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper,
int prefetch) |
<V> reactor.core.publisher.Flux<V> |
concatMapDelayError(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper,
boolean delayUntilEnd,
int prefetch) |
<V> reactor.core.publisher.Flux<V> |
concatMapDelayError(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper,
int prefetch) |
<V> reactor.core.publisher.Flux<V> |
concatMapDelayError(java.util.function.Function<? super T,org.reactivestreams.Publisher<? extends V>> mapper) |
<R> reactor.core.publisher.Flux<R> |
concatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper) |
<R> reactor.core.publisher.Flux<R> |
concatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper,
int prefetch) |
reactor.core.publisher.Flux<T> |
concatWith(org.reactivestreams.Publisher<? extends T> other) |
reactor.core.publisher.Mono<java.lang.Long> |
count() |
reactor.core.publisher.Flux<T> |
defaultIfEmpty(T defaultV) |
reactor.core.publisher.Flux<T> |
delay(java.time.Duration delay) |
reactor.core.publisher.Flux<T> |
delayMillis(long delay) |
reactor.core.publisher.Flux<T> |
delayMillis(long delay,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<T> |
delaySubscription(java.time.Duration delay) |
<U> reactor.core.publisher.Flux<T> |
delaySubscription(org.reactivestreams.Publisher<U> subscriptionDelay) |
reactor.core.publisher.Flux<T> |
delaySubscriptionMillis(long delay) |
reactor.core.publisher.Flux<T> |
delaySubscriptionMillis(long delay,
reactor.core.scheduler.TimedScheduler timer) |
<X> reactor.core.publisher.Flux<X> |
dematerialize() |
reactor.core.publisher.Flux<T> |
distinct() |
<V> reactor.core.publisher.Flux<T> |
distinct(java.util.function.Function<? super T,? extends V> keySelector) |
reactor.core.publisher.Flux<T> |
distinctUntilChanged() |
<V> reactor.core.publisher.Flux<T> |
distinctUntilChanged(java.util.function.Function<? super T,? extends V> keySelector) |
reactor.core.publisher.Flux<T> |
doAfterTerminate(java.lang.Runnable afterTerminate) |
reactor.core.publisher.Flux<T> |
doOnCancel(java.lang.Runnable onCancel) |
reactor.core.publisher.Flux<T> |
doOnComplete(java.lang.Runnable onComplete) |
<E extends java.lang.Throwable> |
doOnError(java.lang.Class<E> exceptionType,
java.util.function.Consumer<? super E> onError) |
reactor.core.publisher.Flux<T> |
doOnError(java.util.function.Consumer<? super java.lang.Throwable> onError) |
reactor.core.publisher.Flux<T> |
doOnError(java.util.function.Predicate<? super java.lang.Throwable> predicate,
java.util.function.Consumer<? super java.lang.Throwable> onError) |
reactor.core.publisher.Flux<T> |
doOnNext(java.util.function.Consumer<? super T> onNext) |
reactor.core.publisher.Flux<T> |
doOnRequest(java.util.function.LongConsumer consumer) |
reactor.core.publisher.Flux<T> |
doOnSubscribe(java.util.function.Consumer<? super org.reactivestreams.Subscription> onSubscribe) |
reactor.core.publisher.Flux<T> |
doOnTerminate(java.lang.Runnable onTerminate) |
reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> |
elapsed() |
reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> |
elapsed(reactor.core.scheduler.TimedScheduler scheduler) |
reactor.core.publisher.Mono<T> |
elementAt(int index) |
reactor.core.publisher.Mono<T> |
elementAt(int index,
T defaultValue) |
static <T> FluxType<T> |
empty() |
boolean |
equals(java.lang.Object obj) |
reactor.core.publisher.Flux<T> |
filter(java.util.function.Predicate<? super T> p) |
reactor.core.publisher.Flux<T> |
firstEmittingWith(org.reactivestreams.Publisher<? extends T> other) |
<R> reactor.core.publisher.Flux<R> |
flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends R>> mapper) |
<R> reactor.core.publisher.Flux<R> |
flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends R>> mapperOnNext,
java.util.function.Function<java.lang.Throwable,? extends org.reactivestreams.Publisher<? extends R>> mapperOnError,
java.util.function.Supplier<? extends org.reactivestreams.Publisher<? extends R>> mapperOnComplete) |
<V> reactor.core.publisher.Flux<V> |
flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper,
boolean delayError,
int concurrency,
int prefetch) |
<V> reactor.core.publisher.Flux<V> |
flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper,
int concurrency) |
<V> reactor.core.publisher.Flux<V> |
flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper,
int concurrency,
int prefetch) |
<R> reactor.core.publisher.Flux<R> |
flatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper) |
<R> reactor.core.publisher.Flux<R> |
flatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper,
int prefetch) |
long |
getPrefetch() |
<K> reactor.core.publisher.Flux<reactor.core.publisher.GroupedFlux<K,T>> |
groupBy(java.util.function.Function<? super T,? extends K> keyMapper) |
<K,V> reactor.core.publisher.Flux<reactor.core.publisher.GroupedFlux<K,V>> |
groupBy(java.util.function.Function<? super T,? extends K> keyMapper,
java.util.function.Function<? super T,? extends V> valueMapper) |
<TRight,TLeftEnd,TRightEnd,R> |
groupJoin(org.reactivestreams.Publisher<? extends TRight> other,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<TLeftEnd>> leftEnd,
java.util.function.Function<? super TRight,? extends org.reactivestreams.Publisher<TRightEnd>> rightEnd,
java.util.function.BiFunction<? super T,? super reactor.core.publisher.Flux<TRight>,? extends R> resultSelector) |
<R> reactor.core.publisher.Flux<R> |
handle(java.util.function.BiConsumer<? super T,reactor.core.publisher.SynchronousSink<R>> handler) |
reactor.core.publisher.Mono<java.lang.Boolean> |
hasElement(T value) |
reactor.core.publisher.Mono<java.lang.Boolean> |
hasElements() |
int |
hashCode() |
reactor.core.publisher.Flux<T> |
hide() |
reactor.core.publisher.Mono<T> |
ignoreElements() |
<TRight,TLeftEnd,TRightEnd,R> |
join(org.reactivestreams.Publisher<? extends TRight> other,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<TLeftEnd>> leftEnd,
java.util.function.Function<? super TRight,? extends org.reactivestreams.Publisher<TRightEnd>> rightEnd,
java.util.function.BiFunction<? super T,? super TRight,? extends R> resultSelector) |
static <T> FluxType<T> |
just(T... values) |
static <T> FluxType<T> |
just(T value)
Construct a HKT encoded completed Flux
|
reactor.core.publisher.Mono<T> |
last() |
reactor.core.publisher.Mono<T> |
last(T defaultValue) |
reactor.core.publisher.Flux<T> |
log() |
reactor.core.publisher.Flux<T> |
log(java.lang.String category) |
reactor.core.publisher.Flux<T> |
log(java.lang.String category,
java.util.logging.Level level,
boolean showOperatorLine,
reactor.core.publisher.SignalType... options) |
reactor.core.publisher.Flux<T> |
log(java.lang.String category,
java.util.logging.Level level,
reactor.core.publisher.SignalType... options) |
<V> reactor.core.publisher.Flux<V> |
map(java.util.function.Function<? super T,? extends V> mapper) |
<E extends java.lang.Throwable> |
mapError(java.lang.Class<E> type,
java.util.function.Function<? super E,? extends java.lang.Throwable> mapper) |
reactor.core.publisher.Flux<T> |
mapError(java.util.function.Function<? super java.lang.Throwable,? extends java.lang.Throwable> mapper) |
reactor.core.publisher.Flux<T> |
mapError(java.util.function.Predicate<? super java.lang.Throwable> predicate,
java.util.function.Function<? super java.lang.Throwable,? extends java.lang.Throwable> mapper) |
reactor.core.publisher.Flux<reactor.core.publisher.Signal<T>> |
materialize() |
reactor.core.publisher.Flux<T> |
mergeWith(org.reactivestreams.Publisher<? extends T> other) |
reactor.core.publisher.Flux<T> |
narrow() |
static <T> reactor.core.publisher.Flux<T> |
narrow(com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T> completableFlux)
Convert the HigherKindedType definition for a Flux into
|
static <T> FluxType<T> |
narrowK(com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T> future)
Convert the raw Higher Kinded Type for FluxType types into the FluxType type definition class
|
reactor.core.publisher.Mono<T> |
next() |
<U> reactor.core.publisher.Flux<U> |
ofType(java.lang.Class<U> clazz) |
reactor.core.publisher.Flux<T> |
onBackpressureBuffer() |
reactor.core.publisher.Flux<T> |
onBackpressureBuffer(int maxSize) |
reactor.core.publisher.Flux<T> |
onBackpressureBuffer(int maxSize,
java.util.function.Consumer<? super T> onOverflow) |
reactor.core.publisher.Flux<T> |
onBackpressureDrop() |
reactor.core.publisher.Flux<T> |
onBackpressureDrop(java.util.function.Consumer<? super T> onDropped) |
reactor.core.publisher.Flux<T> |
onBackpressureError() |
reactor.core.publisher.Flux<T> |
onBackpressureLatest() |
<E extends java.lang.Throwable> |
onErrorResumeWith(java.lang.Class<E> type,
java.util.function.Function<? super E,? extends org.reactivestreams.Publisher<? extends T>> fallback) |
reactor.core.publisher.Flux<T> |
onErrorResumeWith(java.util.function.Function<? super java.lang.Throwable,? extends org.reactivestreams.Publisher<? extends T>> fallback) |
reactor.core.publisher.Flux<T> |
onErrorResumeWith(java.util.function.Predicate<? super java.lang.Throwable> predicate,
java.util.function.Function<? super java.lang.Throwable,? extends org.reactivestreams.Publisher<? extends T>> fallback) |
<E extends java.lang.Throwable> |
onErrorReturn(java.lang.Class<E> type,
T fallbackValue) |
<E extends java.lang.Throwable> |
onErrorReturn(java.util.function.Predicate<? super java.lang.Throwable> predicate,
T fallbackValue) |
reactor.core.publisher.Flux<T> |
onErrorReturn(T fallbackValue) |
reactor.core.publisher.Flux<T> |
onTerminateDetach() |
reactor.core.publisher.ParallelFlux<T> |
parallel() |
reactor.core.publisher.ParallelFlux<T> |
parallel(int parallelism) |
reactor.core.publisher.ParallelFlux<T> |
parallel(int parallelism,
int prefetch) |
reactor.core.publisher.ConnectableFlux<T> |
publish() |
<R> reactor.core.publisher.Flux<R> |
publish(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<? extends R>> transform) |
<R> reactor.core.publisher.Flux<R> |
publish(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<? extends R>> transform,
int prefetch) |
reactor.core.publisher.ConnectableFlux<T> |
publish(int prefetch) |
reactor.core.publisher.Mono<T> |
publishNext() |
reactor.core.publisher.Flux<T> |
publishOn(reactor.core.scheduler.Scheduler scheduler) |
reactor.core.publisher.Flux<T> |
publishOn(reactor.core.scheduler.Scheduler scheduler,
int prefetch) |
<A> reactor.core.publisher.Mono<A> |
reduce(A initial,
java.util.function.BiFunction<A,? super T,A> accumulator) |
reactor.core.publisher.Mono<T> |
reduce(java.util.function.BiFunction<T,T,T> aggregator) |
<A> reactor.core.publisher.Mono<A> |
reduceWith(java.util.function.Supplier<A> initial,
java.util.function.BiFunction<A,? super T,A> accumulator) |
reactor.core.publisher.Flux<T> |
repeat() |
reactor.core.publisher.Flux<T> |
repeat(java.util.function.BooleanSupplier predicate) |
reactor.core.publisher.Flux<T> |
repeat(long numRepeat) |
reactor.core.publisher.Flux<T> |
repeat(long numRepeat,
java.util.function.BooleanSupplier predicate) |
reactor.core.publisher.Flux<T> |
repeatWhen(java.util.function.Function<reactor.core.publisher.Flux<java.lang.Long>,? extends org.reactivestreams.Publisher<?>> whenFactory) |
reactor.core.publisher.ConnectableFlux<T> |
replay() |
reactor.core.publisher.ConnectableFlux<T> |
replay(java.time.Duration ttl) |
reactor.core.publisher.ConnectableFlux<T> |
replay(int history) |
reactor.core.publisher.ConnectableFlux<T> |
replay(int history,
java.time.Duration ttl) |
reactor.core.publisher.ConnectableFlux<T> |
replayMillis(int history,
long ttl,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.ConnectableFlux<T> |
replayMillis(long ttl,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<T> |
retry() |
reactor.core.publisher.Flux<T> |
retry(long numRetries) |
reactor.core.publisher.Flux<T> |
retry(long numRetries,
java.util.function.Predicate<java.lang.Throwable> retryMatcher) |
reactor.core.publisher.Flux<T> |
retry(java.util.function.Predicate<java.lang.Throwable> retryMatcher) |
reactor.core.publisher.Flux<T> |
retryWhen(java.util.function.Function<reactor.core.publisher.Flux<java.lang.Throwable>,? extends org.reactivestreams.Publisher<?>> whenFactory) |
reactor.core.publisher.Flux<T> |
sample(java.time.Duration timespan) |
<U> reactor.core.publisher.Flux<T> |
sample(org.reactivestreams.Publisher<U> sampler) |
reactor.core.publisher.Flux<T> |
sampleFirst(java.time.Duration timespan) |
<U> reactor.core.publisher.Flux<T> |
sampleFirst(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<U>> samplerFactory) |
reactor.core.publisher.Flux<T> |
sampleFirstMillis(long timespan) |
reactor.core.publisher.Flux<T> |
sampleMillis(long timespan) |
<U> reactor.core.publisher.Flux<T> |
sampleTimeout(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<U>> throttlerFactory) |
<U> reactor.core.publisher.Flux<T> |
sampleTimeout(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<U>> throttlerFactory,
int maxConcurrency) |
<A> reactor.core.publisher.Flux<A> |
scan(A initial,
java.util.function.BiFunction<A,? super T,A> accumulator) |
reactor.core.publisher.Flux<T> |
scan(java.util.function.BiFunction<T,T,T> accumulator) |
<A> reactor.core.publisher.Flux<A> |
scanWith(java.util.function.Supplier<A> initial,
java.util.function.BiFunction<A,? super T,A> accumulator) |
reactor.core.publisher.Flux<T> |
share() |
reactor.core.publisher.Mono<T> |
single() |
reactor.core.publisher.Mono<T> |
single(T defaultValue) |
reactor.core.publisher.Mono<T> |
singleOrEmpty() |
reactor.core.publisher.Flux<T> |
skip(java.time.Duration timespan) |
reactor.core.publisher.Flux<T> |
skip(long skipped) |
reactor.core.publisher.Flux<T> |
skipLast(int n) |
reactor.core.publisher.Flux<T> |
skipMillis(long timespan) |
reactor.core.publisher.Flux<T> |
skipMillis(long timespan,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<T> |
skipUntil(java.util.function.Predicate<? super T> untilPredicate) |
reactor.core.publisher.Flux<T> |
skipUntilOther(org.reactivestreams.Publisher<?> other) |
reactor.core.publisher.Flux<T> |
skipWhile(java.util.function.Predicate<? super T> skipPredicate) |
reactor.core.publisher.Flux<T> |
sort() |
reactor.core.publisher.Flux<T> |
sort(java.util.Comparator<? super T> sortFunction) |
reactor.core.publisher.Flux<T> |
startWith(java.lang.Iterable<? extends T> iterable) |
reactor.core.publisher.Flux<T> |
startWith(org.reactivestreams.Publisher<? extends T> publisher) |
reactor.core.publisher.Flux<T> |
startWith(T... values) |
reactor.core.Cancellation |
subscribe() |
reactor.core.Cancellation |
subscribe(java.util.function.Consumer<? super T> consumer) |
reactor.core.Cancellation |
subscribe(java.util.function.Consumer<? super T> consumer,
java.util.function.Consumer<? super java.lang.Throwable> errorConsumer) |
reactor.core.Cancellation |
subscribe(java.util.function.Consumer<? super T> consumer,
java.util.function.Consumer<? super java.lang.Throwable> errorConsumer,
java.lang.Runnable completeConsumer) |
reactor.core.Cancellation |
subscribe(java.util.function.Consumer<? super T> consumer,
java.util.function.Consumer<? super java.lang.Throwable> errorConsumer,
java.lang.Runnable completeConsumer,
int prefetch) |
reactor.core.Cancellation |
subscribe(java.util.function.Consumer<? super T> consumer,
int prefetch) |
reactor.core.Cancellation |
subscribe(int prefetch) |
void |
subscribe(org.reactivestreams.Subscriber<? super T> s) |
reactor.core.publisher.Flux<T> |
subscribeOn(reactor.core.scheduler.Scheduler scheduler) |
<E extends org.reactivestreams.Subscriber<? super T>> |
subscribeWith(E subscriber) |
reactor.core.publisher.Flux<T> |
switchIfEmpty(org.reactivestreams.Publisher<? extends T> alternate) |
<V> reactor.core.publisher.Flux<V> |
switchMap(java.util.function.Function<? super T,org.reactivestreams.Publisher<? extends V>> fn) |
<V> reactor.core.publisher.Flux<V> |
switchMap(java.util.function.Function<? super T,org.reactivestreams.Publisher<? extends V>> fn,
int prefetch) |
<E extends java.lang.Throwable> |
switchOnError(java.lang.Class<E> type,
org.reactivestreams.Publisher<? extends T> fallback) |
reactor.core.publisher.Flux<T> |
switchOnError(java.util.function.Predicate<? super java.lang.Throwable> predicate,
org.reactivestreams.Publisher<? extends T> fallback) |
reactor.core.publisher.Flux<T> |
switchOnError(org.reactivestreams.Publisher<? extends T> fallback) |
reactor.core.publisher.Flux<T> |
take(java.time.Duration timespan) |
reactor.core.publisher.Flux<T> |
take(long n) |
reactor.core.publisher.Flux<T> |
takeLast(int n) |
reactor.core.publisher.Flux<T> |
takeMillis(long timespan) |
reactor.core.publisher.Flux<T> |
takeMillis(long timespan,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<T> |
takeUntil(java.util.function.Predicate<? super T> predicate) |
reactor.core.publisher.Flux<T> |
takeUntilOther(org.reactivestreams.Publisher<?> other) |
reactor.core.publisher.Flux<T> |
takeWhile(java.util.function.Predicate<? super T> continuePredicate) |
reactor.core.publisher.Mono<java.lang.Void> |
then() |
reactor.core.publisher.Mono<java.lang.Void> |
then(org.reactivestreams.Publisher<java.lang.Void> other) |
reactor.core.publisher.Mono<java.lang.Void> |
then(java.util.function.Supplier<? extends org.reactivestreams.Publisher<java.lang.Void>> afterSupplier) |
<V> reactor.core.publisher.Flux<V> |
thenMany(org.reactivestreams.Publisher<V> other) |
<V> reactor.core.publisher.Flux<V> |
thenMany(java.util.function.Supplier<? extends org.reactivestreams.Publisher<V>> afterSupplier) |
reactor.core.publisher.Flux<T> |
timeout(java.time.Duration timeout) |
reactor.core.publisher.Flux<T> |
timeout(java.time.Duration timeout,
org.reactivestreams.Publisher<? extends T> fallback) |
<U> reactor.core.publisher.Flux<T> |
timeout(org.reactivestreams.Publisher<U> firstTimeout) |
<U,V> reactor.core.publisher.Flux<T> |
timeout(org.reactivestreams.Publisher<U> firstTimeout,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<V>> nextTimeoutFactory) |
<U,V> reactor.core.publisher.Flux<T> |
timeout(org.reactivestreams.Publisher<U> firstTimeout,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<V>> nextTimeoutFactory,
org.reactivestreams.Publisher<? extends T> fallback) |
reactor.core.publisher.Flux<T> |
timeoutMillis(long timeout) |
reactor.core.publisher.Flux<T> |
timeoutMillis(long timeout,
org.reactivestreams.Publisher<? extends T> fallback) |
reactor.core.publisher.Flux<T> |
timeoutMillis(long timeout,
org.reactivestreams.Publisher<? extends T> fallback,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<T> |
timeoutMillis(long timeout,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> |
timestamp() |
reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> |
timestamp(reactor.core.scheduler.TimedScheduler scheduler) |
java.lang.Iterable<T> |
toIterable() |
java.lang.Iterable<T> |
toIterable(long batchSize) |
java.lang.Iterable<T> |
toIterable(long batchSize,
java.util.function.Supplier<java.util.Queue<T>> queueProvider) |
ReactiveSeq<T> |
toReactiveSeq() |
java.util.stream.Stream<T> |
toStream() |
java.util.stream.Stream<T> |
toStream(int batchSize) |
java.lang.String |
toString() |
<V> reactor.core.publisher.Flux<V> |
transform(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<V>> transformer) |
static <T> FluxType<T> |
widen(reactor.core.publisher.Flux<T> completableFlux)
Convert a Flux to a simulated HigherKindedType that captures Flux nature
and Flux element data type separately.
|
static <T> FluxType<T> |
widen(org.reactivestreams.Publisher<T> completableFlux) |
static <C2,T> com.aol.cyclops.hkt.alias.Higher<C2,com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T>> |
widen2(com.aol.cyclops.hkt.alias.Higher<C2,FluxType<T>> flux)
Widen a FluxType nested inside another HKT encoded type
|
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window() |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(java.time.Duration timespan) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(java.time.Duration timespan,
java.time.Duration timeshift) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(int maxSize) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(int maxSize,
java.time.Duration timespan) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(int maxSize,
int skip) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(org.reactivestreams.Publisher<?> boundary) |
<U,V> reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
window(org.reactivestreams.Publisher<U> bucketOpening,
java.util.function.Function<? super U,? extends org.reactivestreams.Publisher<V>> closeSelector) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
windowMillis(int maxSize,
long timespan) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
windowMillis(int maxSize,
long timespan,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
windowMillis(long timespan) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
windowMillis(long timespan,
long timeshift,
reactor.core.scheduler.TimedScheduler timer) |
reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> |
windowMillis(long timespan,
reactor.core.scheduler.TimedScheduler timer) |
<U,R> reactor.core.publisher.Flux<R> |
withLatestFrom(org.reactivestreams.Publisher<? extends U> other,
java.util.function.BiFunction<? super T,? super U,? extends R> resultSelector) |
<T2> reactor.core.publisher.Flux<reactor.util.function.Tuple2<T,T2>> |
zipWith(org.reactivestreams.Publisher<? extends T2> source2) |
<T2,V> reactor.core.publisher.Flux<V> |
zipWith(org.reactivestreams.Publisher<? extends T2> source2,
java.util.function.BiFunction<? super T,? super T2,? extends V> combinator) |
<T2> reactor.core.publisher.Flux<reactor.util.function.Tuple2<T,T2>> |
zipWith(org.reactivestreams.Publisher<? extends T2> source2,
int prefetch) |
<T2,V> reactor.core.publisher.Flux<V> |
zipWith(org.reactivestreams.Publisher<? extends T2> source2,
int prefetch,
java.util.function.BiFunction<? super T,? super T2,? extends V> combinator) |
<T2> reactor.core.publisher.Flux<reactor.util.function.Tuple2<T,T2>> |
zipWithIterable(java.lang.Iterable<? extends T2> iterable) |
<T2,V> reactor.core.publisher.Flux<V> |
zipWithIterable(java.lang.Iterable<? extends T2> iterable,
java.util.function.BiFunction<? super T,? super T2,? extends V> zipper) |
public static <T> FluxType<T> just(T value)
value - To encode inside a HKT encoded Fluxpublic static <T> FluxType<T> just(T... values)
public static <T> FluxType<T> empty()
public static <T> FluxType<T> widen(reactor.core.publisher.Flux<T> completableFlux)
Flux - Flux to widen to a FluxTypepublic static <C2,T> com.aol.cyclops.hkt.alias.Higher<C2,com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T>> widen2(com.aol.cyclops.hkt.alias.Higher<C2,FluxType<T>> flux)
flux - HTK encoded type containing a Flux to widenpublic static <T> FluxType<T> widen(org.reactivestreams.Publisher<T> completableFlux)
public static <T> FluxType<T> narrowK(com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T> future)
future - HKT encoded list into a FluxTypepublic static <T> reactor.core.publisher.Flux<T> narrow(com.aol.cyclops.hkt.alias.Higher<FluxType.µ,T> completableFlux)
Flux - Type Constructor to convert back into narrowed typepublic reactor.core.publisher.Flux<T> narrow()
public ReactiveSeq<T> toReactiveSeq()
public void subscribe(org.reactivestreams.Subscriber<? super T> s)
subscribe in interface org.reactivestreams.Publisher<T>s - Publisher.subscribe(org.reactivestreams.Subscriber)public int hashCode()
hashCode in class java.lang.ObjectObject.hashCode()public boolean equals(java.lang.Object obj)
equals in class java.lang.Objectobj - Object.equals(java.lang.Object)public java.lang.String toString()
toString in class java.lang.ObjectFlux.toString()public final reactor.core.publisher.Mono<java.lang.Boolean> all(java.util.function.Predicate<? super T> predicate)
predicate - Flux.all(java.util.function.Predicate)public final reactor.core.publisher.Mono<java.lang.Boolean> any(java.util.function.Predicate<? super T> predicate)
predicate - Flux.any(java.util.function.Predicate)public final <P> P as(java.util.function.Function<? super reactor.core.publisher.Flux<T>,P> transformer)
transformer - Flux.as(java.util.function.Function)public final reactor.core.publisher.Flux<T> awaitOnSubscribe()
Flux.awaitOnSubscribe()public final T blockFirst()
Flux.blockFirst()public final T blockFirst(java.time.Duration d)
d - Flux.blockFirst(java.time.Duration)public final T blockFirstMillis(long timeout)
timeout - Flux.blockFirstMillis(long)public final T blockLast()
Flux.blockLast()public final T blockLast(java.time.Duration d)
d - Flux.blockLast(java.time.Duration)public final T blockLastMillis(long timeout)
timeout - Flux.blockLastMillis(long)public final reactor.core.publisher.Flux<java.util.List<T>> buffer()
Flux.buffer()public final reactor.core.publisher.Flux<java.util.List<T>> buffer(int maxSize)
maxSize - Flux.buffer(int)public final <C extends java.util.Collection<? super T>> reactor.core.publisher.Flux<C> buffer(int maxSize, java.util.function.Supplier<C> bufferSupplier)
maxSize - bufferSupplier - Flux.buffer(int, java.util.function.Supplier)public final reactor.core.publisher.Flux<java.util.List<T>> buffer(int maxSize, int skip)
maxSize - skip - Flux.buffer(int, int)public final <C extends java.util.Collection<? super T>> reactor.core.publisher.Flux<C> buffer(int maxSize, int skip, java.util.function.Supplier<C> bufferSupplier)
maxSize - skip - bufferSupplier - Flux.buffer(int, int, java.util.function.Supplier)public final reactor.core.publisher.Flux<java.util.List<T>> buffer(org.reactivestreams.Publisher<?> other)
other - Flux.buffer(org.reactivestreams.Publisher)public final <C extends java.util.Collection<? super T>> reactor.core.publisher.Flux<C> buffer(org.reactivestreams.Publisher<?> other, java.util.function.Supplier<C> bufferSupplier)
other - bufferSupplier - Flux.buffer(org.reactivestreams.Publisher, java.util.function.Supplier)public final <U,V> reactor.core.publisher.Flux<java.util.List<T>> buffer(org.reactivestreams.Publisher<U> bucketOpening, java.util.function.Function<? super U,? extends org.reactivestreams.Publisher<V>> closeSelector)
bucketOpening - closeSelector - Flux.buffer(org.reactivestreams.Publisher, java.util.function.Function)public final <U,V,C extends java.util.Collection<? super T>> reactor.core.publisher.Flux<C> buffer(org.reactivestreams.Publisher<U> bucketOpening, java.util.function.Function<? super U,? extends org.reactivestreams.Publisher<V>> closeSelector, java.util.function.Supplier<C> bufferSupplier)
bucketOpening - closeSelector - bufferSupplier - Flux.buffer(org.reactivestreams.Publisher, java.util.function.Function, java.util.function.Supplier)public final reactor.core.publisher.Flux<java.util.List<T>> buffer(java.time.Duration timespan)
timespan - Flux.buffer(java.time.Duration)public final reactor.core.publisher.Flux<java.util.List<T>> buffer(java.time.Duration timespan, java.time.Duration timeshift)
timespan - timeshift - Flux.buffer(java.time.Duration, java.time.Duration)public final reactor.core.publisher.Flux<java.util.List<T>> buffer(int maxSize, java.time.Duration timespan)
maxSize - timespan - Flux.buffer(int, java.time.Duration)public final <C extends java.util.Collection<? super T>> reactor.core.publisher.Flux<C> buffer(int maxSize, java.time.Duration timespan, java.util.function.Supplier<C> bufferSupplier)
maxSize - timespan - bufferSupplier - Flux.buffer(int, java.time.Duration, java.util.function.Supplier)public final reactor.core.publisher.Flux<java.util.List<T>> bufferMillis(long timespan)
timespan - Flux.bufferMillis(long)public final reactor.core.publisher.Flux<java.util.List<T>> bufferMillis(long timespan, reactor.core.scheduler.TimedScheduler timer)
timespan - timer - Flux.bufferMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<java.util.List<T>> bufferMillis(long timespan, long timeshift)
timespan - timeshift - Flux.bufferMillis(long, long)public final reactor.core.publisher.Flux<java.util.List<T>> bufferMillis(long timespan, long timeshift, reactor.core.scheduler.TimedScheduler timer)
timespan - timeshift - timer - Flux.bufferMillis(long, long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<java.util.List<T>> bufferMillis(int maxSize, long timespan)
maxSize - timespan - Flux.bufferMillis(int, long)public final reactor.core.publisher.Flux<java.util.List<T>> bufferMillis(int maxSize, long timespan, reactor.core.scheduler.TimedScheduler timer)
maxSize - timespan - timer - Flux.bufferMillis(int, long, reactor.core.scheduler.TimedScheduler)public final <C extends java.util.Collection<? super T>> reactor.core.publisher.Flux<C> bufferMillis(int maxSize, long timespan, reactor.core.scheduler.TimedScheduler timer, java.util.function.Supplier<C> bufferSupplier)
maxSize - timespan - timer - bufferSupplier - Flux.bufferMillis(int, long, reactor.core.scheduler.TimedScheduler, java.util.function.Supplier)public final reactor.core.publisher.Flux<T> cache()
Flux.cache()public final reactor.core.publisher.Flux<T> cache(int history)
history - Flux.cache(int)public final reactor.core.publisher.Flux<T> cache(java.time.Duration ttl)
ttl - Flux.cache(java.time.Duration)public final reactor.core.publisher.Flux<T> cache(int history, java.time.Duration ttl)
history - ttl - Flux.cache(int, java.time.Duration)public final <E> reactor.core.publisher.Flux<E> cast(java.lang.Class<E> clazz)
clazz - Flux.cast(java.lang.Class)public final reactor.core.publisher.Flux<T> cancelOn(reactor.core.scheduler.Scheduler scheduler)
scheduler - Flux.cancelOn(reactor.core.scheduler.Scheduler)public final <E> reactor.core.publisher.Mono<E> collect(java.util.function.Supplier<E> containerSupplier,
java.util.function.BiConsumer<E,? super T> collector)
containerSupplier - collector - Flux.collect(java.util.function.Supplier, java.util.function.BiConsumer)public final <R,A> reactor.core.publisher.Mono<R> collect(java.util.stream.Collector<T,A,R> collector)
collector - Flux.collect(java.util.stream.Collector)public final reactor.core.publisher.Mono<java.util.List<T>> collectList()
Flux.collectList()public final <K> reactor.core.publisher.Mono<java.util.Map<K,T>> collectMap(java.util.function.Function<? super T,? extends K> keyExtractor)
keyExtractor - Flux.collectMap(java.util.function.Function)public final <K,V> reactor.core.publisher.Mono<java.util.Map<K,V>> collectMap(java.util.function.Function<? super T,? extends K> keyExtractor, java.util.function.Function<? super T,? extends V> valueExtractor)
keyExtractor - valueExtractor - Flux.collectMap(java.util.function.Function, java.util.function.Function)public final <K,V> reactor.core.publisher.Mono<java.util.Map<K,V>> collectMap(java.util.function.Function<? super T,? extends K> keyExtractor, java.util.function.Function<? super T,? extends V> valueExtractor, java.util.function.Supplier<java.util.Map<K,V>> mapSupplier)
keyExtractor - valueExtractor - mapSupplier - Flux.collectMap(java.util.function.Function, java.util.function.Function, java.util.function.Supplier)public final <K> reactor.core.publisher.Mono<java.util.Map<K,java.util.Collection<T>>> collectMultimap(java.util.function.Function<? super T,? extends K> keyExtractor)
keyExtractor - Flux.collectMultimap(java.util.function.Function)public final <K,V> reactor.core.publisher.Mono<java.util.Map<K,java.util.Collection<V>>> collectMultimap(java.util.function.Function<? super T,? extends K> keyExtractor, java.util.function.Function<? super T,? extends V> valueExtractor)
keyExtractor - valueExtractor - Flux.collectMultimap(java.util.function.Function, java.util.function.Function)public final <K,V> reactor.core.publisher.Mono<java.util.Map<K,java.util.Collection<V>>> collectMultimap(java.util.function.Function<? super T,? extends K> keyExtractor, java.util.function.Function<? super T,? extends V> valueExtractor, java.util.function.Supplier<java.util.Map<K,java.util.Collection<V>>> mapSupplier)
keyExtractor - valueExtractor - mapSupplier - Flux.collectMultimap(java.util.function.Function, java.util.function.Function, java.util.function.Supplier)public final reactor.core.publisher.Mono<java.util.List<T>> collectSortedList()
Flux.collectSortedList()public final reactor.core.publisher.Mono<java.util.List<T>> collectSortedList(java.util.Comparator<? super T> comparator)
comparator - Flux.collectSortedList(java.util.Comparator)public final <V> reactor.core.publisher.Flux<V> compose(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<V>> transformer)
transformer - Flux.compose(java.util.function.Function)public final <V> reactor.core.publisher.Flux<V> concatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper)
mapper - Flux.concatMap(java.util.function.Function)public final <V> reactor.core.publisher.Flux<V> concatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper, int prefetch)
mapper - prefetch - Flux.concatMap(java.util.function.Function, int)public final <V> reactor.core.publisher.Flux<V> concatMapDelayError(java.util.function.Function<? super T,org.reactivestreams.Publisher<? extends V>> mapper)
mapper - Flux.concatMapDelayError(java.util.function.Function)public final <V> reactor.core.publisher.Flux<V> concatMapDelayError(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper, int prefetch)
mapper - prefetch - Flux.concatMapDelayError(java.util.function.Function, int)public final <V> reactor.core.publisher.Flux<V> concatMapDelayError(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper, boolean delayUntilEnd, int prefetch)
mapper - delayUntilEnd - prefetch - Flux.concatMapDelayError(java.util.function.Function, boolean, int)public final <R> reactor.core.publisher.Flux<R> concatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper)
mapper - Flux.concatMapIterable(java.util.function.Function)public final <R> reactor.core.publisher.Flux<R> concatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper, int prefetch)
mapper - prefetch - Flux.concatMapIterable(java.util.function.Function, int)public final reactor.core.publisher.Flux<T> concatWith(org.reactivestreams.Publisher<? extends T> other)
other - Flux.concatWith(org.reactivestreams.Publisher)public final reactor.core.publisher.Mono<java.lang.Long> count()
Flux.count()public final reactor.core.publisher.Flux<T> defaultIfEmpty(T defaultV)
defaultV - Flux.defaultIfEmpty(java.lang.Object)public final reactor.core.publisher.Flux<T> delay(java.time.Duration delay)
delay - Flux.delay(java.time.Duration)public final reactor.core.publisher.Flux<T> delayMillis(long delay)
delay - Flux.delayMillis(long)public final reactor.core.publisher.Flux<T> delayMillis(long delay, reactor.core.scheduler.TimedScheduler timer)
delay - timer - Flux.delayMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<T> delaySubscription(java.time.Duration delay)
delay - Flux.delaySubscription(java.time.Duration)public final <U> reactor.core.publisher.Flux<T> delaySubscription(org.reactivestreams.Publisher<U> subscriptionDelay)
subscriptionDelay - Flux.delaySubscription(org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> delaySubscriptionMillis(long delay)
delay - Flux.delaySubscriptionMillis(long)public final reactor.core.publisher.Flux<T> delaySubscriptionMillis(long delay, reactor.core.scheduler.TimedScheduler timer)
delay - timer - Flux.delaySubscriptionMillis(long, reactor.core.scheduler.TimedScheduler)public final <X> reactor.core.publisher.Flux<X> dematerialize()
Flux.dematerialize()public final reactor.core.publisher.Flux<T> distinct()
Flux.distinct()public final <V> reactor.core.publisher.Flux<T> distinct(java.util.function.Function<? super T,? extends V> keySelector)
keySelector - Flux.distinct(java.util.function.Function)public final reactor.core.publisher.Flux<T> distinctUntilChanged()
Flux.distinctUntilChanged()public final <V> reactor.core.publisher.Flux<T> distinctUntilChanged(java.util.function.Function<? super T,? extends V> keySelector)
keySelector - Flux.distinctUntilChanged(java.util.function.Function)public final reactor.core.publisher.Flux<T> doAfterTerminate(java.lang.Runnable afterTerminate)
afterTerminate - Flux.doAfterTerminate(java.lang.Runnable)public final reactor.core.publisher.Flux<T> doOnCancel(java.lang.Runnable onCancel)
onCancel - Flux.doOnCancel(java.lang.Runnable)public final reactor.core.publisher.Flux<T> doOnComplete(java.lang.Runnable onComplete)
onComplete - Flux.doOnComplete(java.lang.Runnable)public final reactor.core.publisher.Flux<T> doOnError(java.util.function.Consumer<? super java.lang.Throwable> onError)
onError - Flux.doOnError(java.util.function.Consumer)public final <E extends java.lang.Throwable> reactor.core.publisher.Flux<T> doOnError(java.lang.Class<E> exceptionType, java.util.function.Consumer<? super E> onError)
exceptionType - onError - Flux.doOnError(java.lang.Class, java.util.function.Consumer)public final reactor.core.publisher.Flux<T> doOnError(java.util.function.Predicate<? super java.lang.Throwable> predicate, java.util.function.Consumer<? super java.lang.Throwable> onError)
predicate - onError - Flux.doOnError(java.util.function.Predicate, java.util.function.Consumer)public final reactor.core.publisher.Flux<T> doOnNext(java.util.function.Consumer<? super T> onNext)
onNext - Flux.doOnNext(java.util.function.Consumer)public final reactor.core.publisher.Flux<T> doOnRequest(java.util.function.LongConsumer consumer)
consumer - Flux.doOnRequest(java.util.function.LongConsumer)public final reactor.core.publisher.Flux<T> doOnSubscribe(java.util.function.Consumer<? super org.reactivestreams.Subscription> onSubscribe)
onSubscribe - Flux.doOnSubscribe(java.util.function.Consumer)public final reactor.core.publisher.Flux<T> doOnTerminate(java.lang.Runnable onTerminate)
onTerminate - Flux.doOnTerminate(java.lang.Runnable)public final reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> elapsed()
Flux.elapsed()public final reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> elapsed(reactor.core.scheduler.TimedScheduler scheduler)
scheduler - Flux.elapsed(reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Mono<T> elementAt(int index)
index - Flux.elementAt(int)public final reactor.core.publisher.Mono<T> elementAt(int index, T defaultValue)
index - defaultValue - Flux.elementAt(int, java.lang.Object)public final reactor.core.publisher.Flux<T> filter(java.util.function.Predicate<? super T> p)
p - Flux.filter(java.util.function.Predicate)public final reactor.core.publisher.Flux<T> firstEmittingWith(org.reactivestreams.Publisher<? extends T> other)
other - Flux.firstEmittingWith(org.reactivestreams.Publisher)public final <R> reactor.core.publisher.Flux<R> flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends R>> mapper)
mapper - Flux.flatMap(java.util.function.Function)public final <V> reactor.core.publisher.Flux<V> flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper, int concurrency)
mapper - concurrency - Flux.flatMap(java.util.function.Function, int)public final <V> reactor.core.publisher.Flux<V> flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper, int concurrency, int prefetch)
mapper - concurrency - prefetch - Flux.flatMap(java.util.function.Function, int, int)public final <V> reactor.core.publisher.Flux<V> flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends V>> mapper, boolean delayError, int concurrency, int prefetch)
mapper - delayError - concurrency - prefetch - Flux.flatMap(java.util.function.Function, boolean, int, int)public final <R> reactor.core.publisher.Flux<R> flatMap(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<? extends R>> mapperOnNext, java.util.function.Function<java.lang.Throwable,? extends org.reactivestreams.Publisher<? extends R>> mapperOnError, java.util.function.Supplier<? extends org.reactivestreams.Publisher<? extends R>> mapperOnComplete)
mapperOnNext - mapperOnError - mapperOnComplete - Flux.flatMap(java.util.function.Function, java.util.function.Function, java.util.function.Supplier)public final <R> reactor.core.publisher.Flux<R> flatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper)
mapper - Flux.flatMapIterable(java.util.function.Function)public final <R> reactor.core.publisher.Flux<R> flatMapIterable(java.util.function.Function<? super T,? extends java.lang.Iterable<? extends R>> mapper, int prefetch)
mapper - prefetch - Flux.flatMapIterable(java.util.function.Function, int)public long getPrefetch()
Flux.getPrefetch()public final <K> reactor.core.publisher.Flux<reactor.core.publisher.GroupedFlux<K,T>> groupBy(java.util.function.Function<? super T,? extends K> keyMapper)
keyMapper - Flux.groupBy(java.util.function.Function)public final <K,V> reactor.core.publisher.Flux<reactor.core.publisher.GroupedFlux<K,V>> groupBy(java.util.function.Function<? super T,? extends K> keyMapper, java.util.function.Function<? super T,? extends V> valueMapper)
keyMapper - valueMapper - Flux.groupBy(java.util.function.Function, java.util.function.Function)public final <TRight,TLeftEnd,TRightEnd,R> reactor.core.publisher.Flux<R> groupJoin(org.reactivestreams.Publisher<? extends TRight> other,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<TLeftEnd>> leftEnd,
java.util.function.Function<? super TRight,? extends org.reactivestreams.Publisher<TRightEnd>> rightEnd,
java.util.function.BiFunction<? super T,? super reactor.core.publisher.Flux<TRight>,? extends R> resultSelector)
other - leftEnd - rightEnd - resultSelector - Flux.groupJoin(org.reactivestreams.Publisher, java.util.function.Function, java.util.function.Function, java.util.function.BiFunction)public final <R> reactor.core.publisher.Flux<R> handle(java.util.function.BiConsumer<? super T,reactor.core.publisher.SynchronousSink<R>> handler)
handler - Flux.handle(java.util.function.BiConsumer)public final reactor.core.publisher.Mono<java.lang.Boolean> hasElement(T value)
value - Flux.hasElement(java.lang.Object)public final reactor.core.publisher.Mono<java.lang.Boolean> hasElements()
Flux.hasElements()public final reactor.core.publisher.Flux<T> hide()
Flux.hide()public final reactor.core.publisher.Mono<T> ignoreElements()
Flux.ignoreElements()public final <TRight,TLeftEnd,TRightEnd,R> reactor.core.publisher.Flux<R> join(org.reactivestreams.Publisher<? extends TRight> other,
java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<TLeftEnd>> leftEnd,
java.util.function.Function<? super TRight,? extends org.reactivestreams.Publisher<TRightEnd>> rightEnd,
java.util.function.BiFunction<? super T,? super TRight,? extends R> resultSelector)
other - leftEnd - rightEnd - resultSelector - Flux.join(org.reactivestreams.Publisher, java.util.function.Function, java.util.function.Function, java.util.function.BiFunction)public final reactor.core.publisher.Mono<T> last()
Flux.last()public final reactor.core.publisher.Mono<T> last(T defaultValue)
defaultValue - Flux.last(java.lang.Object)public final reactor.core.publisher.Flux<T> log()
Flux.log()public final reactor.core.publisher.Flux<T> log(java.lang.String category)
category - Flux.log(java.lang.String)public final reactor.core.publisher.Flux<T> log(java.lang.String category, java.util.logging.Level level, reactor.core.publisher.SignalType... options)
category - level - options - Flux.log(java.lang.String, java.util.logging.Level, reactor.core.publisher.SignalType[])public final reactor.core.publisher.Flux<T> log(java.lang.String category, java.util.logging.Level level, boolean showOperatorLine, reactor.core.publisher.SignalType... options)
category - level - showOperatorLine - options - Flux.log(java.lang.String, java.util.logging.Level, boolean, reactor.core.publisher.SignalType[])public final <V> reactor.core.publisher.Flux<V> map(java.util.function.Function<? super T,? extends V> mapper)
mapper - Flux.map(java.util.function.Function)public final reactor.core.publisher.Flux<T> mapError(java.util.function.Function<? super java.lang.Throwable,? extends java.lang.Throwable> mapper)
mapper - Flux.mapError(java.util.function.Function)public final <E extends java.lang.Throwable> reactor.core.publisher.Flux<T> mapError(java.lang.Class<E> type, java.util.function.Function<? super E,? extends java.lang.Throwable> mapper)
type - mapper - Flux.mapError(java.lang.Class, java.util.function.Function)public final reactor.core.publisher.Flux<T> mapError(java.util.function.Predicate<? super java.lang.Throwable> predicate, java.util.function.Function<? super java.lang.Throwable,? extends java.lang.Throwable> mapper)
predicate - mapper - Flux.mapError(java.util.function.Predicate, java.util.function.Function)public final reactor.core.publisher.Flux<reactor.core.publisher.Signal<T>> materialize()
Flux.materialize()public final reactor.core.publisher.Flux<T> mergeWith(org.reactivestreams.Publisher<? extends T> other)
other - Flux.mergeWith(org.reactivestreams.Publisher)public final reactor.core.publisher.Mono<T> next()
Flux.next()public final <U> reactor.core.publisher.Flux<U> ofType(java.lang.Class<U> clazz)
clazz - Flux.ofType(java.lang.Class)public final reactor.core.publisher.Flux<T> onBackpressureBuffer()
Flux.onBackpressureBuffer()public final reactor.core.publisher.Flux<T> onBackpressureBuffer(int maxSize)
maxSize - Flux.onBackpressureBuffer(int)public final reactor.core.publisher.Flux<T> onBackpressureBuffer(int maxSize, java.util.function.Consumer<? super T> onOverflow)
maxSize - onOverflow - Flux.onBackpressureBuffer(int, java.util.function.Consumer)public final reactor.core.publisher.Flux<T> onBackpressureDrop()
Flux.onBackpressureDrop()public final reactor.core.publisher.Flux<T> onBackpressureDrop(java.util.function.Consumer<? super T> onDropped)
onDropped - Flux.onBackpressureDrop(java.util.function.Consumer)public final reactor.core.publisher.Flux<T> onBackpressureError()
Flux.onBackpressureError()public final reactor.core.publisher.Flux<T> onBackpressureLatest()
Flux.onBackpressureLatest()public final reactor.core.publisher.Flux<T> onErrorResumeWith(java.util.function.Function<? super java.lang.Throwable,? extends org.reactivestreams.Publisher<? extends T>> fallback)
fallback - Flux.onErrorResumeWith(java.util.function.Function)public final <E extends java.lang.Throwable> reactor.core.publisher.Flux<T> onErrorResumeWith(java.lang.Class<E> type, java.util.function.Function<? super E,? extends org.reactivestreams.Publisher<? extends T>> fallback)
type - fallback - Flux.onErrorResumeWith(java.lang.Class, java.util.function.Function)public final reactor.core.publisher.Flux<T> onErrorResumeWith(java.util.function.Predicate<? super java.lang.Throwable> predicate, java.util.function.Function<? super java.lang.Throwable,? extends org.reactivestreams.Publisher<? extends T>> fallback)
predicate - fallback - Flux.onErrorResumeWith(java.util.function.Predicate, java.util.function.Function)public final reactor.core.publisher.Flux<T> onErrorReturn(T fallbackValue)
fallbackValue - Flux.onErrorReturn(java.lang.Object)public final <E extends java.lang.Throwable> reactor.core.publisher.Flux<T> onErrorReturn(java.lang.Class<E> type, T fallbackValue)
type - fallbackValue - Flux.onErrorReturn(java.lang.Class, java.lang.Object)public final <E extends java.lang.Throwable> reactor.core.publisher.Flux<T> onErrorReturn(java.util.function.Predicate<? super java.lang.Throwable> predicate, T fallbackValue)
predicate - fallbackValue - Flux.onErrorReturn(java.util.function.Predicate, java.lang.Object)public final reactor.core.publisher.Flux<T> onTerminateDetach()
Flux.onTerminateDetach()public final reactor.core.publisher.ParallelFlux<T> parallel()
Flux.parallel()public final reactor.core.publisher.ParallelFlux<T> parallel(int parallelism)
parallelism - Flux.parallel(int)public final reactor.core.publisher.ParallelFlux<T> parallel(int parallelism, int prefetch)
parallelism - prefetch - Flux.parallel(int, int)public final reactor.core.publisher.ConnectableFlux<T> publish()
Flux.publish()public final reactor.core.publisher.ConnectableFlux<T> publish(int prefetch)
prefetch - Flux.publish(int)public final <R> reactor.core.publisher.Flux<R> publish(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<? extends R>> transform)
transform - Flux.publish(java.util.function.Function)public final <R> reactor.core.publisher.Flux<R> publish(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<? extends R>> transform, int prefetch)
transform - prefetch - Flux.publish(java.util.function.Function, int)public final reactor.core.publisher.Mono<T> publishNext()
Flux.publishNext()public final reactor.core.publisher.Flux<T> publishOn(reactor.core.scheduler.Scheduler scheduler)
scheduler - Flux.publishOn(reactor.core.scheduler.Scheduler)public final reactor.core.publisher.Flux<T> publishOn(reactor.core.scheduler.Scheduler scheduler, int prefetch)
scheduler - prefetch - Flux.publishOn(reactor.core.scheduler.Scheduler, int)public final reactor.core.publisher.Mono<T> reduce(java.util.function.BiFunction<T,T,T> aggregator)
aggregator - Flux.reduce(java.util.function.BiFunction)public final <A> reactor.core.publisher.Mono<A> reduce(A initial,
java.util.function.BiFunction<A,? super T,A> accumulator)
initial - accumulator - Flux.reduce(java.lang.Object, java.util.function.BiFunction)public final <A> reactor.core.publisher.Mono<A> reduceWith(java.util.function.Supplier<A> initial,
java.util.function.BiFunction<A,? super T,A> accumulator)
initial - accumulator - Flux.reduceWith(java.util.function.Supplier, java.util.function.BiFunction)public final reactor.core.publisher.Flux<T> repeat()
Flux.repeat()public final reactor.core.publisher.Flux<T> repeat(java.util.function.BooleanSupplier predicate)
predicate - Flux.repeat(java.util.function.BooleanSupplier)public final reactor.core.publisher.Flux<T> repeat(long numRepeat)
numRepeat - Flux.repeat(long)public final reactor.core.publisher.Flux<T> repeat(long numRepeat, java.util.function.BooleanSupplier predicate)
numRepeat - predicate - Flux.repeat(long, java.util.function.BooleanSupplier)public final reactor.core.publisher.Flux<T> repeatWhen(java.util.function.Function<reactor.core.publisher.Flux<java.lang.Long>,? extends org.reactivestreams.Publisher<?>> whenFactory)
whenFactory - Flux.repeatWhen(java.util.function.Function)public final reactor.core.publisher.ConnectableFlux<T> replay()
Flux.replay()public final reactor.core.publisher.ConnectableFlux<T> replay(int history)
history - Flux.replay(int)public final reactor.core.publisher.ConnectableFlux<T> replay(java.time.Duration ttl)
ttl - Flux.replay(java.time.Duration)public final reactor.core.publisher.ConnectableFlux<T> replay(int history, java.time.Duration ttl)
history - ttl - Flux.replay(int, java.time.Duration)public final reactor.core.publisher.ConnectableFlux<T> replayMillis(long ttl, reactor.core.scheduler.TimedScheduler timer)
ttl - timer - Flux.replayMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.ConnectableFlux<T> replayMillis(int history, long ttl, reactor.core.scheduler.TimedScheduler timer)
history - ttl - timer - Flux.replayMillis(int, long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<T> retry()
Flux.retry()public final reactor.core.publisher.Flux<T> retry(long numRetries)
numRetries - Flux.retry(long)public final reactor.core.publisher.Flux<T> retry(java.util.function.Predicate<java.lang.Throwable> retryMatcher)
retryMatcher - Flux.retry(java.util.function.Predicate)public final reactor.core.publisher.Flux<T> retry(long numRetries, java.util.function.Predicate<java.lang.Throwable> retryMatcher)
numRetries - retryMatcher - Flux.retry(long, java.util.function.Predicate)public final reactor.core.publisher.Flux<T> retryWhen(java.util.function.Function<reactor.core.publisher.Flux<java.lang.Throwable>,? extends org.reactivestreams.Publisher<?>> whenFactory)
whenFactory - Flux.retryWhen(java.util.function.Function)public final reactor.core.publisher.Flux<T> sample(java.time.Duration timespan)
timespan - Flux.sample(java.time.Duration)public final <U> reactor.core.publisher.Flux<T> sample(org.reactivestreams.Publisher<U> sampler)
sampler - Flux.sample(org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> sampleFirst(java.time.Duration timespan)
timespan - Flux.sampleFirst(java.time.Duration)public final <U> reactor.core.publisher.Flux<T> sampleFirst(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<U>> samplerFactory)
samplerFactory - Flux.sampleFirst(java.util.function.Function)public final reactor.core.publisher.Flux<T> sampleFirstMillis(long timespan)
timespan - Flux.sampleFirstMillis(long)public final reactor.core.publisher.Flux<T> sampleMillis(long timespan)
timespan - Flux.sampleMillis(long)public final <U> reactor.core.publisher.Flux<T> sampleTimeout(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<U>> throttlerFactory)
throttlerFactory - Flux.sampleTimeout(java.util.function.Function)public final <U> reactor.core.publisher.Flux<T> sampleTimeout(java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<U>> throttlerFactory, int maxConcurrency)
throttlerFactory - maxConcurrency - Flux.sampleTimeout(java.util.function.Function, int)public final reactor.core.publisher.Flux<T> scan(java.util.function.BiFunction<T,T,T> accumulator)
accumulator - Flux.scan(java.util.function.BiFunction)public final <A> reactor.core.publisher.Flux<A> scan(A initial,
java.util.function.BiFunction<A,? super T,A> accumulator)
initial - accumulator - Flux.scan(java.lang.Object, java.util.function.BiFunction)public final <A> reactor.core.publisher.Flux<A> scanWith(java.util.function.Supplier<A> initial,
java.util.function.BiFunction<A,? super T,A> accumulator)
initial - accumulator - Flux.scanWith(java.util.function.Supplier, java.util.function.BiFunction)public final reactor.core.publisher.Flux<T> share()
Flux.share()public final reactor.core.publisher.Mono<T> single()
Flux.single()public final reactor.core.publisher.Mono<T> single(T defaultValue)
defaultValue - Flux.single(java.lang.Object)public final reactor.core.publisher.Mono<T> singleOrEmpty()
Flux.singleOrEmpty()public final reactor.core.publisher.Flux<T> skip(long skipped)
skipped - Flux.skip(long)public final reactor.core.publisher.Flux<T> skip(java.time.Duration timespan)
timespan - Flux.skip(java.time.Duration)public final reactor.core.publisher.Flux<T> skipLast(int n)
n - Flux.skipLast(int)public final reactor.core.publisher.Flux<T> skipMillis(long timespan)
timespan - Flux.skipMillis(long)public final reactor.core.publisher.Flux<T> skipMillis(long timespan, reactor.core.scheduler.TimedScheduler timer)
timespan - timer - Flux.skipMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<T> skipUntil(java.util.function.Predicate<? super T> untilPredicate)
untilPredicate - Flux.skipUntil(java.util.function.Predicate)public final reactor.core.publisher.Flux<T> skipUntilOther(org.reactivestreams.Publisher<?> other)
other - Flux.skipUntilOther(org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> skipWhile(java.util.function.Predicate<? super T> skipPredicate)
skipPredicate - Flux.skipWhile(java.util.function.Predicate)public final reactor.core.publisher.Flux<T> sort()
Flux.sort()public final reactor.core.publisher.Flux<T> sort(java.util.Comparator<? super T> sortFunction)
sortFunction - Flux.sort(java.util.Comparator)public final reactor.core.publisher.Flux<T> startWith(java.lang.Iterable<? extends T> iterable)
iterable - Flux.startWith(java.lang.Iterable)public final reactor.core.publisher.Flux<T> startWith(T... values)
values - Flux.startWith(java.lang.Object[])public final reactor.core.publisher.Flux<T> startWith(org.reactivestreams.Publisher<? extends T> publisher)
publisher - Flux.startWith(org.reactivestreams.Publisher)public final reactor.core.Cancellation subscribe()
Flux.subscribe()public final reactor.core.Cancellation subscribe(int prefetch)
prefetch - Flux.subscribe(int)public final reactor.core.Cancellation subscribe(java.util.function.Consumer<? super T> consumer)
consumer - Flux.subscribe(java.util.function.Consumer)public final reactor.core.Cancellation subscribe(java.util.function.Consumer<? super T> consumer, int prefetch)
consumer - prefetch - Flux.subscribe(java.util.function.Consumer, int)public final reactor.core.Cancellation subscribe(java.util.function.Consumer<? super T> consumer, java.util.function.Consumer<? super java.lang.Throwable> errorConsumer)
consumer - errorConsumer - Flux.subscribe(java.util.function.Consumer, java.util.function.Consumer)public final reactor.core.Cancellation subscribe(java.util.function.Consumer<? super T> consumer, java.util.function.Consumer<? super java.lang.Throwable> errorConsumer, java.lang.Runnable completeConsumer)
consumer - errorConsumer - completeConsumer - Flux.subscribe(java.util.function.Consumer, java.util.function.Consumer, java.lang.Runnable)public final reactor.core.Cancellation subscribe(java.util.function.Consumer<? super T> consumer, java.util.function.Consumer<? super java.lang.Throwable> errorConsumer, java.lang.Runnable completeConsumer, int prefetch)
consumer - errorConsumer - completeConsumer - prefetch - Flux.subscribe(java.util.function.Consumer, java.util.function.Consumer, java.lang.Runnable, int)public final reactor.core.publisher.Flux<T> subscribeOn(reactor.core.scheduler.Scheduler scheduler)
scheduler - Flux.subscribeOn(reactor.core.scheduler.Scheduler)public final <E extends org.reactivestreams.Subscriber<? super T>> E subscribeWith(E subscriber)
subscriber - Flux.subscribeWith(org.reactivestreams.Subscriber)public final reactor.core.publisher.Flux<T> switchIfEmpty(org.reactivestreams.Publisher<? extends T> alternate)
alternate - Flux.switchIfEmpty(org.reactivestreams.Publisher)public final <V> reactor.core.publisher.Flux<V> switchMap(java.util.function.Function<? super T,org.reactivestreams.Publisher<? extends V>> fn)
fn - Flux.switchMap(java.util.function.Function)public final <V> reactor.core.publisher.Flux<V> switchMap(java.util.function.Function<? super T,org.reactivestreams.Publisher<? extends V>> fn, int prefetch)
fn - prefetch - Flux.switchMap(java.util.function.Function, int)public final <E extends java.lang.Throwable> reactor.core.publisher.Flux<T> switchOnError(java.lang.Class<E> type, org.reactivestreams.Publisher<? extends T> fallback)
type - fallback - Flux.switchOnError(java.lang.Class, org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> switchOnError(java.util.function.Predicate<? super java.lang.Throwable> predicate, org.reactivestreams.Publisher<? extends T> fallback)
predicate - fallback - Flux.switchOnError(java.util.function.Predicate, org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> switchOnError(org.reactivestreams.Publisher<? extends T> fallback)
fallback - Flux.switchOnError(org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> take(long n)
n - Flux.take(long)public final reactor.core.publisher.Flux<T> take(java.time.Duration timespan)
timespan - Flux.take(java.time.Duration)public final reactor.core.publisher.Flux<T> takeLast(int n)
n - Flux.takeLast(int)public final reactor.core.publisher.Flux<T> takeMillis(long timespan)
timespan - Flux.takeMillis(long)public final reactor.core.publisher.Flux<T> takeMillis(long timespan, reactor.core.scheduler.TimedScheduler timer)
timespan - timer - Flux.takeMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<T> takeUntil(java.util.function.Predicate<? super T> predicate)
predicate - Flux.takeUntil(java.util.function.Predicate)public final reactor.core.publisher.Flux<T> takeUntilOther(org.reactivestreams.Publisher<?> other)
other - Flux.takeUntilOther(org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> takeWhile(java.util.function.Predicate<? super T> continuePredicate)
continuePredicate - Flux.takeWhile(java.util.function.Predicate)public final reactor.core.publisher.Mono<java.lang.Void> then()
Flux.then()public final reactor.core.publisher.Mono<java.lang.Void> then(org.reactivestreams.Publisher<java.lang.Void> other)
other - Flux.then(org.reactivestreams.Publisher)public final reactor.core.publisher.Mono<java.lang.Void> then(java.util.function.Supplier<? extends org.reactivestreams.Publisher<java.lang.Void>> afterSupplier)
afterSupplier - Flux.then(java.util.function.Supplier)public final <V> reactor.core.publisher.Flux<V> thenMany(org.reactivestreams.Publisher<V> other)
other - Flux.thenMany(org.reactivestreams.Publisher)public final <V> reactor.core.publisher.Flux<V> thenMany(java.util.function.Supplier<? extends org.reactivestreams.Publisher<V>> afterSupplier)
afterSupplier - Flux.thenMany(java.util.function.Supplier)public final reactor.core.publisher.Flux<T> timeout(java.time.Duration timeout)
timeout - Flux.timeout(java.time.Duration)public final reactor.core.publisher.Flux<T> timeout(java.time.Duration timeout, org.reactivestreams.Publisher<? extends T> fallback)
timeout - fallback - Flux.timeout(java.time.Duration, org.reactivestreams.Publisher)public final <U> reactor.core.publisher.Flux<T> timeout(org.reactivestreams.Publisher<U> firstTimeout)
firstTimeout - Flux.timeout(org.reactivestreams.Publisher)public final <U,V> reactor.core.publisher.Flux<T> timeout(org.reactivestreams.Publisher<U> firstTimeout, java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<V>> nextTimeoutFactory)
firstTimeout - nextTimeoutFactory - Flux.timeout(org.reactivestreams.Publisher, java.util.function.Function)public final <U,V> reactor.core.publisher.Flux<T> timeout(org.reactivestreams.Publisher<U> firstTimeout, java.util.function.Function<? super T,? extends org.reactivestreams.Publisher<V>> nextTimeoutFactory, org.reactivestreams.Publisher<? extends T> fallback)
firstTimeout - nextTimeoutFactory - fallback - Flux.timeout(org.reactivestreams.Publisher, java.util.function.Function, org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> timeoutMillis(long timeout)
timeout - Flux.timeoutMillis(long)public final reactor.core.publisher.Flux<T> timeoutMillis(long timeout, reactor.core.scheduler.TimedScheduler timer)
timeout - timer - Flux.timeoutMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<T> timeoutMillis(long timeout, org.reactivestreams.Publisher<? extends T> fallback)
timeout - fallback - Flux.timeoutMillis(long, org.reactivestreams.Publisher)public final reactor.core.publisher.Flux<T> timeoutMillis(long timeout, org.reactivestreams.Publisher<? extends T> fallback, reactor.core.scheduler.TimedScheduler timer)
timeout - fallback - timer - Flux.timeoutMillis(long, org.reactivestreams.Publisher, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> timestamp()
Flux.timestamp()public final reactor.core.publisher.Flux<reactor.util.function.Tuple2<java.lang.Long,T>> timestamp(reactor.core.scheduler.TimedScheduler scheduler)
scheduler - Flux.timestamp(reactor.core.scheduler.TimedScheduler)public final java.lang.Iterable<T> toIterable()
Flux.toIterable()public final java.lang.Iterable<T> toIterable(long batchSize)
batchSize - Flux.toIterable(long)public final java.lang.Iterable<T> toIterable(long batchSize, java.util.function.Supplier<java.util.Queue<T>> queueProvider)
batchSize - queueProvider - Flux.toIterable(long, java.util.function.Supplier)public java.util.stream.Stream<T> toStream()
Flux.toStream()public java.util.stream.Stream<T> toStream(int batchSize)
batchSize - Flux.toStream(int)public final <V> reactor.core.publisher.Flux<V> transform(java.util.function.Function<? super reactor.core.publisher.Flux<T>,? extends org.reactivestreams.Publisher<V>> transformer)
transformer - Flux.transform(java.util.function.Function)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window()
Flux.window()public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(int maxSize)
maxSize - Flux.window(int)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(int maxSize, int skip)
maxSize - skip - Flux.window(int, int)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(org.reactivestreams.Publisher<?> boundary)
boundary - Flux.window(org.reactivestreams.Publisher)public final <U,V> reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(org.reactivestreams.Publisher<U> bucketOpening, java.util.function.Function<? super U,? extends org.reactivestreams.Publisher<V>> closeSelector)
bucketOpening - closeSelector - Flux.window(org.reactivestreams.Publisher, java.util.function.Function)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(java.time.Duration timespan)
timespan - Flux.window(java.time.Duration)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(java.time.Duration timespan, java.time.Duration timeshift)
timespan - timeshift - Flux.window(java.time.Duration, java.time.Duration)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> window(int maxSize, java.time.Duration timespan)
maxSize - timespan - Flux.window(int, java.time.Duration)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> windowMillis(long timespan)
timespan - Flux.windowMillis(long)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> windowMillis(long timespan, reactor.core.scheduler.TimedScheduler timer)
timespan - timer - Flux.windowMillis(long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> windowMillis(long timespan, long timeshift, reactor.core.scheduler.TimedScheduler timer)
timespan - timeshift - timer - Flux.windowMillis(long, long, reactor.core.scheduler.TimedScheduler)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> windowMillis(int maxSize, long timespan)
maxSize - timespan - Flux.windowMillis(int, long)public final reactor.core.publisher.Flux<reactor.core.publisher.Flux<T>> windowMillis(int maxSize, long timespan, reactor.core.scheduler.TimedScheduler timer)
maxSize - timespan - timer - Flux.windowMillis(int, long, reactor.core.scheduler.TimedScheduler)public final <U,R> reactor.core.publisher.Flux<R> withLatestFrom(org.reactivestreams.Publisher<? extends U> other,
java.util.function.BiFunction<? super T,? super U,? extends R> resultSelector)
other - resultSelector - Flux.withLatestFrom(org.reactivestreams.Publisher, java.util.function.BiFunction)public final <T2,V> reactor.core.publisher.Flux<V> zipWith(org.reactivestreams.Publisher<? extends T2> source2,
java.util.function.BiFunction<? super T,? super T2,? extends V> combinator)
source2 - combinator - Flux.zipWith(org.reactivestreams.Publisher, java.util.function.BiFunction)public final <T2,V> reactor.core.publisher.Flux<V> zipWith(org.reactivestreams.Publisher<? extends T2> source2,
int prefetch,
java.util.function.BiFunction<? super T,? super T2,? extends V> combinator)
source2 - prefetch - combinator - Flux.zipWith(org.reactivestreams.Publisher, int, java.util.function.BiFunction)public final <T2> reactor.core.publisher.Flux<reactor.util.function.Tuple2<T,T2>> zipWith(org.reactivestreams.Publisher<? extends T2> source2)
source2 - Flux.zipWith(org.reactivestreams.Publisher)public final <T2> reactor.core.publisher.Flux<reactor.util.function.Tuple2<T,T2>> zipWith(org.reactivestreams.Publisher<? extends T2> source2, int prefetch)
source2 - prefetch - Flux.zipWith(org.reactivestreams.Publisher, int)public final <T2> reactor.core.publisher.Flux<reactor.util.function.Tuple2<T,T2>> zipWithIterable(java.lang.Iterable<? extends T2> iterable)
iterable - Flux.zipWithIterable(java.lang.Iterable)public final <T2,V> reactor.core.publisher.Flux<V> zipWithIterable(java.lang.Iterable<? extends T2> iterable,
java.util.function.BiFunction<? super T,? super T2,? extends V> zipper)
iterable - zipper - Flux.zipWithIterable(java.lang.Iterable, java.util.function.BiFunction)