4.x: Streamable operator checklist
Dieses Issue hat noch niemand übernommen.
- 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'sDisposable-returningsubscribe(...)overloads probably map toCompletionStageDisposable<Void>onStreamable(asforEachalready does).- Overloads that take a
Schedulermay want an additionalExecutorServicetwin, following the existingStreamable.timer/intervalRange/createpattern. Those twins are not listed here because they have noFlowablecounterpart. - A few methods already exist on
Streamablewith a different shape (e.g.flatMap(mapper, StandardConcurrentConfig)vs.Flowable'sStandardConcurrentBufferedConfig,retry(BiPredicate<Long, Throwable>)vs.BiPredicate<Integer, Throwable>,repeatWhen/retryWhentakingCompletionStage<Boolean>functions,create(VirtualGenerator)vs.virtualCreate). These are listed in the missing overloads section only where theFlowableoverload has no exact-arity/equivalent-type match; adjust or tick them off as appropriate. - Return-type differences on otherwise matching signatures are not listed (e.g.
ignoreElements()returnsCompletableonFlowablebutStreamable<T>onStreamable;collect(Collector)returnsSingle<R>vs.Streamable<R>;forEachreturnsDisposablevs.CompletionStageDisposable<Void>).
Methods that exist on Streamable but are missing overloads (23 methods, 69 overloads)
blockingFirst
-
T blockingFirst(T defaultItem)
blockingLast
-
T blockingLast(T defaultItem)
collect
-
<U> Single<U> collect(Supplier<? extends U> initialItemSupplier, BiConsumer<? super U, ? super T> collector)
concat
-
static <T> Streamable<T> concat(Streamable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> concat(Iterable<? extends Streamable<? extends T>> sources, StandardBufferedConfig config) -
static <T> Streamable<T> concat(Streamable<? extends Streamable<? extends T>> sources, StandardBufferedConfig config)
delay
-
<U> Streamable<T> delay(Function<? super T, ? extends Streamable<U>> itemDelayIndicator) -
Streamable<T> delay(long time, TimeUnit unit) -
<U, V> Streamable<T> delay(Streamable<U> subscriptionIndicator, Function<? super T, ? extends Streamable<V>> itemDelayIndicator) -
Streamable<T> delay(long time, TimeUnit unit, Scheduler scheduler, boolean delayError)
error
-
static <T> Streamable<T> error(Supplier<? extends Throwable> supplier)
flatMap
-
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardConcurrentBufferedConfig config) -
<U, R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends R> combiner) -
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> onNextMapper, Function<? super Throwable, ? extends Streamable<? extends R>> onErrorMapper, Supplier<? extends Streamable<? extends R>> onCompleteSupplier) -
<U, R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends R> combiner, StandardConcurrentBufferedConfig config) -
<R> Streamable<R> flatMap(Function<? super T, ? extends Streamable<? extends R>> onNextMapper, Function<Throwable, ? extends Streamable<? extends R>> onErrorMapper, Supplier<? extends Streamable<? extends R>> onCompleteSupplier, StandardConcurrentBufferedConfig config)
groupBy
-
<K> Streamable<GroupedStreamable<K, T>> groupBy(Function<? super T, ? extends K> keySelector, StandardBufferedConfig config) -
<K, V> Streamable<GroupedStreamable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) -
<K, V> Streamable<GroupedStreamable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, StandardBufferedConfig config) -
<K, V> Streamable<GroupedStreamable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Function<? super Consumer<Object>, ? extends Map<K, Object>> evictingMapFactory, StandardBufferedConfig config)
intervalRange
-
static Streamable<Long> intervalRange(long start, long count, long initialDelay, long period, TimeUnit unit)
just
-
static <T> Streamable<T> just(T item1, T item2) -
static <T> Streamable<T> just(T item1, T item2, T item3) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8, T item9) -
static <T> Streamable<T> just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8, T item9, T item10)
repeat
-
Streamable<T> repeat()
repeatWhen
-
Streamable<T> repeatWhen(Function<? super Streamable<Object>, ? extends Streamable<?>> handler)
retry
-
Streamable<T> retry() -
Streamable<T> retry(BiPredicate<? super Integer, ? super Throwable> predicate) -
Streamable<T> retry(long times, Predicate<? super Throwable> predicate)
retryWhen
-
Streamable<T> retryWhen(Function<? super Streamable<Throwable>, ? extends Streamable<?>> handler)
skip
-
Streamable<T> skip(long time, TimeUnit unit) -
Streamable<T> skip(long time, TimeUnit unit, Scheduler scheduler)
subscribe
-
Disposable subscribe() -
Disposable subscribe(Consumer<? super T> onNext) -
Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError) -
Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete) -
Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete, DisposableContainer container)
take
-
Streamable<T> take(long time, TimeUnit unit) -
Streamable<T> take(long time, TimeUnit unit, Scheduler scheduler)
takeUntil
-
Streamable<T> takeUntil(Predicate<? super T> stopPredicate)
test
-
TestSubscriber<T> test(long initialRequest) -
TestSubscriber<T> test(long initialRequest, boolean cancel)
timeout
-
<V> Streamable<T> timeout(Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator) -
<V> Streamable<T> timeout(Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator, Streamable<? extends T> fallback) -
Streamable<T> timeout(long timeout, TimeUnit unit) -
<U, V> Streamable<T> timeout(Streamable<U> firstTimeoutIndicator, Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator) -
Streamable<T> timeout(long timeout, TimeUnit unit, Streamable<? extends T> fallback) -
Streamable<T> timeout(long timeout, TimeUnit unit, Scheduler scheduler) -
<U, V> Streamable<T> timeout(Streamable<U> firstTimeoutIndicator, Function<? super T, ? extends Streamable<V>> itemTimeoutIndicator, Streamable<? extends T> fallback)
timer
-
static Streamable<Long> timer(long delay, TimeUnit unit)
using
-
static <T, D> Streamable<T> using(Supplier<? extends D> resourceSupplier, Function<? super D, ? extends Streamable<? extends T>> sourceSupplier, Consumer<? super D> resourceCleanup, boolean eager)
zip
-
static <T, R> Streamable<R> zip(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> zipper) -
static <T, R> Streamable<R> zip(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> zipper, StandardBufferedConfig config) -
static <T1, T2, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, BiFunction<? super T1, ? super T2, ? extends R> zipper) -
static <T1, T2, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, BiFunction<? super T1, ? super T2, ? extends R> zipper, StandardBufferedConfig config) -
static <T1, T2, T3, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Function3<? super T1, ? super T2, ? super T3, ? extends R> zipper) -
static <T1, T2, T3, T4, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Function4<? super T1, ? super T2, ? super T3, ? super T4, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Function5<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Function6<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, T7, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Function7<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, T7, T8, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Function8<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? extends R> zipper) -
static <T1, T2, T3, T4, T5, T6, T7, T8, T9, R> Streamable<R> zip(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Streamable<? extends T9> source9, Function9<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? super T9, ? extends R> zipper)
Methods that are entirely missing from Streamable (140 methods, 317 overloads)
all
-
Single<Boolean> all(Predicate<? super T> predicate)
amb
-
static <T> Streamable<T> amb(Iterable<? extends Streamable<? extends T>> sources)
ambArray
-
static <T> Streamable<T> ambArray(Streamable<? extends T>... sources)
ambWith
-
Streamable<T> ambWith(Streamable<? extends T> other)
any
-
Single<Boolean> any(Predicate<? super T> predicate)
blockingForEach
-
void blockingForEach(Consumer<? super T> onNext) -
void blockingForEach(Consumer<? super T> onNext, int bufferSize)
blockingIterable
-
Iterable<T> blockingIterable() -
Iterable<T> blockingIterable(int bufferSize)
blockingLatest
-
Iterable<T> blockingLatest()
blockingMostRecent
-
Iterable<T> blockingMostRecent(T initialItem)
blockingNext
-
Iterable<T> blockingNext()
blockingSingle
-
T blockingSingle() -
T blockingSingle(T defaultItem)
blockingStream
-
Stream<T> blockingStream() -
Stream<T> blockingStream(int prefetch)
blockingSubscribe
-
void blockingSubscribe() -
void blockingSubscribe(Consumer<? super T> onNext) -
void blockingSubscribe(Subscriber<? super T> subscriber) -
void blockingSubscribe(Consumer<? super T> onNext, int bufferSize) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, int bufferSize) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete) -
void blockingSubscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete, int bufferSize)
buffer
-
Streamable<List<T>> buffer(int count) -
<B> Streamable<List<T>> buffer(Streamable<B> boundaryIndicator) -
Streamable<List<T>> buffer(int count, int skip) -
<U extends Collection<? super T>> Streamable<U> buffer(int count, Supplier<U> bufferSupplier) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit) -
<TOpening, TClosing> Streamable<List<T>> buffer(Streamable<? extends TOpening> openingIndicator, Function<? super TOpening, ? extends Streamable<? extends TClosing>> closingIndicator) -
<B> Streamable<List<T>> buffer(Streamable<B> boundaryIndicator, int initialCapacity) -
<B, U extends Collection<? super T>> Streamable<U> buffer(Streamable<B> boundaryIndicator, Supplier<U> bufferSupplier) -
<U extends Collection<? super T>> Streamable<U> buffer(int count, int skip, Supplier<U> bufferSupplier) -
Streamable<List<T>> buffer(long timespan, long timeskip, TimeUnit unit) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit, int count) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit, Scheduler scheduler) -
<TOpening, TClosing, U extends Collection<? super T>> Streamable<U> buffer(Streamable<? extends TOpening> openingIndicator, Function<? super TOpening, ? extends Streamable<? extends TClosing>> closingIndicator, Supplier<U> bufferSupplier) -
Streamable<List<T>> buffer(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler) -
Streamable<List<T>> buffer(long timespan, TimeUnit unit, Scheduler scheduler, int count) -
<U extends Collection<? super T>> Streamable<U> buffer(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler, Supplier<U> bufferSupplier) -
<U extends Collection<? super T>> Streamable<U> buffer(long timespan, TimeUnit unit, Scheduler scheduler, int count, Supplier<U> bufferSupplier, boolean restartTimerOnMaxSize)
cache
-
Streamable<T> cache()
cacheWithInitialCapacity
-
Streamable<T> cacheWithInitialCapacity(int initialCapacity)
cast
-
<U> Streamable<U> cast(Class<U> clazz)
collectInto
-
<U> Single<U> collectInto(U initialItem, BiConsumer<? super U, ? super T> collector)
combineLatest
-
static <T, R> Streamable<R> combineLatest(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> combiner) -
static <T, R> Streamable<R> combineLatest(Iterable<? extends Streamable<? extends T>> sources, Function<? super Object[], ? extends R> combiner, StandardBufferedConfig config) -
static <T1, T2, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, BiFunction<? super T1, ? super T2, ? extends R> combiner) -
static <T1, T2, T3, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Function3<? super T1, ? super T2, ? super T3, ? extends R> combiner) -
static <T1, T2, T3, T4, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Function4<? super T1, ? super T2, ? super T3, ? super T4, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Function5<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Function6<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, T7, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Function7<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, T7, T8, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Function8<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? extends R> combiner) -
static <T1, T2, T3, T4, T5, T6, T7, T8, T9, R> Streamable<R> combineLatest(Streamable<? extends T1> source1, Streamable<? extends T2> source2, Streamable<? extends T3> source3, Streamable<? extends T4> source4, Streamable<? extends T5> source5, Streamable<? extends T6> source6, Streamable<? extends T7> source7, Streamable<? extends T8> source8, Streamable<? extends T9> source9, Function9<? super T1, ? super T2, ? super T3, ? super T4, ? super T5, ? super T6, ? super T7, ? super T8, ? super T9, ? extends R> combiner)
combineLatestArray
-
static <T, R> Streamable<R> combineLatestArray(Streamable<? extends T>[] sources, Function<? super Object[], ? extends R> combiner) -
static <T, R> Streamable<R> combineLatestArray(Streamable<? extends T>[] sources, Function<? super Object[], ? extends R> combiner, StandardBufferedConfig config)
compose
-
<R> Streamable<R> compose(StreamableTransformer<? super T, ? extends R> composer)
concatArray
-
static <T> Streamable<T> concatArray(Streamable<? extends T>... sources) -
static <T> Streamable<T> concatArray(StandardBufferedConfig config, Streamable<? extends T>... sources)
concatArrayEager
-
static <T> Streamable<T> concatArrayEager(Streamable<? extends T>... sources)
concatEager
-
static <T> Streamable<T> concatEager(Iterable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> concatEager(Streamable<? extends Streamable<? extends T>> sources)
concatMap
-
<R> Streamable<R> concatMap(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> concatMap(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardBufferedConfig config) -
<R> Streamable<R> concatMap(Function<? super T, ? extends Streamable<? extends R>> mapper, Scheduler scheduler, StandardBufferedConfig config)
concatMapCompletable
-
Completable concatMapCompletable(Function<? super T, ? extends CompletableSource> mapper) -
Completable concatMapCompletable(Function<? super T, ? extends CompletableSource> mapper, StandardBufferedConfig config)
concatMapEager
-
<R> Streamable<R> concatMapEager(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> concatMapEager(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardConcurrentBufferedConfig config)
concatMapIterable
-
<U> Streamable<U> concatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper) -
<U> Streamable<U> concatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, StandardBufferedConfig config)
concatMapMaybe
-
<R> Streamable<R> concatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper) -
<R> Streamable<R> concatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper, StandardBufferedConfig config)
concatMapSingle
-
<R> Streamable<R> concatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper) -
<R> Streamable<R> concatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper, StandardBufferedConfig config)
concatMapStream
-
<R> Streamable<R> concatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper) -
<R> Streamable<R> concatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper, StandardBufferedConfig config)
concatWith
-
Streamable<T> concatWith(Streamable<? extends T> other) -
Streamable<T> concatWith(SingleSource<? extends T> other) -
Streamable<T> concatWith(MaybeSource<? extends T> other) -
Streamable<T> concatWith(CompletableSource other)
contains
-
Single<Boolean> contains(Object item)
count
-
Single<Long> count()
debounce
-
<U> Streamable<T> debounce(Function<? super T, ? extends Streamable<U>> debounceIndicator) -
Streamable<T> debounce(long timeout, TimeUnit unit) -
Streamable<T> debounce(long timeout, TimeUnit unit, Scheduler scheduler) -
Streamable<T> debounce(long timeout, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
defaultIfEmpty
-
Streamable<T> defaultIfEmpty(T defaultItem)
delaySubscription
-
<U> Streamable<T> delaySubscription(Streamable<U> subscriptionIndicator) -
Streamable<T> delaySubscription(long time, TimeUnit unit) -
Streamable<T> delaySubscription(long time, TimeUnit unit, Scheduler scheduler)
dematerialize
-
<R> Streamable<R> dematerialize(Function<? super T, Notification<R>> selector)
distinct
-
Streamable<T> distinct() -
<K> Streamable<T> distinct(Function<? super T, K> keySelector) -
<K> Streamable<T> distinct(Function<? super T, K> keySelector, Supplier<? extends Collection<? super K>> collectionSupplier)
distinctUntilChanged
-
Streamable<T> distinctUntilChanged() -
<K> Streamable<T> distinctUntilChanged(Function<? super T, K> keySelector) -
Streamable<T> distinctUntilChanged(BiPredicate<? super T, ? super T> comparer)
doAfterNext
-
Streamable<T> doAfterNext(Consumer<? super T> onAfterNext)
doAfterTerminate
-
Streamable<T> doAfterTerminate(Action onAfterTerminate)
doFinally
-
Streamable<T> doFinally(Action onFinally)
doOnCancel
-
Streamable<T> doOnCancel(Action onCancel)
doOnComplete
-
Streamable<T> doOnComplete(Action onComplete)
doOnEach
-
Streamable<T> doOnEach(Consumer<? super Notification<T>> onNotification) -
Streamable<T> doOnEach(Subscriber<? super T> subscriber)
doOnLifecycle
-
Streamable<T> doOnLifecycle(Consumer<? super Subscription> onSubscribe, LongConsumer onRequest, Action onCancel)
doOnSubscribe
-
Streamable<T> doOnSubscribe(Consumer<? super Subscription> onSubscribe)
doOnTerminate
-
Streamable<T> doOnTerminate(Action onTerminate)
elementAt
-
Maybe<T> elementAt(long index) -
Single<T> elementAt(long index, T defaultItem)
elementAtOrError
-
Single<T> elementAtOrError(long index)
first
-
Single<T> first(T defaultItem)
firstElement
-
Maybe<T> firstElement()
firstOrError
-
Single<T> firstOrError()
firstOrErrorStage
-
CompletionStage<T> firstOrErrorStage()
firstStage
-
CompletionStage<T> firstStage(T defaultItem)
flatMapCompletable
-
Completable flatMapCompletable(Function<? super T, ? extends CompletableSource> mapper) -
Completable flatMapCompletable(Function<? super T, ? extends CompletableSource> mapper, StandardConcurrentConfig config)
flatMapIterable
-
<U> Streamable<U> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper) -
<U> Streamable<U> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, StandardBufferedConfig config) -
<U, V> Streamable<V> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends V> combiner) -
<U, V> Streamable<V> flatMapIterable(Function<? super T, ? extends Iterable<? extends U>> mapper, BiFunction<? super T, ? super U, ? extends V> combiner, StandardConcurrentBufferedConfig config)
flatMapMaybe
-
<R> Streamable<R> flatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper) -
<R> Streamable<R> flatMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper, StandardConcurrentConfig config)
flatMapSingle
-
<R> Streamable<R> flatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper) -
<R> Streamable<R> flatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper, StandardConcurrentConfig config)
flatMapStream
-
<R> Streamable<R> flatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper) -
<R> Streamable<R> flatMapStream(Function<? super T, ? extends Stream<? extends R>> mapper, StandardBufferedConfig config)
forEachWhile
-
Disposable forEachWhile(Predicate<? super T> onNext) -
Disposable forEachWhile(Predicate<? super T> onNext, Consumer<? super Throwable> onError) -
Disposable forEachWhile(Predicate<? super T> onNext, Consumer<? super Throwable> onError, Action onComplete)
fromAction
-
static <T> Streamable<T> fromAction(Action action)
fromCallable
-
static <T> Streamable<T> fromCallable(Callable<? extends T> callable)
fromCompletionStage
-
static <T> Streamable<T> fromCompletionStage(CompletionStage<T> stage)
fromFuture
-
static <T> Streamable<T> fromFuture(Future<? extends T> future) -
static <T> Streamable<T> fromFuture(Future<? extends T> future, long timeout, TimeUnit unit)
fromObservable
-
static <T> Streamable<T> fromObservable(ObservableSource<T> source, BackpressureStrategy strategy)
fromOptional
-
static <T> Streamable<T> fromOptional(Optional<T> optional)
fromRunnable
-
static <T> Streamable<T> fromRunnable(Runnable run)
fromSupplier
-
static <T> Streamable<T> fromSupplier(Supplier<? extends T> supplier)
generate
-
static <T> Streamable<T> generate(Consumer<Emitter<T>> generator) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiConsumer<S, Emitter<T>> generator) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiFunction<S, Emitter<T>, S> generator) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiConsumer<S, Emitter<T>> generator, Consumer<? super S> disposeState) -
static <T, S> Streamable<T> generate(Supplier<S> initialState, BiFunction<S, Emitter<T>, S> generator, Consumer<? super S> disposeState)
groupJoin
-
<TRight, TLeftEnd, TRightEnd, R> Streamable<R> groupJoin(Streamable<? extends TRight> other, Function<? super T, ? extends Streamable<TLeftEnd>> leftEnd, Function<? super TRight, ? extends Streamable<TRightEnd>> rightEnd, BiFunction<? super T, ? super Streamable<TRight>, ? extends R> resultSelector)
interval
-
static Streamable<Long> interval(long period, TimeUnit unit) -
static Streamable<Long> interval(long initialDelay, long period, TimeUnit unit) -
static Streamable<Long> interval(long period, TimeUnit unit, Scheduler scheduler) -
static Streamable<Long> interval(long initialDelay, long period, TimeUnit unit, Scheduler scheduler)
isEmpty
-
Single<Boolean> isEmpty()
join
-
<TRight, TLeftEnd, TRightEnd, R> Streamable<R> join(Streamable<? extends TRight> other, Function<? super T, ? extends Streamable<TLeftEnd>> leftEnd, Function<? super TRight, ? extends Streamable<TRightEnd>> rightEnd, BiFunction<? super T, ? super TRight, ? extends R> resultSelector)
lastElement
-
Maybe<T> lastElement()
lastOrErrorStage
-
CompletionStage<T> lastOrErrorStage()
lastStage
-
CompletionStage<T> lastStage(T defaultItem)
materialize
-
Streamable<Notification<T>> materialize()
merge
-
static <T> Streamable<T> merge(Iterable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> merge(Streamable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> merge(Iterable<? extends Streamable<? extends T>> sources, StandardConcurrentBufferedConfig config) -
static <T> Streamable<T> merge(Streamable<? extends Streamable<? extends T>> sources, StandardConcurrentBufferedConfig config)
mergeArray
-
static <T> Streamable<T> mergeArray(Streamable<? extends T>... sources) -
static <T> Streamable<T> mergeArray(StandardConcurrentBufferedConfig config, Streamable<? extends T>... sources)
mergeWith
-
Streamable<T> mergeWith(Streamable<? extends T> other) -
Streamable<T> mergeWith(SingleSource<? extends T> other) -
Streamable<T> mergeWith(MaybeSource<? extends T> other) -
Streamable<T> mergeWith(CompletableSource other)
observeOn
-
Streamable<T> observeOn(Scheduler scheduler) -
Streamable<T> observeOn(Scheduler scheduler, StandardBufferedConfig config)
ofType
-
<U> Streamable<U> ofType(Class<U> clazz)
onErrorComplete
-
Streamable<T> onErrorComplete() -
Streamable<T> onErrorComplete(Predicate<? super Throwable> predicate)
onErrorResumeWith
-
Streamable<T> onErrorResumeWith(Streamable<? extends T> fallback)
onErrorReturn
-
Streamable<T> onErrorReturn(Function<? super Throwable, ? extends T> itemSupplier)
onErrorReturnItem
-
Streamable<T> onErrorReturnItem(T item)
publish
-
ConnectableStreamable<T> publish() -
<R> Streamable<R> publish(Function<? super Streamable<T>, ? extends Streamable<R>> selector) -
ConnectableStreamable<T> publish(int bufferSize) -
<R> Streamable<R> publish(Function<? super Streamable<T>, ? extends Streamable<? extends R>> selector, int prefetch)
reduce
-
Maybe<T> reduce(BiFunction<T, T, T> reducer) -
<R> Single<R> reduce(R seed, BiFunction<R, ? super T, R> reducer)
reduceWith
-
<R> Single<R> reduceWith(Supplier<R> seedSupplier, BiFunction<R, ? super T, R> reducer)
repeatUntil
-
Streamable<T> repeatUntil(BooleanSupplier stop)
replay
-
ConnectableStreamable<T> replay() -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector) -
ConnectableStreamable<T> replay(int bufferSize) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize) -
ConnectableStreamable<T> replay(int bufferSize, boolean eagerTruncate) -
ConnectableStreamable<T> replay(long time, TimeUnit unit) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, boolean eagerTruncate) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, long time, TimeUnit unit) -
ConnectableStreamable<T> replay(int bufferSize, long time, TimeUnit unit) -
ConnectableStreamable<T> replay(long time, TimeUnit unit, Scheduler scheduler) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, long time, TimeUnit unit) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, long time, TimeUnit unit, Scheduler scheduler) -
ConnectableStreamable<T> replay(int bufferSize, long time, TimeUnit unit, Scheduler scheduler) -
ConnectableStreamable<T> replay(long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, long time, TimeUnit unit, Scheduler scheduler) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate) -
ConnectableStreamable<T> replay(int bufferSize, long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate) -
<R> Streamable<R> replay(Function<? super Streamable<T>, ? extends Streamable<R>> selector, int bufferSize, long time, TimeUnit unit, Scheduler scheduler, boolean eagerTruncate)
retryUntil
-
Streamable<T> retryUntil(BooleanSupplier stop)
sample
-
<U> Streamable<T> sample(Streamable<U> sampler) -
Streamable<T> sample(long period, TimeUnit unit) -
<U> Streamable<T> sample(Streamable<U> sampler, boolean emitLast) -
Streamable<T> sample(long period, TimeUnit unit, Scheduler scheduler) -
Streamable<T> sample(long period, TimeUnit unit, Scheduler scheduler, SampleConfig<? super T> config)
scan
-
Streamable<T> scan(BiFunction<T, T, T> accumulator) -
<R> Streamable<R> scan(R initialValue, BiFunction<R, ? super T, R> accumulator)
scanWith
-
<R> Streamable<R> scanWith(Supplier<R> seedSupplier, BiFunction<R, ? super T, R> accumulator)
sequenceEqual
-
static <T> Single<Boolean> sequenceEqual(Streamable<? extends T> source1, Streamable<? extends T> source2) -
static <T> Single<Boolean> sequenceEqual(Streamable<? extends T> source1, Streamable<? extends T> source2, SequenceEqualConfig<T> config)
share
-
Streamable<T> share()
single
-
Single<T> single(T defaultItem)
singleElement
-
Maybe<T> singleElement()
singleOrError
-
Single<T> singleOrError()
singleOrErrorStage
-
CompletionStage<T> singleOrErrorStage()
singleStage
-
CompletionStage<T> singleStage(T defaultItem)
skipLast
-
Streamable<T> skipLast(int count) -
Streamable<T> skipLast(long time, TimeUnit unit) -
Streamable<T> skipLast(long time, TimeUnit unit, Scheduler scheduler) -
Streamable<T> skipLast(long time, TimeUnit unit, Scheduler scheduler, StandardBufferedConfig config)
skipUntil
-
<U> Streamable<T> skipUntil(Streamable<U> other)
skipWhile
-
Streamable<T> skipWhile(Predicate<? super T> predicate)
sorted
-
Streamable<T> sorted() -
Streamable<T> sorted(Comparator<? super T> comparator)
startWith
-
Streamable<T> startWith(CompletableSource other) -
Streamable<T> startWith(SingleSource<T> other) -
Streamable<T> startWith(MaybeSource<T> other) -
Streamable<T> startWith(Streamable<? extends T> other)
startWithArray
-
Streamable<T> startWithArray(T... items)
startWithItem
-
Streamable<T> startWithItem(T item)
startWithIterable
-
Streamable<T> startWithIterable(Iterable<? extends T> items)
subscribeOn
-
Streamable<T> subscribeOn(Scheduler scheduler) -
Streamable<T> subscribeOn(Scheduler scheduler, boolean requestOn)
subscribeWith
-
<E extends Subscriber<? super T>> E subscribeWith(E subscriber)
switchIfEmpty
-
Streamable<T> switchIfEmpty(Streamable<? extends T> other)
switchMap
-
<R> Streamable<R> switchMap(Function<? super T, ? extends Streamable<? extends R>> mapper) -
<R> Streamable<R> switchMap(Function<? super T, ? extends Streamable<? extends R>> mapper, StandardBufferedConfig config)
switchMapCompletable
-
Completable switchMapCompletable(Function<? super T, ? extends CompletableSource> mapper)
switchMapCompletableDelayError
-
Completable switchMapCompletableDelayError(Function<? super T, ? extends CompletableSource> mapper)
switchMapMaybe
-
<R> Streamable<R> switchMapMaybe(Function<? super T, ? extends MaybeSource<? extends R>> mapper)
switchMapMaybeDelayError
-
<R> Streamable<R> switchMapMaybeDelayError(Function<? super T, ? extends MaybeSource<? extends R>> mapper)
switchMapSingle
-
<R> Streamable<R> switchMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper)
switchMapSingleDelayError
-
<R> Streamable<R> switchMapSingleDelayError(Function<? super T, ? extends SingleSource<? extends R>> mapper)
switchOnNext
-
static <T> Streamable<T> switchOnNext(Streamable<? extends Streamable<? extends T>> sources) -
static <T> Streamable<T> switchOnNext(Streamable<? extends Streamable<? extends T>> sources, StandardBufferedConfig config)
takeLast
-
Streamable<T> takeLast(int count) -
Streamable<T> takeLast(long time, TimeUnit unit) -
Streamable<T> takeLast(long count, long time, TimeUnit unit) -
Streamable<T> takeLast(long time, TimeUnit unit, Scheduler scheduler) -
Streamable<T> takeLast(long count, long time, TimeUnit unit, Scheduler scheduler) -
Streamable<T> takeLast(long time, TimeUnit unit, Scheduler scheduler, StandardBufferedConfig config) -
Streamable<T> takeLast(long count, long time, TimeUnit unit, Scheduler scheduler, StandardBufferedConfig config)
throttleFirst
-
Streamable<T> throttleFirst(long windowDuration, TimeUnit unit) -
Streamable<T> throttleFirst(long skipDuration, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleFirst(long skipDuration, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
throttleLast
-
Streamable<T> throttleLast(long intervalDuration, TimeUnit unit) -
Streamable<T> throttleLast(long intervalDuration, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleLast(long intervalDuration, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
throttleLatest
-
Streamable<T> throttleLatest(long timeout, TimeUnit unit) -
Streamable<T> throttleLatest(long timeout, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleLatest(long timeout, TimeUnit unit, Scheduler scheduler, SampleConfig<? super T> config)
throttleWithTimeout
-
Streamable<T> throttleWithTimeout(long timeout, TimeUnit unit) -
Streamable<T> throttleWithTimeout(long timeout, TimeUnit unit, Scheduler scheduler) -
Streamable<T> throttleWithTimeout(long timeout, TimeUnit unit, Scheduler scheduler, Consumer<? super T> onDropped)
timeInterval
-
Streamable<Timed<T>> timeInterval() -
Streamable<Timed<T>> timeInterval(Scheduler scheduler) -
Streamable<Timed<T>> timeInterval(TimeUnit unit) -
Streamable<Timed<T>> timeInterval(TimeUnit unit, Scheduler scheduler)
timestamp
-
Streamable<Timed<T>> timestamp() -
Streamable<Timed<T>> timestamp(Scheduler scheduler) -
Streamable<Timed<T>> timestamp(TimeUnit unit) -
Streamable<Timed<T>> timestamp(TimeUnit unit, Scheduler scheduler)
toFuture
-
Future<T> toFuture()
toList
-
Single<List<T>> toList() -
Single<List<T>> toList(int capacityHint) -
<U extends Collection<? super T>> Single<U> toList(Supplier<U> collectionSupplier)
toMap
-
<K> Single<Map<K, T>> toMap(Function<? super T, ? extends K> keySelector) -
<K, V> Single<Map<K, V>> toMap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) -
<K, V> Single<Map<K, V>> toMap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Supplier<? extends Map<K, V>> mapSupplier)
toMultimap
-
<K> Single<Map<K, Collection<T>>> toMultimap(Function<? super T, ? extends K> keySelector) -
<K, V> Single<Map<K, Collection<V>>> toMultimap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) -
<K, V> Single<Map<K, Collection<V>>> toMultimap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Supplier<Map<K, Collection<V>>> mapSupplier) -
<K, V> Single<Map<K, Collection<V>>> toMultimap(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, Supplier<? extends Map<K, Collection<V>>> mapSupplier, Function<? super K, ? extends Collection<? super V>> collectionFactory)
toSortedList
-
Single<List<T>> toSortedList() -
Single<List<T>> toSortedList(Comparator<? super T> comparator) -
Single<List<T>> toSortedList(int capacityHint) -
Single<List<T>> toSortedList(Comparator<? super T> comparator, int capacityHint)
unsubscribeOn
-
Streamable<T> unsubscribeOn(Scheduler scheduler)
window
-
Streamable<Streamable<T>> window(long count) -
<B> Streamable<Streamable<T>> window(Streamable<B> boundaryIndicator) -
Streamable<Streamable<T>> window(long count, long skip) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit) -
<B> Streamable<Streamable<T>> window(Streamable<B> boundaryIndicator, int bufferSize) -
<U, V> Streamable<Streamable<T>> window(Streamable<U> openingIndicator, Function<? super U, ? extends Streamable<V>> closingIndicator) -
Streamable<Streamable<T>> window(long count, long skip, int bufferSize) -
Streamable<Streamable<T>> window(long timespan, long timeskip, TimeUnit unit) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, long count) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler) -
<U, V> Streamable<Streamable<T>> window(Streamable<U> openingIndicator, Function<? super U, ? extends Streamable<V>> closingIndicator, int bufferSize) -
Streamable<Streamable<T>> window(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, long count, boolean restart) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler, long count) -
Streamable<Streamable<T>> window(long timespan, long timeskip, TimeUnit unit, Scheduler scheduler, int bufferSize) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler, long count, boolean restart) -
Streamable<Streamable<T>> window(long timespan, TimeUnit unit, Scheduler scheduler, long count, boolean restart, int bufferSize)
withLatestFrom
-
<U, R> Streamable<R> withLatestFrom(Streamable<? extends U> other, BiFunction<? super T, ? super U, ? extends R> combiner) -
<R> Streamable<R> withLatestFrom(Streamable<?>[] others, Function<? super Object[], R> combiner) -
<R> Streamable<R> withLatestFrom(Iterable<? extends Streamable<?>> others, Function<? super Object[], R> combiner) -
<T1, T2, R> Streamable<R> withLatestFrom(Streamable<T1> source1, Streamable<T2> source2, Function3<? super T, ? super T1, ? super T2, R> combiner) -
<T1, T2, T3, R> Streamable<R> withLatestFrom(Streamable<T1> source1, Streamable<T2> source2, Streamable<T3> source3, Function4<? super T, ? super T1, ? super T2, ? super T3, R> combiner) -
<T1, T2, T3, T4, R> Streamable<R> withLatestFrom(Streamable<T1> source1, Streamable<T2> source2, Streamable<T3> source3, Streamable<T4> source4, Function5<? super T, ? super T1, ? super T2, ? super T3, ? super T4, R> combiner)
zipArray
-
static <T, R> Streamable<R> zipArray(Streamable<? extends T>[] sources, Function<? super Object[], ? extends R> zipper, StandardBufferedConfig config)
zipWith
-
<U, R> Streamable<R> zipWith(Iterable<U> other, BiFunction<? super T, ? super U, ? extends R> zipper) -
<U, R> Streamable<R> zipWith(Streamable<? extends U> other, BiFunction<? super T, ? super U, ? extends R> zipper) -
<U, R> Streamable<R> zipWith(Streamable<? extends U> other, BiFunction<? super T, ? super U, ? extends R> zipper, StandardBufferedConfig config)
Probably not applicable to a pull-based Streamable (listed for completeness) (16 methods, 30 overloads)
These are backpressure-, Subscription-, or Flowable-plumbing specific. Tick them off as won't do or replace them with a Streamable-native design where one makes sense (e.g. create(VirtualGenerator) already covers virtualCreate, transform(VirtualTransformer) covers virtualTransform). Signatures in this section are kept verbatim from Flowable.
bufferSize
-
static int bufferSize()
create
-
static <T> Flowable<T> create(FlowableOnSubscribe<T> source, BackpressureStrategy mode)
doOnRequest
-
Flowable<T> doOnRequest(LongConsumer onRequest)
onBackpressureBuffer
-
Flowable<T> onBackpressureBuffer() -
Flowable<T> onBackpressureBuffer(OnBackpressureBufferConfig<? super T> config) -
Flowable<T> onBackpressureBuffer(long capacity, BackpressureOverflowStrategy overflowStrategy) -
Flowable<T> onBackpressureBuffer(long capacity, BackpressureOverflowStrategy overflowStrategy, Consumer<? super T> onDropped)
onBackpressureDrop
-
Flowable<T> onBackpressureDrop() -
Flowable<T> onBackpressureDrop(Consumer<? super T> onDrop)
onBackpressureLatest
-
Flowable<T> onBackpressureLatest() -
Flowable<T> onBackpressureLatest(Consumer<? super T> onDropped)
onBackpressureReduce
-
Flowable<T> onBackpressureReduce(BiFunction<T, T, T> reducer) -
<R> Flowable<R> onBackpressureReduce(Supplier<R> supplier, BiFunction<R, ? super T, R> reducer)
onTerminateDetach
-
Flowable<T> onTerminateDetach()
parallel
-
ParallelFlowable<T> parallel() -
ParallelFlowable<T> parallel(int parallelism) -
ParallelFlowable<T> parallel(int parallelism, int prefetch)
rebatchRequests
-
Flowable<T> rebatchRequests(int n)
safeSubscribe
-
void safeSubscribe(Subscriber<? super T> subscriber)
serialize
-
Flowable<T> serialize()
toStreamable
-
Streamable<T> toStreamable() -
Streamable<T> toStreamable(ExecutorService executor)
unsafeCreate
-
static <T> Flowable<T> unsafeCreate(Publisher<T> onSubscribe)
virtualCreate
-
static <T> Flowable<T> virtualCreate(VirtualGenerator<T> generator) -
static <T> Flowable<T> virtualCreate(VirtualGenerator<T> generator, ExecutorService executor) -
static <T> Flowable<T> virtualCreate(VirtualGenerator<T> generator, Scheduler scheduler)
virtualTransform
-
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer) -
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer, ExecutorService executor) -
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer, Scheduler scheduler, int prefetch) -
<R> Flowable<R> virtualTransform(VirtualTransformer<T, R> transformer, ExecutorService executor, int prefetch)
Beitragsleitfaden
Erste Schritte
- Lies das ganze Issue und danach den Beitragsleitfaden des Projekts.
- Schreib ins Issue, dass du es übernimmst — das erspart doppelte Arbeit.
- Forke das Repository und arbeite in einem Branch.
- Ö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