【问题标题】:RxJava zip operation on a varying length of Retrofit Observable array对不同长度的 Retrofit Observable 数组进行 RxJava zip 操作
【发布时间】:2017-07-13 17:44:32
【问题描述】:

我有不同长度的 Observable 数组。我想压缩请求(即发出一堆 API 请求并等到所有请求都完成),但我不知道如何实现 zip 功能。

Observable.zip(observables, new FuncN<List<ResponseBody>>() {
    @Override
    public List<ResponseBody> call(Object... args) {
        return Arrays.asList(args); <- compile error here
    }
});

这里的obserablesList&lt;Observable&lt;ResponseBody&gt;&gt; 的数组,它的长度是先验未知的。

zip函数中调用的参数不能更正为ResponseBody...。如何让它返回Observable&lt;List&lt;ResponseBody&gt;&gt;

FuncNRxJava 1.x.x 的设计中是否有约束?

附:我正在使用 RxJava 1.1.6

【问题讨论】:

  • 到底是什么问题? FuncN 方法返回的类型是将由创建的运算符使用 zip 方法发出的项目类型,这意味着在这里您将拥有 Observable>
  • @yosriz,问题是Arrays.asList(args)不是List&lt;ResponseBody&gt;类型,导致编译出错。

标签: rx-java observable retrofit2 rx-android


【解决方案1】:

只需 merge 您的 observables 并使用 toList 收集结果:

Observable.merge(observables).toList()

【讨论】:

  • 请问merge 操作是否会阻止每个请求?我希望同时触发请求异步,以便请求不会相互阻塞。
  • merge 不会阻止,但Observable.concat
  • 我只是试了一下,我发现请求是一个一个发送的,即在第二个请求被触发之前只有第一个请求有响应...
  • 您的意思是您收到了这个订单:1st request, 1st response, 2nd request, 2nd response, 3rd request, 3rd response...?
  • 是的,我在使用合并时得到了这个订单。
【解决方案2】:

我确认merge 操作员工作并找到导致订单的罪魁祸首:1st request, 1st response, 2nd request, 2nd response, 3rd request, 3rd response

observables 列表中的 observable 在添加过程中未在Scheduler.io() 上订阅。

以前:

observables.add(mViewModelDelegate.get().getApiService()
                    .rxGetCategoryBrandList(baseUrl, categoryId));

之后:

observables.add(mViewModelDelegate.get().getApiService()
                    .rxGetCategoryBrandList(baseUrl, categoryId)
                    .subscribeOn(Schedulers.io())
                    .observeOn(AndroidSchedulers.mainThread());

【讨论】:

    【解决方案3】:

    用户合并并确保所有可观察的请求都有自己的线程。 你可以试试这个:

    private void runMyTest() {
        List<Single<String>> singleObservableList = new ArrayList<>();
        singleObservableList.add(getSingleObservable(500, "AAA"));
        singleObservableList.add(getSingleObservable(300, "BBB"));
        singleObservableList.add(getSingleObservable(100, "CCC"));
        Single.merge(singleObservableList)
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe(System.out::println);
    }
    
    private Single<String> getSingleObservable(long waitMilliSeconds, String name) {
        return Single
                .create((SingleOnSubscribe<String>) e -> {
                        try {
                            Thread.sleep(waitMilliSeconds);
                        } catch (InterruptedException exception) {
                            exception.printStackTrace();
                        }
                        System.out.println("name = " +name+ ", waitMilliSeconds = " +waitMilliSeconds+ ", thread name = " +Thread.currentThread().getName()+ ", id =" +Thread.currentThread().getId());
                        if(!e.isDisposed()) e.onSuccess(name);
                    })
                .subscribeOn(Schedulers.io());
    }
    

    输出:

    System.out:名称 = CCC,waitMilliSeconds = 100,线程名称 = RxCachedThreadScheduler-4, id =463

    System.out: CCC

    System.out:名称 = BBB,waitMilliSeconds = 300,线程名称 = RxCachedThreadScheduler-3, id =462

    System.out: BBB

    System.out:名称 = AAA,waitMilliSeconds = 500,线程名称 = RxCachedThreadScheduler-2, id =461

    System.out:AAA

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-03-28
      • 2017-12-05
      • 1970-01-01
      • 2020-04-05
      • 2017-02-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多