Java 类io.reactivex.observables.GroupedObservable 实例源码
项目:code-examples-android-expert
文件:lessonD_AdvancedStreams.java
@Test
public void shouldBeNecessaryToSubscribetoStreamAfterSplitting() {
final double[] averages = {0, 0};
Observable<Integer> numbers = Observable.just(22, 22, 99, 22, 101, 22);
Function<Integer, Integer> keySelector = integer -> integer % 2;
Observable<GroupedObservable<Integer, Integer>> split = numbers.groupBy(keySelector);
split.subscribe(
group -> {
Observable<Double> convertToDouble = group.map(integer -> (double) integer);
Function<Double, Double> insertIntoAveragesArray = aDouble -> averages[group.getKey()] = aDouble;
convertToDouble.reduce((t1, t2) -> t1+t2).map(insertIntoAveragesArray).subscribe();
}
);
assertThat(averages[0]).isEqualTo(0);
assertThat(averages[1]).isEqualTo(0);
}
项目:RxJava2-Android-Sample
文件:GroupByExampleActivity.java
private void doSomeWork() {
Observable.range(0, 8).groupBy(new Function<Integer, String>() {
@Override
public String apply(@NonNull Integer integer) throws Exception {
return integer % 2 == 0 ? "偶数" : "奇数";
}
}).subscribe(new Consumer<GroupedObservable<String, Integer>>() {
@Override
public void accept(@NonNull GroupedObservable<String, Integer> stringIntegerGroupedObservable) throws Exception {
String key = stringIntegerGroupedObservable.getKey();
Log.i(TAG, "accept: key=" + key);
if (key.equals("偶数")) {
stringIntegerGroupedObservable.subscribe(getObserver(key));
} else {
stringIntegerGroupedObservable.subscribe(getObserver(key));
}
}
});
}
项目:RxWindowIfChanged
文件:WindowIfChangedTest.java
@Test public void completeCompletesInner() {
Observable<Message> messages = Observable.just(new Message("Bob", "Hello"));
final AtomicInteger seen = new AtomicInteger();
WindowIfChanged.create(messages, userSelector)
.switchMap(
new Function<GroupedObservable<String, Message>, Observable<Notification<String>>>() {
@Override public Observable<Notification<String>> apply(
GroupedObservable<String, Message> group) {
final int count = seen.incrementAndGet();
return group.map(new Function<Message, String>() {
@Override public String apply(Message message) throws Exception {
return count + " " + message;
}
}).materialize();
}
})
.test()
.assertValues( //
Notification.createOnNext("1 Bob Hello"), //
Notification.<String>createOnComplete()) //
.assertComplete();
}
项目:RxWindowIfChanged
文件:WindowIfChangedTest.java
@Test public void errorCompletesInner() {
RuntimeException error = new RuntimeException("boom!");
Observable<Message> messages = Observable.just( //
Notification.createOnNext(new Message("Bob", "Hello")),
Notification.createOnError(error)
).dematerialize();
final AtomicInteger seen = new AtomicInteger();
WindowIfChanged.create(messages, userSelector)
.switchMap(
new Function<GroupedObservable<String, Message>, Observable<Notification<String>>>() {
@Override public Observable<Notification<String>> apply(
GroupedObservable<String, Message> group) {
final int count = seen.incrementAndGet();
return group.map(new Function<Message, String>() {
@Override public String apply(Message message) throws Exception {
return count + " " + message;
}
}).materialize();
}
})
.test()
.assertValues( //
Notification.createOnNext("1 Bob Hello"), //
Notification.<String>createOnComplete()) //
.assertError(error);
}
项目:Learning-RxJava
文件:Ch4_19.java
public static void main(String[] args) {
Observable<String> source =
Observable.just("Alpha", "Beta", "Gamma", "Delta",
"Epsilon");
Observable<GroupedObservable<Integer, String>> byLengths =
source.groupBy(s -> s.length());
byLengths.flatMapSingle(grp ->
grp.reduce("", (x, y) -> x.equals("") ? y : x + ", " + y)
.map(s -> grp.getKey() + ": " + s)
).subscribe(System.out::println);
}
项目:Learning-RxJava
文件:Ch4_18.java
public static void main(String[] args) {
Observable<String> source =
Observable.just("Alpha", "Beta", "Gamma", "Delta",
"Epsilon");
Observable<GroupedObservable<Integer, String>> byLengths =
source.groupBy(s -> s.length());
byLengths.flatMapSingle(grp -> grp.toList())
.subscribe(System.out::println);
}
项目:javarx-study
文件:ObservableFunctionalTest.java
@Test
public void testBasicGroupByObservable() throws InterruptedException {
Observable<GroupedObservable<String, Integer>> grouped =
Observable.range(1, 100)
.groupBy(integer -> {
if (integer % 2 == 0) return "Even";
else return "Odd";
});
grouped.subscribe(g -> g.subscribe(x ->
System.out.println
("g:" + g.getKey() + ", value:" + x)));
Thread.sleep(4000);
}
项目:RxWindowIfChanged
文件:WindowIfChangedTest.java
@Test public void splits() {
Observable<Message> messages = Observable.just( //
new Message("Bob", "Hello"), //
new Message("Bob", "World"), //
new Message("Alice", "Hey"), //
new Message("Bob", "What's"), //
new Message("Bob", "Up?"), //
new Message("Eve", "Hey") //
);
final AtomicInteger seen = new AtomicInteger();
WindowIfChanged.create(messages, userSelector)
.switchMap(
new Function<GroupedObservable<String, Message>, Observable<Notification<String>>>() {
@Override public Observable<Notification<String>> apply(
GroupedObservable<String, Message> group) {
final int count = seen.incrementAndGet();
return group.map(new Function<Message, String>() {
@Override public String apply(Message message) throws Exception {
return count + " " + message;
}
}).materialize();
}
})
.test()
.assertValues( //
Notification.createOnNext("1 Bob Hello"), //
Notification.createOnNext("1 Bob World"), //
Notification.<String>createOnComplete(), //
Notification.createOnNext("2 Alice Hey"), //
Notification.<String>createOnComplete(), //
Notification.createOnNext("3 Bob What's"), //
Notification.createOnNext("3 Bob Up?"), //
Notification.<String>createOnComplete(), //
Notification.createOnNext("4 Eve Hey"), //
Notification.<String>createOnComplete()); //
}
项目:Learning-RxJava
文件:Ch4_17.java
public static void main(String[] args) {
Observable<String> source =
Observable.just("Alpha", "Beta", "Gamma", "Delta", "Epsilon");
Observable<GroupedObservable<Integer, String>> byLengths =
source.groupBy(s -> s.length());
}
项目:RxWindowIfChanged
文件:WindowIfChangedObservable.java
@Override protected void subscribeActual(Observer<? super GroupedObservable<K, T>> observer) {
upstream.subscribe(new WindowIfChangedObserver<>(keySelector, observer));
}
项目:RxWindowIfChanged
文件:WindowIfChangedObserver.java
WindowIfChangedObserver(Function<? super T, ? extends K> keySelector,
Observer<? super GroupedObservable<K, T>> observer) {
this.keySelector = keySelector;
this.observer = observer;
}
项目:RxWindowIfChanged
文件:WindowIfChanged.java
public static <T, K> Observable<GroupedObservable<K, T>> create(Observable<T> upstream,
Function<? super T, ? extends K> keySelector) {
return new WindowIfChangedObservable<>(upstream, keySelector);
}
项目:cyclops
文件:ObservableKind.java
@CheckReturnValue
@SchedulerSupport("none")
public <K> Observable<GroupedObservable<K, T>> groupBy(Function<? super T, ? extends K> keySelector) {
return boxed.groupBy(keySelector);
}
项目:cyclops
文件:ObservableKind.java
@CheckReturnValue
@SchedulerSupport("none")
public <K> Observable<GroupedObservable<K, T>> groupBy(Function<? super T, ? extends K> keySelector, boolean delayError) {
return boxed.groupBy(keySelector, delayError);
}
项目:cyclops
文件:ObservableKind.java
@CheckReturnValue
@SchedulerSupport("none")
public <K, V> Observable<GroupedObservable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector) {
return boxed.groupBy(keySelector, valueSelector);
}
项目:cyclops
文件:ObservableKind.java
@CheckReturnValue
@SchedulerSupport("none")
public <K, V> Observable<GroupedObservable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, boolean delayError) {
return boxed.groupBy(keySelector, valueSelector, delayError);
}
项目:cyclops
文件:ObservableKind.java
@CheckReturnValue
@SchedulerSupport("none")
public <K, V> Observable<GroupedObservable<K, V>> groupBy(Function<? super T, ? extends K> keySelector, Function<? super T, ? extends V> valueSelector, boolean delayError, int bufferSize) {
return boxed.groupBy(keySelector, valueSelector, delayError, bufferSize);
}