-
Notifications
You must be signed in to change notification settings - Fork 7.6k
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
2.x: add Flowable.parallel() and parallel operators #4974
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,165 @@ | ||
/** | ||
* Copyright (c) 2016-present, RxJava Contributors. | ||
* | ||
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in | ||
* compliance with the License. You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software distributed under the License is | ||
* distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See | ||
* the License for the specific language governing permissions and limitations under the License. | ||
*/ | ||
|
||
package io.reactivex.internal.operators.parallel; | ||
|
||
import java.util.concurrent.Callable; | ||
|
||
import org.reactivestreams.*; | ||
|
||
import io.reactivex.exceptions.Exceptions; | ||
import io.reactivex.functions.BiConsumer; | ||
import io.reactivex.internal.subscribers.DeferredScalarSubscriber; | ||
import io.reactivex.internal.subscriptions.*; | ||
import io.reactivex.parallel.ParallelFlowable; | ||
import io.reactivex.plugins.RxJavaPlugins; | ||
|
||
/** | ||
* Reduce the sequence of values in each 'rail' to a single value. | ||
* | ||
* @param <T> the input value type | ||
* @param <C> the collection type | ||
*/ | ||
public final class ParallelCollect<T, C> extends ParallelFlowable<C> { | ||
|
||
final ParallelFlowable<? extends T> source; | ||
|
||
final Callable<C> initialCollection; | ||
|
||
final BiConsumer<C, T> collector; | ||
|
||
public ParallelCollect(ParallelFlowable<? extends T> source, | ||
Callable<C> initialCollection, BiConsumer<C, T> collector) { | ||
this.source = source; | ||
this.initialCollection = initialCollection; | ||
this.collector = collector; | ||
} | ||
|
||
@Override | ||
public void subscribe(Subscriber<? super C>[] subscribers) { | ||
if (!validate(subscribers)) { | ||
return; | ||
} | ||
|
||
int n = subscribers.length; | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Same in other similar places would be good. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I don't do those unless the variable has to be accessed from an inner class. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. In long methods reader has to spend extra time to check that it's not modified anywhere, but ok |
||
@SuppressWarnings("unchecked") | ||
Subscriber<T>[] parents = new Subscriber[n]; | ||
|
||
for (int i = 0; i < n; i++) { | ||
|
||
C initialValue; | ||
|
||
try { | ||
initialValue = initialCollection.call(); | ||
} catch (Throwable ex) { | ||
Exceptions.throwIfFatal(ex); | ||
reportError(subscribers, ex); | ||
return; | ||
} | ||
|
||
if (initialValue == null) { | ||
reportError(subscribers, new NullPointerException("The initialSupplier returned a null value")); | ||
return; | ||
} | ||
|
||
parents[i] = new ParallelCollectSubscriber<T, C>(subscribers[i], initialValue, collector); | ||
} | ||
|
||
source.subscribe(parents); | ||
} | ||
|
||
void reportError(Subscriber<?>[] subscribers, Throwable ex) { | ||
for (Subscriber<?> s : subscribers) { | ||
EmptySubscription.error(ex, s); | ||
} | ||
} | ||
|
||
@Override | ||
public int parallelism() { | ||
return source.parallelism(); | ||
} | ||
|
||
static final class ParallelCollectSubscriber<T, C> extends DeferredScalarSubscriber<T, C> { | ||
|
||
|
||
private static final long serialVersionUID = -4767392946044436228L; | ||
|
||
final BiConsumer<C, T> collector; | ||
|
||
C collection; | ||
|
||
boolean done; | ||
|
||
ParallelCollectSubscriber(Subscriber<? super C> subscriber, | ||
C initialValue, BiConsumer<C, T> collector) { | ||
super(subscriber); | ||
this.collection = initialValue; | ||
this.collector = collector; | ||
} | ||
|
||
@Override | ||
public void onSubscribe(Subscription s) { | ||
if (SubscriptionHelper.validate(this.s, s)) { | ||
this.s = s; | ||
|
||
actual.onSubscribe(this); | ||
|
||
s.request(Long.MAX_VALUE); | ||
} | ||
} | ||
|
||
@Override | ||
public void onNext(T t) { | ||
if (done) { | ||
return; | ||
} | ||
|
||
try { | ||
collector.accept(collection, t); | ||
} catch (Throwable ex) { | ||
Exceptions.throwIfFatal(ex); | ||
cancel(); | ||
onError(ex); | ||
return; | ||
} | ||
} | ||
|
||
@Override | ||
public void onError(Throwable t) { | ||
if (done) { | ||
RxJavaPlugins.onError(t); | ||
return; | ||
} | ||
done = true; | ||
collection = null; | ||
actual.onError(t); | ||
} | ||
|
||
@Override | ||
public void onComplete() { | ||
if (done) { | ||
return; | ||
} | ||
done = true; | ||
C c = collection; | ||
collection = null; | ||
complete(c); | ||
} | ||
|
||
@Override | ||
public void cancel() { | ||
super.cancel(); | ||
s.cancel(); | ||
} | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,88 @@ | ||
/** | ||
* Copyright (c) 2016-present, RxJava Contributors. | ||
* | ||
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in | ||
* compliance with the License. You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software distributed under the License is | ||
* distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See | ||
* the License for the specific language governing permissions and limitations under the License. | ||
*/ | ||
|
||
package io.reactivex.internal.operators.parallel; | ||
|
||
import org.reactivestreams.*; | ||
|
||
import io.reactivex.functions.Function; | ||
import io.reactivex.internal.functions.ObjectHelper; | ||
import io.reactivex.internal.operators.flowable.FlowableConcatMap; | ||
import io.reactivex.internal.util.ErrorMode; | ||
import io.reactivex.parallel.ParallelFlowable; | ||
|
||
/** | ||
* Concatenates the generated Publishers on each rail. | ||
* | ||
* @param <T> the input value type | ||
* @param <R> the output value type | ||
*/ | ||
public final class ParallelConcatMap<T, R> extends ParallelFlowable<R> { | ||
|
||
final ParallelFlowable<T> source; | ||
|
||
final Function<? super T, ? extends Publisher<? extends R>> mapper; | ||
|
||
final int prefetch; | ||
|
||
final ErrorMode errorMode; | ||
|
||
public ParallelConcatMap( | ||
ParallelFlowable<T> source, | ||
Function<? super T, ? extends Publisher<? extends R>> mapper, | ||
int prefetch, ErrorMode errorMode) { | ||
this.source = source; | ||
this.mapper = ObjectHelper.requireNonNull(mapper, "mapper"); | ||
this.prefetch = prefetch; | ||
this.errorMode = ObjectHelper.requireNonNull(errorMode, "errorMode"); | ||
} | ||
|
||
@Override | ||
public int parallelism() { | ||
return source.parallelism(); | ||
} | ||
|
||
@Override | ||
public void subscribe(Subscriber<? super R>[] subscribers) { | ||
if (!validate(subscribers)) { | ||
return; | ||
} | ||
|
||
int n = subscribers.length; | ||
|
||
@SuppressWarnings("unchecked") | ||
final Subscriber<T>[] parents = new Subscriber[n]; | ||
|
||
// FIXME cheat until we have support from RxJava2 internals | ||
Publisher<T> p = new Publisher<T>() { | ||
int i; | ||
|
||
@SuppressWarnings("unchecked") | ||
@Override | ||
public void subscribe(Subscriber<? super T> s) { | ||
parents[i++] = (Subscriber<T>)s; | ||
} | ||
}; | ||
|
||
FlowableConcatMap<T, R> op = new FlowableConcatMap<T, R>(p, mapper, prefetch, errorMode); | ||
|
||
for (int i = 0; i < n; i++) { | ||
|
||
op.subscribe(subscribers[i]); | ||
// FIXME needs a FlatMap subscriber | ||
// parents[i] = FlowableConcatMap.createSubscriber(s, mapper, prefetch, errorMode); | ||
} | ||
|
||
source.subscribe(parents); | ||
} | ||
} |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
What about remove
runOn
and addScheduler
as a parameter toparallel()
?There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
It is the same logic as with regular factory methods such as
just
,range
,fromIterable
don't take aScheduler
, plus you can apply multiplerunOn
's on a sequence at different stages. For example create a pipeline with stages of parallelism=2 and 3 stages in total.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Ah, that's nice, got it.