【问题标题】:how to divide the flow of events in the two emitted simultaneously and process them?如何划分两个同时发出的事件流并进行处理?
【发布时间】:2016-08-26 07:07:22
【问题描述】:

有一个随机数流。

rx.Observable
.range (0, 1000)
.map (() -> 200d * Math.random ())

我需要将流程分成两部分。小于 100 的数字和大于 100 的数字。

之后,对于小于 100 的数字(链 1): 我需要向网络执行request1,等待答复并继续其他运营商的流程chain1。

对于大于 100 的数字(链 2): 我必须再发一个request2,等待答复,然后继续处理链操作。

request1request2 不会互相等待,链是并行执行的。但是内链处理必须等待请求的响应。

怎么做?

【问题讨论】:

    标签: java android rx-java


    【解决方案1】:
    rx.Observable
                            .create(subscriber -> {
                                for (int i = 0; i < 100; i++) {
                                    subscriber.onNext(i);
                                    Log.i("Iniop", "Create thread name: " + Thread.currentThread().getName());
                                }
                                subscriber.onCompleted();
                            })
                            .onBackpressureBuffer()
                            .observeOn(Schedulers.computation())
                            .subscribeOn(Schedulers.computation())
                            .map(v -> {
                                Log.i("Iniop", "Map thread name: " + Thread.currentThread().getName());
                                return 200d * Math.random();
                            })
                            .groupBy(k -> {
                                        Log.i("Iniop", "Group thread name: " + Thread.currentThread().getName());
                                        return k > 100 ? "yes" : "no";
                                    }
                                    , v -> v)
                            .forEach(gO -> gO.observeOn(Schedulers.newThread())
                                            .map(v -> new Pair<String, Double>(gO.getKey(), v))
                                            .subscribe(v -> {
                                                        Log.i("Iniop", "Key: " + v.first);
                                                        Log.i("Iniop", "Value: " + v.second);
                                                        Log.i("Iniop", "Thread name: " + Thread.currentThread().getName());
                                                    }
                                                    , e -> Log.e("Iniop", "Err", e))
                                    , e -> Log.e("Iniop", "Err", e));
    

    【讨论】:

    • 您可能希望将 subscribeOn 移动到 mapgroupBy 运算符之前,以使它们在计算调度程序上运行。
    猜你喜欢
    • 2017-12-28
    • 1970-01-01
    • 2021-02-14
    • 2021-02-20
    • 1970-01-01
    • 1970-01-01
    • 2011-11-01
    • 2020-11-11
    相关资源
    最近更新 更多