4.x: Streamable operator checklist
Nobody has claimed this yet.
- Dominant language
- Java
- Stars
- 48.2k
- Forks
- 7.6k
- Avg merge
- 39m
- Merged PRs (30d)
- 11
Description
Via Claude Fable 5 High
Streamable: missing Flowable operators and overloads (4.x parity checklist)
This is a mechanical comparison of the public API of Flowable against Streamable on the 4.x branch (as of commit bc5d8b8). Flowable currently exposes 464 public methods/overloads (206 distinct names); Streamable exposes 75 (55 distinct names).
Signatures below are written with the Streamable equivalents of the Flowable types substituted in:
Flowable type |
Streamable equivalent used below |
|---|---|
Flowable<T> / Publisher<T> (parameters, return types) |
Streamable<T> |
GroupedFlowable<K, V> |
GroupedStreamable<K, V> |
FlowableOperator<R, T> |
StreamableOperator<T, R> |
FlowableConverter<T, R> |
StreamableConverter<T, R> |
FlowableTransformer<T, R> |
StreamableTransformer<T, R> (type does not exist yet) |
ConnectableFlowable<T> |
ConnectableStreamable<T> (type does not exist yet) |
FlowableEmitter<T> |
VirtualEmitter<T> |
Subscriber<T> / FlowableSubscriber<T> |
Flow.Subscriber<T> |
Things to keep in mind while working through the list:
Flowable'sDisposable-returningsubscribe(...)overloads probably map toCompletionStageDisposable<Void>onStreamable(asforEachalready does).- Overloads that take a
Schedulermay want an additionalExecutorServicetwin, following the existingStreamable.timer/intervalRange/createpattern. Those twins are not listed here because they have noFlowablecounterpart. - A few methods already exist on
Streamablewith a different shape (e.g.flatMap(mapper, StandardConcurrentConfig)vs.Flowable'sStandardConcurrentBufferedConfig,retry(BiPredicate<Long, Throwable>)vs.BiPredicate<Integer, Throwable>,repeatWhen/retryWhentakingCompletionStage<Boolean>functions,create(VirtualGenerator)vs.virtualCreate). These are listed in the missing overloads section only where theFlowableoverload has no exact-arity/equivalent-type match; adjust or tick them off as appropriate. - Return-type differences on otherwise matching signatures are not listed (e.g.
ignoreElements()returnsCompletableonFlowablebutStreamable<T>onStreamable;collect(Collector)returnsSingle<R>vs.Streamable<R>;forEachreturnsDisposablevs.CompletionStageDisposable<Void>).
Methods that exist on Streamable but are missing overloads (23 methods, 69 overloads)
blockingFirst
-
T blockingFirst(T defaultItem)
blockingLast
-
T blockingLast(T defaultItem)
collect
-
<U> Single<U> collect(Supplier<? extends U> initialItemSupplier, BiConsumer<? super U, ? super T> collector)
concat
-
static <T> Streamable<T> concat(Streamable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> concat(Iterable<? extends Streamable<? extends T>> sources, StandardBufferedConfig config) -
static <T> Streamable<T> concat(Streamable<? extends Streamable<? extends T>> sources, StandardBufferedConfig config)
delay
-
<U> Streamable<T> delay(Function<? super T, ? extends Streamable<U>> itemDelayIndicator) -
Streamable<T> delay(long time, TimeUnit unit) -
<U, V> Streamable<T> delay(Streamable<U> subscriptionIndicator, Function<? super T, ? extends Streamable<V>> itemDelayIndicator) -
Streamable<T> delay(long time, TimeUnit unit, Scheduler scheduler, boolean delayError)
error
-
static <T> Streamable<T> error(Supplier<? extends Throwable> supplier)
flatMap
-
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardConcurrentBufferedConfig config) -
<U, R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends R> combiner) -
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> onNextMapper, Function<? super Throwable, ? extends Streamable<? extends R>> onErrorMapper, Supplier<? extends Streamable<? extends R>> onCompleteSupplier) -
<U, R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends R> combiner, StandardConcurrentBufferedConfig config) -
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> onNextMapper, Function<Throwable, ? extends Streamable<? extends R>> onErrorMapper, Supplier<? extends Streamable<? extends R>> onCompleteSupplier, StandardConcurrentBufferedConfig config)
groupBy
-
<K> Streamable<GroupedStreamable<K, T>> groupBy(Function<? super T, ? extends K> keySelector, StandardBufferedConfig config) -
<K, V> Streamable<GroupedStreamable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) -
<K, V> Streamable<GroupedStreamable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, StandardBufferedConfig config) -
<K, V> Streamable<GroupedStreamable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Function<? super Consumer<Object>, ? extends Map<K, Object>> evictingMapFactory, StandardBufferedConfig config)
intervalRange
-
static Streamable<Long> intervalRange(long start, long count, long initialDelay, long period, TimeUnit unit)
just
-
static <T> Streamable<T> just(T item1, T item2) -
static <T> Streamable<T> just(T item1, T item2, T item3) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8, T item9) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8, T item9, T item10)
repeat
-
Streamable<T> repeat()
repeatWhen
-
Streamable<T> repeatWhen(Function<? super Streamable<Object>, ? extends Streamable<?>> handler)
retry
-
Streamable<T> retry() -
Streamable<T> retry(BiPredicate<? super Integer, ? super Throwable> predicate) -
Streamable<T> retry(long times, Predicate<? super Throwable> predicate)
retryWhen
-
Streamable<T> retryWhen(Function<? super Streamable<Throwable>, ? extends Streamable<?>> handler)
skip
-
Streamable<T> skip(long time, TimeUnit unit) -
Streamable<T> skip(long time, TimeUnit unit, Scheduler scheduler)
subscribe
-
Disposable subscribe() -
Disposable subscribe(Consumer<? super T> onNext) -
Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError) -
Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete) -
Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete, DisposableContainer container)
take
-
Streamable<T> take(long time, TimeUnit unit) -
Streamable<T> take(long time, TimeUnit unit, Scheduler scheduler)
takeUntil
-
Streamable<T> takeUntil(Predicate<? super T> stopPredicate)
test
-
TestSubscriber<T> test(long initialRequest) -
TestSubscriber<T> test(long initialRequest, boolean cancel)
timeout
-
<V> Streamable<T> timeout(Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator) -
<V> Streamable<T> timeout(Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator, Streamable<? extends T> fallback) -
Streamable<T> timeout(long timeout, TimeUnit unit) -
<U, V> Streamable<T> timeout(Streamable<U> firstTimeoutIndicator, Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator) -
Streamable<T> timeout(long timeout, TimeUnit unit, Streamable<? extends T> fallback) -
Streamable<T> timeout(long timeout, TimeUnit unit, Scheduler scheduler) -
<U, V> Streamable<T> timeout(Streamable<U> firstTimeoutIndicator, Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator, Streamable<? extends T> fallback)
timer
-
static Streamable<Long> timer(long delay, TimeUnit unit)
using
-
static <T, D> Streamable<T> using(Supplier<? extends D> resourceSupplier, Function<? super D, ? extends Streamable<? extends T>> sourceSupplier, Consumer<? super D> resourceCleanup, boolean eager)
zip
-
static <T, R> Streamable<R> zip(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> zipper) -
static <T, R> Streamable<R> zip(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> zipper, StandardBufferedConfig config) -
static <T1, T2, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, BiFunction<? super T1, ? super T2, ? extends R> zipper) -
static <T1, T2, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, BiFunction<? super T1, ? super T2, ? extends R> zipper, StandardBufferedConfig config) -
static <T1, T2, T3, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Function3<? super T1, ? super T2, ? super T3, ? extends R> zipper) -
static <T1, T2, T3, T4, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Function4<? super T1, ? super T2, ? super T3, ? super T4, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Function5<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Function6<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, T7, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Function7<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, T7, T8, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Function8<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, T7, T8, T9, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Streamable<? extends T9> source9, Function9<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? super T9, ? extends R> zipper)
Methods that are entirely missing from Streamable (140 methods, 317 overloads)
all
-
Single<Boolean> all(Predicate<? super T> predicate)
amb
-
static <T> Streamable<T> amb(Iterable<? extends Streamable<? extends T>> sources)
ambArray
-
static <T> Streamable<T> ambArray(Streamable<? extends T>... sources)
ambWith
-
Streamable<T> ambWith(Streamable<? extends T> other)
any
-
Single<Boolean> any(Predicate<? super T> predicate)
blockingForEach
-
void blockingForEach(Consumer<? super T> onNext) -
void blockingForEach(Consumer<? super T> onNext, int bufferSize)
blockingIterable
-
Iterable<T> blockingIterable() -
Iterable<T> blockingIterable(int bufferSize)
blockingLatest
-
Iterable<T> blockingLatest()
blockingMostRecent
-
Iterable<T> blockingMostRecent(T initialItem)
blockingNext
-
Iterable<T> blockingNext()
blockingSingle
-
T blockingSingle() -
T blockingSingle(T defaultItem)
blockingStream
-
Stream<T> blockingStream() -
Stream<T> blockingStream(int prefetch)
blockingSubscribe
-
void blockingSubscribe() -
void blockingSubscribe(Consumer<? super T> onNext) -
void blockingSubscribe(Subscriber<? super T> subscriber) -
void blockingSubscribe(Consumer<? super T> onNext, int bufferSize) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, int bufferSize) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete, int bufferSize)
buffer
-
Streamable<List<T>> buffer(int count) -
<B> Streamable<List<T>> buffer(Streamable<B> boundaryIndicator) -
Streamable<List<T>> buffer(int count, int skip) -
<U extends Collection<? super T>> Streamable<U> buffer(int count, Supplier<U> bufferSupplier) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit) -
<TOpening, TClosing> Streamable<List<T>> buffer(Streamable<? extends TOpening> openingIndicator, Function<? super TOpening, ? extends Streamable<? extends TClosing>> closingIndicator) -
<B> Streamable<List<T>> buffer(Streamable<B> boundaryIndicator, int initialCapacity) -
<B, U extends Collection<? super T>> Streamable<U> buffer(Streamable<B> boundaryIndicator, Supplier<U> bufferSupplier) -
<U extends Collection<? super T>> Streamable<U> buffer(int count, int skip, Supplier<U> bufferSupplier) -
Streamable<List<T>> buffer(long timespan, long timeskip, TimeUnit unit) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit, int count) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit, Scheduler scheduler) -
<TOpening, TClosing, U extends Collection<? super T>> Streamable<U> buffer(Streamable<? extends TOpening> openingIndicator, Function<? super TOpening, ? extends Streamable<? extends TClosing>> closingIndicator, Supplier<U> bufferSupplier) -
Streamable<List<T>> buffer(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit, Scheduler scheduler, int count) -
<U extends Collection<? super T>> Streamable<U> buffer(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler, Supplier<U> bufferSupplier) -
<U extends Collection<? super T>> Streamable<U> buffer(long timespan, TimeUnit unit, Scheduler scheduler, int count, Supplier<U> bufferSupplier, boolean restartTimerOnMaxSize)
cache
-
Streamable<T> cache()
cacheWithInitialCapacity
-
Streamable<T> cacheWithInitialCapacity(int initialCapacity)
cast
-
<U> Streamable<U> cast(Class<U> clazz)
collectInto
-
<U> Single<U> collectInto(U initialItem, BiConsumer<? super U, ? super T> collector)
combineLatest
-
static <T, R> Streamable<R> combineLatest(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> combiner) -
static <T, R> Streamable<R> combineLatest(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> combiner, StandardBufferedConfig config) -
static <T1, T2, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, BiFunction<? super T1, ? super T2, ? extends R> combiner) -
static <T1, T2, T3, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Function3<? super T1, ? super T2, ? super T3, ? extends R> combiner) -
static <T1, T2, T3, T4, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Function4<? super T1, ? super T2, ? super T3, ? super T4, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Function5<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Function6<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, T7, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Function7<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, T7, T8, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Function8<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, T7, T8, T9, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Streamable<? extends T9> source9, Function9<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? super T9, ? extends R> combiner)
combineLatestArray
-
static <T, R> Streamable<R> combineLatestArray(Streamable<? extends T>[] sources, Function<? super Object[], ? extends R> combiner) -
static <T, R> Streamable<R> combineLatestArray(Streamable<? extends T>[] sources, Function<? super Object[], ? extends R> combiner, StandardBufferedConfig config)
compose
-
<R> Streamable<R> compose(StreamableTransformer<? super T, ? extends R> composer)
concatArray
-
static <T> Streamable<T> concatArray(Streamable<? extends T>... sources) -
static <T> Streamable<T> concatArray(StandardBufferedConfig config, Streamable<? extends T>... sources)
concatArrayEager
-
static <T> Streamable<T> concatArrayEager(Streamable<? extends T>... sources)
concatEager
-
static <T> Streamable<T> concatEager(Iterable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> concatEager(Streamable<? extends Streamable<? extends T>> sources)
concatMap
-
<R> Streamable<R> concatMap(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> concatMap(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardBufferedConfig config) -
<R> Streamable<R> concatMap(Function<? super T, ? extends Streamable<? extends R>> mapper, Scheduler scheduler, StandardBufferedConfig config)
concatMapCompletable
-
Completable concatMapCompletable(Function<? super T, ? extends CompletableSource> mapper) -
Completable concatMapCompletable(Function<? super T, ? extends CompletableSource> mapper, StandardBufferedConfig config)
concatMapEager
-
<R> Streamable<R> concatMapEager(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> concatMapEager(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardConcurrentBufferedConfig config)
concatMapIterable
-
<U> Streamable<U> concatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper) -
<U> Streamable<U> concatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, StandardBufferedConfig config)
concatMapMaybe
-
<R> Streamable<R> concatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper) -
<R> Streamable<R> concatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper, StandardBufferedConfig config)
concatMapSingle
-
<R> Streamable<R> concatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper) -
<R> Streamable<R> concatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper, StandardBufferedConfig config)
concatMapStream
-
<R> Streamable<R> concatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper) -
<R> Streamable<R> concatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper, StandardBufferedConfig config)
concatWith
-
Streamable<T> concatWith(Streamable<? extends T> other) -
Streamable<T> concatWith(SingleSource<? extends T> other) -
Streamable<T> concatWith(MaybeSource<? extends T> other) -
Streamable<T> concatWith(CompletableSource other)
contains
-
Single<Boolean> contains(Object item)
count
-
Single<Long> count()
debounce
-
<U> Streamable<T> debounce(Function<? super T, ? extends Streamable<U>> debounceIndicator) -
Streamable<T> debounce(long timeout, TimeUnit unit) -
Streamable<T> debounce(long timeout, TimeUnit unit, Scheduler scheduler) -
Streamable<T> debounce(long timeout, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
defaultIfEmpty
-
Streamable<T> defaultIfEmpty(T defaultItem)
delaySubscription
-
<U> Streamable<T> delaySubscription(Streamable<U> subscriptionIndicator) -
Streamable<T> delaySubscription(long time, TimeUnit unit) -
Streamable<T> delaySubscription(long time, TimeUnit unit, Scheduler scheduler)
dematerialize
-
<R> Streamable<R> dematerialize(Function<? super T, Notification<R>> selector)
distinct
-
Streamable<T> distinct() -
<K> Streamable<T> distinct(Function<? super T, K> keySelector) -
<K> Streamable<T> distinct(Function<? super T, K> keySelector, Supplier<? extends Collection<? super K>> collectionSupplier)
distinctUntilChanged
-
Streamable<T> distinctUntilChanged() -
<K> Streamable<T> distinctUntilChanged(Function<? super T, K> keySelector) -
Streamable<T> distinctUntilChanged(BiPredicate<? super T, ? super T> comparer)
doAfterNext
-
Streamable<T> doAfterNext(Consumer<? super T> onAfterNext)
doAfterTerminate
-
Streamable<T> doAfterTerminate(Action onAfterTerminate)
doFinally
-
Streamable<T> doFinally(Action onFinally)
doOnCancel
-
Streamable<T> doOnCancel(Action onCancel)
doOnComplete
-
Streamable<T> doOnComplete(Action onComplete)
doOnEach
-
Streamable<T> doOnEach(Consumer<? super Notification<T>> onNotification) -
Streamable<T> doOnEach(Subscriber<? super T> subscriber)
doOnLifecycle
-
Streamable<T> doOnLifecycle(Consumer<? super Subscription> onSubscribe, LongConsumer onRequest, Action onCancel)
doOnSubscribe
-
Streamable<T> doOnSubscribe(Consumer<? super Subscription> onSubscribe)
doOnTerminate
-
Streamable<T> doOnTerminate(Action onTerminate)
elementAt
-
Maybe<T> elementAt(long index) -
Single<T> elementAt(long index, T defaultItem)
elementAtOrError
-
Single<T> elementAtOrError(long index)
first
-
Single<T> first(T defaultItem)
firstElement
-
Maybe<T> firstElement()
firstOrError
-
Single<T> firstOrError()
firstOrErrorStage
-
CompletionStage<T> firstOrErrorStage()
firstStage
-
CompletionStage<T> firstStage(T defaultItem)
flatMapCompletable
-
Completable flatMapCompletable(Function<? super T, ? extends CompletableSource> mapper) -
Completable flatMapCompletable(Function<? super T, ? extends CompletableSource> mapper, StandardConcurrentConfig config)
flatMapIterable
-
<U> Streamable<U> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper) -
<U> Streamable<U> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, StandardBufferedConfig config) -
<U, V> Streamable<V> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends V> combiner) -
<U, V> Streamable<V> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends V> combiner, StandardConcurrentBufferedConfig config)
flatMapMaybe
-
<R> Streamable<R> flatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper) -
<R> Streamable<R> flatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper, StandardConcurrentConfig config)
flatMapSingle
-
<R> Streamable<R> flatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper) -
<R> Streamable<R> flatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper, StandardConcurrentConfig config)
flatMapStream
-
<R> Streamable<R> flatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper) -
<R> Streamable<R> flatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper, StandardBufferedConfig config)
forEachWhile
-
Disposable forEachWhile(Predicate<? super T> onNext) -
Disposable forEachWhile(Predicate<? super T> onNext, Consumer<? super Throwable> onError) -
Disposable forEachWhile(Predicate<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete)
fromAction
-
static <T> Streamable<T> fromAction(Action action)
fromCallable
-
static <T> Streamable<T> fromCallable(Callable<? extends T> callable)
fromCompletionStage
-
static <T> Streamable<T> fromCompletionStage(CompletionStage<T> stage)
fromFuture
-
static <T> Streamable<T> fromFuture(Future<? extends T> future) -
static <T> Streamable<T> fromFuture(Future<? extends T> future, long timeout, TimeUnit unit)
fromObservable
-
static <T> Streamable<T> fromObservable(ObservableSource<T> source, BackpressureStrategy strategy)
fromOptional
-
static <T> Streamable<T> fromOptional(Optional<T> optional)
fromRunnable
-
static <T> Streamable<T> fromRunnable(Runnable run)
fromSupplier
-
static <T> Streamable<T> fromSupplier(Supplier<? extends T> supplier)
generate
-
static <T> Streamable<T> generate(Consumer<Emitter<T>> generator) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiConsumer<S, Emitter<T>> generator) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiFunction<S, Emitter<T>, S> generator) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiConsumer<S, Emitter<T>> generator, Consumer<? super S> disposeState) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiFunction<S, Emitter<T>, S> generator, Consumer<? super S> disposeState)
groupJoin
-
<TRight, TLeftEnd, TRightEnd, R> Streamable<R> groupJoin(Streamable<? extends TRight> other, Function<? super T, ? extends Streamable<TLeftEnd>> leftEnd, Function<? super TRight, ? extends Streamable<TRightEnd>> rightEnd, BiFunction<? super T, ? super Streamable<TRight>, ? extends R> resultSelector)
interval
-
static Streamable<Long> interval(long period, TimeUnit unit) -
static Streamable<Long> interval(long initialDelay, long period, TimeUnit unit) -
static Streamable<Long> interval(long period, TimeUnit unit, Scheduler scheduler) -
static Streamable<Long> interval(long initialDelay, long period, TimeUnit unit, Scheduler scheduler)
isEmpty
-
Single<Boolean> isEmpty()
join
-
<TRight, TLeftEnd, TRightEnd, R> Streamable<R> join(Streamable<? extends TRight> other, Function<? super T, ? extends Streamable<TLeftEnd>> leftEnd, Function<? super TRight, ? extends Streamable<TRightEnd>> rightEnd, BiFunction<? super T, ? super TRight, ? extends R> resultSelector)
lastElement
-
Maybe<T> lastElement()
lastOrErrorStage
-
CompletionStage<T> lastOrErrorStage()
lastStage
-
CompletionStage<T> lastStage(T defaultItem)
materialize
-
Streamable<Notification<T>> materialize()
merge
-
static <T> Streamable<T> merge(Iterable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> merge(Streamable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> merge(Iterable<? extends Streamable<? extends T>> sources, StandardConcurrentBufferedConfig config) -
static <T> Streamable<T> merge(Streamable<? extends Streamable<? extends T>> sources, StandardConcurrentBufferedConfig config)
mergeArray
-
static <T> Streamable<T> mergeArray(Streamable<? extends T>... sources) -
static <T> Streamable<T> mergeArray(StandardConcurrentBufferedConfig config, Streamable<? extends T>... sources)
mergeWith
-
Streamable<T> mergeWith(Streamable<? extends T> other) -
Streamable<T> mergeWith(SingleSource<? extends T> other) -
Streamable<T> mergeWith(MaybeSource<? extends T> other) -
Streamable<T> mergeWith(CompletableSource other)
observeOn
-
Streamable<T> observeOn(Scheduler scheduler) -
Streamable<T> observeOn(Scheduler scheduler, StandardBufferedConfig config)
ofType
-
<U> Streamable<U> ofType(Class<U> clazz)
onErrorComplete
-
Streamable<T> onErrorComplete() -
Streamable<T> onErrorComplete(Predicate<? super Throwable> predicate)
onErrorResumeWith
-
Streamable<T> onErrorResumeWith(Streamable<? extends T> fallback)
onErrorReturn
-
Streamable<T> onErrorReturn(Function<? super Throwable, ? extends T> itemSupplier)
onErrorReturnItem
-
Streamable<T> onErrorReturnItem(T item)
publish
-
ConnectableStreamable<T> publish() -
<R> Streamable<R> publish(Function<? super Streamable<T>, ? extends Streamable<R>> selector) -
ConnectableStreamable<T> publish(int bufferSize) -
<R> Streamable<R> publish(Function<? super Streamable<T>, ? extends Streamable<? extends R>> selector, int prefetch)
reduce
-
Maybe<T> reduce(BiFunction<T, T, T> reducer) -
<R> Single<R> reduce(R seed, BiFunction<R, ? super T, R> reducer)
reduceWith
-
<R> Single<R> reduceWith(Supplier<R> seedSupplier, BiFunction<R, ? super T, R> reducer)
repeatUntil
-
Streamable<T> repeatUntil(BooleanSupplier stop)
replay
-
ConnectableStreamable<T> replay() -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector) -
ConnectableStreamable<T> replay(int bufferSize) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize) -
ConnectableStreamable<T> replay(int bufferSize, boolean eagerTruncate) -
ConnectableStreamable<T> replay(long time, TimeUnit unit) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, boolean eagerTruncate) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, long time, TimeUnit unit) -
ConnectableStreamable<T> replay(int bufferSize, long time, TimeUnit unit) -
ConnectableStreamable<T> replay(long time, TimeUnit unit, Scheduler scheduler) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, long time, TimeUnit unit) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, long time, TimeUnit unit, Scheduler scheduler) -
ConnectableStreamable<T> replay(int bufferSize, long time, TimeUnit unit, Scheduler scheduler) -
ConnectableStreamable<T> replay(long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, long time, TimeUnit unit, Scheduler scheduler) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate) -
ConnectableStreamable<T> replay(int bufferSize, long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate)
retryUntil
-
Streamable<T> retryUntil(BooleanSupplier stop)
sample
-
<U> Streamable<T> sample(Streamable<U> sampler) -
Streamable<T> sample(long period, TimeUnit unit) -
<U> Streamable<T> sample(Streamable<U> sampler, boolean emitLast) -
Streamable<T> sample(long period, TimeUnit unit, Scheduler scheduler) -
Streamable<T> sample(long period, TimeUnit unit, Scheduler scheduler, SampleConfig<? super T> config)
scan
-
Streamable<T> scan(BiFunction<T, T, T> accumulator) -
<R> Streamable<R> scan(R initialValue, BiFunction<R, ? super T, R> accumulator)
scanWith
-
<R> Streamable<R> scanWith(Supplier<R> seedSupplier, BiFunction<R, ? super T, R> accumulator)
sequenceEqual
-
static <T> Single<Boolean> sequenceEqual(Streamable<? extends T> source1, Streamable<? extends T> source2) -
static <T> Single<Boolean> sequenceEqual(Streamable<? extends T> source1, Streamable<? extends T> source2, SequenceEqualConfig<T> config)
share
-
Streamable<T> share()
single
-
Single<T> single(T defaultItem)
singleElement
-
Maybe<T> singleElement()
singleOrError
-
Single<T> singleOrError()
singleOrErrorStage
-
CompletionStage<T> singleOrErrorStage()
singleStage
-
CompletionStage<T> singleStage(T defaultItem)
skipLast
-
Streamable<T> skipLast(int count) -
Streamable<T> skipLast(long time, TimeUnit unit) -
Streamable<T> skipLast(long time, TimeUnit unit, Scheduler scheduler) -
Streamable<T> skipLast(long time, TimeUnit unit, Scheduler scheduler, StandardBufferedConfig config)
skipUntil
-
<U> Streamable<T> skipUntil(Streamable<U> other)
skipWhile
-
Streamable<T> skipWhile(Predicate<? super T> predicate)
sorted
-
Streamable<T> sorted() -
Streamable<T> sorted(Comparator<? super T> comparator)
startWith
-
Streamable<T> startWith(CompletableSource other) -
Streamable<T> startWith(SingleSource<T> other) -
Streamable<T> startWith(MaybeSource<T> other) -
Streamable<T> startWith(Streamable<? extends T> other)
startWithArray
-
Streamable<T> startWithArray(T... items)
startWithItem
-
Streamable<T> startWithItem(T item)
startWithIterable
-
Streamable<T> startWithIterable(Iterable<? extends T> items)
subscribeOn
-
Streamable<T> subscribeOn(Scheduler scheduler) -
Streamable<T> subscribeOn(Scheduler scheduler, boolean requestOn)
subscribeWith
-
<E extends Subscriber<? super T>> E subscribeWith(E subscriber)
switchIfEmpty
-
Streamable<T> switchIfEmpty(Streamable<? extends T> other)
switchMap
-
<R> Streamable<R> switchMap(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> switchMap(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardBufferedConfig config)
switchMapCompletable
-
Completable switchMapCompletable(Function<? super T, ? extends CompletableSource> mapper)
switchMapCompletableDelayError
-
Completable switchMapCompletableDelayError(Function<? super T, ? extends CompletableSource> mapper)
switchMapMaybe
-
<R> Streamable<R> switchMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper)
switchMapMaybeDelayError
-
<R> Streamable<R> switchMapMaybeDelayError(Function<? super T, ? extends MaybeSource<? extends R>> mapper)
switchMapSingle
-
<R> Streamable<R> switchMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper)
switchMapSingleDelayError
-
<R> Streamable<R> switchMapSingleDelayError(Function<? super T, ? extends SingleSource<? extends R>> mapper)
switchOnNext
-
static <T> Streamable<T> switchOnNext(Streamable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> switchOnNext(Streamable<? extends Streamable<? extends T>> sources, StandardBufferedConfig config)
takeLast
-
Streamable<T> takeLast(int count) -
Streamable<T> takeLast(long time, TimeUnit unit) -
Streamable<T> takeLast(long count, long time, TimeUnit unit) -
Streamable<T> takeLast(long time, TimeUnit unit, Scheduler scheduler) -
Streamable<T> takeLast(long count, long time, TimeUnit unit, Scheduler scheduler) -
Streamable<T> takeLast(long time, TimeUnit unit, Scheduler scheduler, StandardBufferedConfig config) -
Streamable<T> takeLast(long count, long time, TimeUnit unit, Scheduler scheduler, StandardBufferedConfig config)
throttleFirst
-
Streamable<T> throttleFirst(long windowDuration, TimeUnit unit) -
Streamable<T> throttleFirst(long skipDuration, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleFirst(long skipDuration, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
throttleLast
-
Streamable<T> throttleLast(long intervalDuration, TimeUnit unit) -
Streamable<T> throttleLast(long intervalDuration, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleLast(long intervalDuration, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
throttleLatest
-
Streamable<T> throttleLatest(long timeout, TimeUnit unit) -
Streamable<T> throttleLatest(long timeout, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleLatest(long timeout, TimeUnit unit, Scheduler scheduler, SampleConfig<? super T> config)
throttleWithTimeout
-
Streamable<T> throttleWithTimeout(long timeout, TimeUnit unit) -
Streamable<T> throttleWithTimeout(long timeout, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleWithTimeout(long timeout, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
timeInterval
-
Streamable<Timed<T>> timeInterval() -
Streamable<Timed<T>> timeInterval(Scheduler scheduler) -
Streamable<Timed<T>> timeInterval(TimeUnit unit) -
Streamable<Timed<T>> timeInterval(TimeUnit unit, Scheduler scheduler)
timestamp
-
Streamable<Timed<T>> timestamp() -
Streamable<Timed<T>> timestamp(Scheduler scheduler) -
Streamable<Timed<T>> timestamp(TimeUnit unit) -
Streamable<Timed<T>> timestamp(TimeUnit unit, Scheduler scheduler)
toFuture
-
Future<T> toFuture()
toList
-
Single<List<T>> toList() -
Single<List<T>> toList(int capacityHint) -
<U extends Collection<? super T>> Single<U> toList(Supplier<U> collectionSupplier)
toMap
-
<K> Single<Map<K, T>> toMap(Function<? super T, ? extends K> keySelector) -
<K, V> Single<Map<K, V>> toMap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) -
<K, V> Single<Map<K, V>> toMap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Supplier<? extends Map<K, V>> mapSupplier)
toMultimap
-
<K> Single<Map<K, Collection<T>>> toMultimap(Function<? super T, ? extends K> keySelector) -
<K, V> Single<Map<K, Collection<V>>> toMultimap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) -
<K, V> Single<Map<K, Collection<V>>> toMultimap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Supplier<Map<K, Collection<V>>> mapSupplier) -
<K, V> Single<Map<K, Collection<V>>> toMultimap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Supplier<? extends Map<K, Collection<V>>> mapSupplier, Function<? super K, ? extends Collection<? super V>> collectionFactory)
toSortedList
-
Single<List<T>> toSortedList() -
Single<List<T>> toSortedList(Comparator<? super T> comparator) -
Single<List<T>> toSortedList(int capacityHint) -
Single<List<T>> toSortedList(Comparator<? super T> comparator, int capacityHint)
unsubscribeOn
-
Streamable<T> unsubscribeOn(Scheduler scheduler)
window
-
Streamable<Streamable<T>> window(long count) -
<B> Streamable<Streamable<T>> window(Streamable<B> boundaryIndicator) -
Streamable<Streamable<T>> window(long count, long skip) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit) -
<B> Streamable<Streamable<T>> window(Streamable<B> boundaryIndicator, int bufferSize) -
<U, V> Streamable<Streamable<T>> window(Streamable<U> openingIndicator, Function<? super U, ? extends Streamable<V>> closingIndicator) -
Streamable<Streamable<T>> window(long count, long skip, int bufferSize) -
Streamable<Streamable<T>> window(long timespan, long timeskip, TimeUnit unit) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, long count) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler) -
<U, V> Streamable<Streamable<T>> window(Streamable<U> openingIndicator, Function<? super U, ? extends Streamable<V>> closingIndicator, int bufferSize) -
Streamable<Streamable<T>> window(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, long count, boolean restart) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler, long count) -
Streamable<Streamable<T>> window(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler, int bufferSize) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler, long count, boolean restart) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler, long count, boolean restart, int bufferSize)
withLatestFrom
-
<U, R> Streamable<R> withLatestFrom(Streamable<? extends U> other, BiFunction<? super T, ? super U, ? extends R> combiner) -
<R> Streamable<R> withLatestFrom(Streamable<?>[] others, Function<? super Object[], R> combiner) -
<R> Streamable<R> withLatestFrom(Iterable<? extends Streamable<?>> others, Function<? super Object[], R> combiner) -
<T1, T2, R> Streamable<R> withLatestFrom(Streamable<T1> source1, Streamable<T2> source2, Function3<? super T, ? super T1, ? super T2, R> combiner) -
<T1, T2, T3, R> Streamable<R> withLatestFrom(Streamable<T1> source1, Streamable<T2> source2, Streamable<T3> source3, Function4<? super T, ? super T1, ? super T2, ? super T3, R> combiner) -
<T1, T2, T3, T4, R> Streamable<R> withLatestFrom(Streamable<T1> source1, Streamable<T2> source2, Streamable<T3> source3, Streamable<T4> source4, Function5<? super T, ? super T1, ? super T2, ? super T3, ? super T4, R> combiner)
zipArray
-
static <T, R> Streamable<R> zipArray(Streamable<? extends T>[] sources, Function<? super Object[], ? extends R> zipper, StandardBufferedConfig config)
zipWith
-
<U, R> Streamable<R> zipWith(Iterable<U> other, BiFunction<? super T, ? super U, ? extends R> zipper) -
<U, R> Streamable<R> zipWith(Streamable<? extends U> other, BiFunction<? super T, ? super U, ? extends R> zipper) -
<U, R> Streamable<R> zipWith(Streamable<? extends U> other, BiFunction<? super T, ? super U, ? extends R> zipper, StandardBufferedConfig config)
Probably not applicable to a pull-based Streamable (listed for completeness) (16 methods, 30 overloads)
These are backpressure-, Subscription-, or Flowable-plumbing specific. Tick them off as won't do or replace them with a Streamable-native design where one makes sense (e.g. create(VirtualGenerator) already covers virtualCreate, transform(VirtualTransformer) covers virtualTransform). Signatures in this section are kept verbatim from Flowable.
bufferSize
-
static int bufferSize()
create
-
static <T> Flowable<T> create(FlowableOnSubscribe<T> source, BackpressureStrategy mode)
doOnRequest
-
Flowable<T> doOnRequest(LongConsumer onRequest)
onBackpressureBuffer
-
Flowable<T> onBackpressureBuffer() -
Flowable<T> onBackpressureBuffer(OnBackpressureBufferConfig<? super T> config) -
Flowable<T> onBackpressureBuffer(long capacity, BackpressureOverflowStrategy overflowStrategy) -
Flowable<T> onBackpressureBuffer(long capacity, BackpressureOverflowStrategy overflowStrategy, Consumer<? super T> onDropped)
onBackpressureDrop
-
Flowable<T> onBackpressureDrop() -
Flowable<T> onBackpressureDrop(Consumer<? super T> onDrop)
onBackpressureLatest
-
Flowable<T> onBackpressureLatest() -
Flowable<T> onBackpressureLatest(Consumer<? super T> onDropped)
onBackpressureReduce
-
Flowable<T> onBackpressureReduce(BiFunction<T, T, T> reducer) -
<R> Flowable<R> onBackpressureReduce(Supplier<R> supplier, BiFunction<R, ? super T, R> reducer)
onTerminateDetach
-
Flowable<T> onTerminateDetach()
parallel
-
ParallelFlowable<T> parallel() -
ParallelFlowable<T> parallel(int parallelism) -
ParallelFlowable<T> parallel(int parallelism, int prefetch)
rebatchRequests
-
Flowable<T> rebatchRequests(int n)
safeSubscribe
-
void safeSubscribe(Subscriber<? super T> subscriber)
serialize
-
Flowable<T> serialize()
toStreamable
-
Streamable<T> toStreamable() -
Streamable<T> toStreamable(ExecutorService executor)
unsafeCreate
-
static <T> Flowable<T> unsafeCreate(Publisher<T> onSubscribe)
virtualCreate
-
static <T> Flowable<T> virtualCreate(VirtualGenerator<T> generator) -
static <T> Flowable<T> virtualCreate(VirtualGenerator<T> generator, ExecutorService executor) -
static <T> Flowable<T> virtualCreate(VirtualGenerator<T> generator, Scheduler scheduler)
virtualTransform
-
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer) -
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer, ExecutorService executor) -
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer, Scheduler scheduler, int prefetch) -
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer, ExecutorService executor, int prefetch)
Contributor guide
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
Research direction
Compare the public APIs in src/main/java/io/reactivex/rxjava4/core/Flowable.java and Streamable.java on the 4.x branch, using the checklist to group the missing overloads by operator. Start with one operator group and account for the stated Streamable type mappings and existing shape differences. Done means the listed Streamable overloads are implemented or explicitly resolved and the public API matches the intended Flowable parity.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- api
- Issue type
- Feature
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Active
- Clarity
- Clearly specified
- Newbie friendliness
- 45/100