io.reactivex.Observable.concatMapEager()方法的使用及代码示例

x33g5p2x  于2022-01-25 转载在 其他  
字(7.3k)|赞(0)|评价(0)|浏览(144)

本文整理了Java中io.reactivex.Observable.concatMapEager()方法的一些代码示例,展示了Observable.concatMapEager()的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Observable.concatMapEager()方法的具体详情如下:
包路径:io.reactivex.Observable
类名称:Observable
方法名:concatMapEager

Observable.concatMapEager介绍

[英]Maps a sequence of values into ObservableSources and concatenates these ObservableSources eagerly into a single ObservableSource.

Eager concatenation means that once a subscriber subscribes, this operator subscribes to all of the source ObservableSources. The operator buffers the values emitted by these ObservableSources and then drains them in order, each one after the previous one completes.

Scheduler: This method does not operate by default on a particular Scheduler.
[中]将一系列值映射到可观测资源中,并将这些可观测资源急切地连接到单个可观测资源中。
即时连接意味着一旦订户订阅,该操作符订阅所有源可观测资源。运算符缓冲这些可观察资源发出的值,然后依次将其耗尽,每一个都在前一个完成之后。
调度器:默认情况下,该方法不会在特定的调度器上运行。

代码示例

代码示例来源:origin: ReactiveX/RxJava

@Override
  public ObservableSource<Object> apply(Observable<Object> o) throws Exception {
    return o.concatMapEager(new Function<Object, ObservableSource<Object>>() {
      @Override
      public ObservableSource<Object> apply(Object v) throws Exception {
        return Observable.just(v);
      }
    });
  }
});

代码示例来源:origin: ReactiveX/RxJava

@Test(expected = IllegalArgumentException.class)
public void testInvalidMaxConcurrent() {
  Observable.just(1).concatMapEager(toJust, 0, Observable.bufferSize());
}

代码示例来源:origin: ReactiveX/RxJava

@SuppressWarnings({ "unchecked", "rawtypes" })
@Test
public void mappingBadCapacityHint() throws Exception {
  Observable<Integer> source = Observable.just(1);
  try {
    Observable.just(source, source, source).concatMapEager((Function)Functions.identity(), 10, -99);
  } catch (IllegalArgumentException ex) {
    assertEquals("prefetch > 0 required but it was -99", ex.getMessage());
  }
}

代码示例来源:origin: ReactiveX/RxJava

@Test(expected = IllegalArgumentException.class)
public void testInvalidCapacityHint() {
  Observable.just(1).concatMapEager(toJust, Observable.bufferSize(), 0);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void dispose() {
  TestHelper.checkDisposed(Observable.just(1).hide().concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      return Observable.range(1, 2);
    }
  }));
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void testMapperThrows() {
  Observable.just(1).concatMapEager(new Function<Integer, Observable<Integer>>() {
    @Override
    public Observable<Integer> apply(Integer t) {
      throw new TestException();
    }
  }).subscribe(to);
  to.assertNoValues();
  to.assertNotComplete();
  to.assertError(TestException.class);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void oneDelayed() {
  Observable.just(1, 2, 3, 4, 5)
  .concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer i) throws Exception {
      return i == 3 ? Observable.just(i) : Observable
          .just(i)
          .delay(1, TimeUnit.MILLISECONDS, Schedulers.io());
    }
  })
  .observeOn(Schedulers.io())
  .test()
  .awaitDone(5, TimeUnit.SECONDS)
  .assertResult(1, 2, 3, 4, 5)
  ;
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void mapperCancels() {
  final TestObserver<Integer> to = new TestObserver<Integer>();
  Observable.just(1).hide()
  .concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      to.cancel();
      return Observable.never();
    }
  }, 1, 128)
  .subscribe(to);
  to.assertEmpty();
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void innerErrorMaxConcurrency() {
  Observable.<Integer>just(1).hide().concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      return Observable.error(new TestException());
    }
  }, 1, 128)
  .test()
  .assertFailure(TestException.class);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void innerError() {
  Observable.<Integer>just(1).hide().concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      return Observable.error(new TestException());
    }
  })
  .test()
  .assertFailure(TestException.class);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
@Ignore("Null values are not allowed in RS")
public void testInnerNull() {
  Observable.just(1).concatMapEager(new Function<Integer, Observable<Integer>>() {
    @Override
    public Observable<Integer> apply(Integer t) {
      return Observable.just(null);
    }
  }).subscribe(to);
  to.assertNoErrors();
  to.assertComplete();
  to.assertValue(null);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void innerErrorFused() {
  Observable.<Integer>just(1).hide().concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      return Observable.range(1, 2).map(new Function<Integer, Integer>() {
        @Override
        public Integer apply(Integer v) throws Exception {
          throw new TestException();
        }
      });
    }
  })
  .test()
  .assertFailure(TestException.class);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void testAsynchronousRun() {
  Observable.range(1, 2).concatMapEager(new Function<Integer, Observable<Integer>>() {
    @Override
    public Observable<Integer> apply(Integer t) {
      return Observable.range(1, 1000).subscribeOn(Schedulers.computation());
    }
  }).observeOn(Schedulers.newThread()).subscribe(to);
  to.awaitTerminalEvent(5, TimeUnit.SECONDS);
  to.assertNoErrors();
  to.assertValueCount(2000);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void testSimple() {
  Observable.range(1, 100).concatMapEager(toJust).subscribe(to);
  to.assertNoErrors();
  to.assertValueCount(100);
  to.assertComplete();
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void testSimple2() {
  Observable.range(1, 100).concatMapEager(toRange).subscribe(to);
  to.assertNoErrors();
  to.assertValueCount(200);
  to.assertComplete();
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void innerCallableThrows() {
  Observable.<Integer>just(1).hide().concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      return Observable.fromCallable(new Callable<Integer>() {
        @Override
        public Integer call() throws Exception {
          throw new TestException();
        }
      });
    }
  })
  .test()
  .assertFailure(TestException.class);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void normal() {
  Observable.range(1, 5)
  .concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer t) {
      return Observable.range(t, 2);
    }
  })
  .test()
  .assertResult(1, 2, 2, 3, 3, 4, 4, 5, 5, 6);
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void testMainError() {
  Observable.<Integer>error(new TestException()).concatMapEager(toJust).subscribe(to);
  to.assertNoValues();
  to.assertError(TestException.class);
  to.assertNotComplete();
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void empty() {
  Observable.<Integer>empty().hide().concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) throws Exception {
      return Observable.range(1, 2);
    }
  })
  .test()
  .assertResult();
}

代码示例来源:origin: ReactiveX/RxJava

@Test
public void longEager() {
  Observable.range(1, 2 * Observable.bufferSize())
  .concatMapEager(new Function<Integer, ObservableSource<Integer>>() {
    @Override
    public ObservableSource<Integer> apply(Integer v) {
      return Observable.just(1);
    }
  })
  .test()
  .assertValueCount(2 * Observable.bufferSize())
  .assertNoErrors()
  .assertComplete();
}

相关文章

Observable类方法