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);
}