【问题标题】:Infinite observable from another observable来自另一个可观察对象的无限可观察对象
【发布时间】:2015-05-15 14:53:28
【问题描述】:

我有一个表示从数据库表中选择的序列的 Observable, 所以它是有限的。

Observable<Item> selectResults() { ... }

我想实现一个指定间隔的拉动,所以最后我会得到另一个可观察的对象,它将包裹我原来的对象并无限期地拉动。

我只是不知道该怎么做:(


好的,这是我的想法,围绕可观察的区间建模,可能需要错误处理和取消订阅逻辑。

public class OnSubscribePeriodicObservable implements OnSubscribe<Item> {
...

  @Override
  public void call(final Subscriber<? super Item> subscriber) {
      final Worker worker = scheduler.createWorker();
      subscriber.add( worker );

      worker.schedulePeriodically(new Action0() {
          @Override
          public void call() {
    selectResults().subscribe( new Observer<Item>() {
            @Override
            public void onCompleted() {
              //continue
            }

            @Override
            public void onError(Throwable e) {
              subscriber.onError( e );
            }

            @Override
            public void onNext(Item t) {
              subscriber.onNext( t );
            }
         }); 
          }

      }, initialDelay, period, unit);
}

【问题讨论】:

    标签: java rx-java


    【解决方案1】:

    您可以使用标准运算符完成此操作,这将为您提供错误传播、取消订阅和廉价的背压:

    Observable<Integer> databaseQuery = Observable
        .just(1, 2, 3, 4)
        .delay(500, TimeUnit.MILLISECONDS);
    
    Observable<Integer> result = Observable
            .timer(1, 2, TimeUnit.SECONDS)
            .onBackpressureDrop()
            .concatMap(t -> databaseQuery);
    
    result.subscribe(System.out::println);
    
    Thread.sleep(10000);
    

    【讨论】:

      猜你喜欢
      • 2023-04-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-08-19
      • 2017-07-12
      • 1970-01-01
      • 2019-02-18
      相关资源
      最近更新 更多