【问题标题】:RxJava - combine concat and merge but check ALL concated observablesRxJava - 结合 concat 和 merge 但检查所有连接的 observables
【发布时间】:2016-01-22 15:24:22
【问题描述】:

我正在从多个来源加载数据并缓存它们。

我所做的如下:

  • 从多个来源加载数据
  • 将结果缓存在内存中

我想要的是:

  1. 加载所有需要的数据
  2. 合并结果并仅在所有数据可用时传播结果

一个简单的例子如下:

mObservable1 = Observable
    .concat(mObservable1Cached, mObservable1)
    .first();

我怎样才能将许多缓存而不是像上面示例中那样缓存的 observables 组合在一起,只有 1 个?这是我的想法,但这行不通,因为一旦两个可观察对象中的一个具有缓存数据,它就会传播结果...

mObservable1And2  = Observable
    .concat(
            Observable.merge(mObservable1Cached, mObservable2Cached),
            Observable.merge(mObservable1, mObservable2)
    )
    .first();

Observable 示例

可观察的缓存数据

mObservable1Cached = Observable.create(new Observable.OnSubscribe<List<Data>>() {
    @Override
    public void call(Subscriber<? super List<Data>> subscriber) {
        subscriber.onNext(mData1Cached);
        subscriber.onCompleted();
    }
})
        .subscribeOn(HandlerScheduler.from(mBackgroundHandler))
        .observeOn(AndroidSchedulers.mainThread());

加载可观察的数据

mObservable1 = Observable.create(new Observable.OnSubscribe<List<Data>>() {
        @Override
        public void call(Subscriber<? super List<Data>> subscriber) {
            subscriber.onNext(...load data...);
            subscriber.onCompleted();
        }
    })
    .doOnNext(new Action1<List<Data>>() {
        @Override
        public void call(List<Data> data) {
            mData1Cached = data;
        }
    })
    .subscribeOn(HandlerScheduler.from(mBackgroundHandler))
    .observeOn(AndroidSchedulers.mainThread());

【问题讨论】:

  • Zip 运算符对您有帮助吗?这只会在所有被压缩的 Observable 发射后发射?

标签: java android merge concat rx-java


【解决方案1】:

您要求阻塞直到连接的可观察对象完成,这可以像这样实现:

mObservable1 = Observable
  .concat(mObservable1Cached, mObservable1)
  .toList()
  .flatMapIterable(x -> x)
  .first();

【讨论】:

    【解决方案2】:

    zip() 可以解决问题。

    这是完整的 (java8) 解决方案。 如果您受限于 java7,那么重构认为它会更加冗长是微不足道的。 同样使用 zip,您可以处理任意数量的来源(不仅仅是固定数量)。

      //Dummy data
      static class Data {
        static Map<String, AtomicInteger> dataVersion = new ConcurrentHashMap<>();
        int sourceId;
        int itemId;
        int version;
    
        Data(int sourceId, int itemId) {
          this.sourceId = sourceId;
          this.itemId = itemId;
          //for testing purpose we going to track how many items created for each source/item pair
          //we always expect one
          AtomicInteger versionCounter = dataVersion.computeIfAbsent(sourceId + "." + itemId, key->new AtomicInteger());
          this.version = versionCounter.incrementAndGet();
        }
        public String toString() {
          return sourceId + "." + itemId + "; version: " + version;
        }
      };
    
      //data cache (per source)
      ConcurrentMap<Integer, List<Data>> sourceCache = new ConcurrentHashMap<>();
    
      //data loader
      List<Data> loadData(int sourceId) {
        return Arrays.asList(new Data(sourceId,1), new Data(sourceId,2), new Data(sourceId,3));
      }
    
      //source observable factory method
      Observable<List<Data>> getSource(int sourceId) {
        return Observable.<List<Data>>create(subscriber->{
          subscriber.onNext(sourceCache.computeIfAbsent(sourceId, key->loadData(key)));
          subscriber.onCompleted();
        });
      }
    
      public void stackOverflow33296442() {
    
        Observable<List<Data>> result = Observable.zip(getSource(1), getSource(2), (src1,src2)->Arrays.asList(src1,src2))
            .concatMap(listOfLists->Observable.from(listOfLists));
    
        //Test
        Iterator<List<Data>> iter1 = result.toBlocking().toIterable().iterator();
        while (iter1.hasNext()) {
          List<Data> next = iter1.next();
          System.out.println("data: " + next);
          for (Data data : next) {
            assert 1 == data.version;
          }
        }
    
        Iterator<List<Data>> iter2 = result.toBlocking().toIterable().iterator();
        while (iter2.hasNext()) {
          List<Data> next = iter2.next();
          System.out.println("data: " + next);
          for (Data data : next) {
            assert 1 == data.version;
          }
        }
      }
    

    输出:

    data: [1.1; version: 1, 1.2; version: 1, 1.3; version: 1]
    data: [2.1; version: 1, 2.2; version: 1, 2.3; version: 1]
    data: [1.1; version: 1, 1.2; version: 1, 1.3; version: 1]
    data: [2.1; version: 1, 2.2; version: 1, 2.3; version: 1]
    

    【讨论】:

      【解决方案3】:

      这感觉类似于如何获取表单提交的 rx 模式。使用包含 2 个表单字段的 UI 时,您只想在两者都有值时启用提交按钮。 RxJava 允许你使用 combineLatest 操作符来做到这一点。每当一个 observable 发出一个项目时,它都会与另一个发出的 observable 中最后一个发出的项目结合。如果到目前为止只有一个 observable 发出了,它会等待另一个发出,然后传递这两个值。

      CombineLatest 是我对如何让它发挥作用的猜测。

      【讨论】:

        猜你喜欢
        • 2017-09-07
        • 1970-01-01
        • 1970-01-01
        • 2017-05-18
        • 2016-02-21
        • 1970-01-01
        • 2015-09-28
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多