ReactiveX / ReactiveX/RxJava

4.x: Streamable operator checklist

Aperta
#8,299 0 commenti 0 reazioni 0 assegnatari Vedi su GitHub

Nessuno ha ancora preso questa issue.

4.x AI Enhancement Streamable
Lingua principale
Java
Stelle
48.2k
Fork
7.6k
Merge medio
39m
PR unite (30g)
11

Descrizione

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)

Guida per i contributori

Apri la guida per i contributori

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Direzione di ricerca

Confronta le API pubbliche in src/main/java/io/reactivex/rxjava4/core/Flowable.java e Streamable.java nel ramo 4.x, usando la checklist per raggruppare gli overload mancanti per operatore. Inizia con un gruppo di operatori e considera i mapping dei tipi Streamable indicati e le differenze di forma esistenti. Il lavoro è completato quando gli overload Streamable elencati sono implementati o risolti esplicitamente e l’API pubblica corrisponde alla parità Flowable prevista.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
java
Ambito
api
Tipo di issue
Funzionalità
Difficoltà
4/5
Tempo stimato
3-5 giorni
Stato di attività
Attiva
Chiarezza
Specificata chiaramente
Idoneità per principianti
45/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.