ReactiveX / ReactiveX/RxJava

4.x: Streamable operator checklist

Offen
#8,299 0 Kommentare 0 Reaktionen 0 zugewiesene Personen Auf GitHub ansehen

Dieses Issue hat noch niemand übernommen.

4.x AI Enhancement Streamable
Vorherrschende Sprache
Java
Sterne
48.2k
Forks
7.6k
Ø Merge
39 Min.
Gemergte PRs (30 T.)
11

Beschreibung

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's Disposable-returning subscribe(...) overloads probably map to CompletionStageDisposable<Void> on Streamable (as forEach already does).
  • Overloads that take a Scheduler may want an additional ExecutorService twin, following the existing Streamable.timer / intervalRange / create pattern. Those twins are not listed here because they have no Flowable counterpart.
  • A few methods already exist on Streamable with a different shape (e.g. flatMap(mapper, StandardConcurrentConfig) vs. Flowable's StandardConcurrentBufferedConfig, retry(BiPredicate<Long, Throwable>) vs. BiPredicate<Integer, Throwable>, repeatWhen/retryWhen taking CompletionStage<Boolean> functions, create(VirtualGenerator) vs. virtualCreate). These are listed in the missing overloads section only where the Flowable overload 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() returns Completable on Flowable but Streamable<T> on Streamable; collect(Collector) returns Single<R> vs. Streamable<R>; forEach returns Disposable vs. 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)

Beitragsleitfaden

Beitragsleitfaden öffnen

Erste Schritte

  1. Lies das ganze Issue und danach den Beitragsleitfaden des Projekts.
  2. Schreib ins Issue, dass du es übernimmst — das erspart doppelte Arbeit.
  3. Forke das Repository und arbeite in einem Branch.
  4. Öffne einen Pull Request, der die Issue-Nummer nennt.

Rechercherichtung

Vergleiche die öffentlichen APIs in src/main/java/io/reactivex/rxjava4/core/Flowable.java und Streamable.java im 4.x-Zweig und verwende die Checkliste, um die fehlenden Überladungen nach Operator zu gruppieren. Beginne mit einer Operatorgruppe und berücksichtige die angegebenen Streamable-Typzuordnungen sowie die bestehenden Formunterschiede. Als erledigt gilt die Aufgabe, wenn die aufgeführten Streamable-Überladungen implementiert oder ausdrücklich geklärt sind und die öffentliche API der vorgesehenen Flowable-Parität entspricht.

Vom Indexierungsmodell aus dem Issue-Text verfasst.

Bewertung

Tech-Stack
java
Bereich
api
Issue-Typ
Feature
Schwierigkeit
4/5
Geschätzter Aufwand
3-5 Tage
Aktivitätsstatus
Aktiv
Klarheit
Klar beschrieben
Anfängerfreundlichkeit
45/100

Neue Issues direkt in Ihr Postfach

Eine kurze Übersicht über anfängerfreundliche GitHub-Issues.