【问题标题】:rxjava2 producer-consumer, 'downstream' request one, the 'upstream' emit onerxjava2 producer-consumer,“下游”请求一,“上游”发出一
【发布时间】:2017-12-01 09:01:59
【问题描述】:

我想用 rxjava2 实现简单的生产者-消费者模型,当下游请求一个,上游发出一个。

我知道flatMapobserveOn 的默认缓冲区大小为128,所以我将缓冲区大小设置为1,但它也不起作用。

Flowable.defer((Callable<Publisher<Integer>>) () -> Flowable.range(1, 5))
            .flatMap((Function<Integer, Publisher<Integer>>) integer -> {
                //do something with long time.
                System.out.println("flatMap:" + integer);
                return Flowable.just(integer);
            }, false, 1) //=====> 1
            .subscribeOn(Schedulers.io())
            .observeOn(Schedulers.computation(), false, 1) //=====> 2
            .subscribe(new Subscriber<Integer>() {
                @Override
                public void onSubscribe(Subscription s) {
                    //request one
                    s.request(1);
                }

                @Override
                public void onNext(Integer integer) {
                    System.out.println("onNext:" + integer);
                }

                @Override
                public void onError(Throwable t) {

                }

                @Override
                public void onComplete() {

                }
            });

实际输出:

flatMap:1
flatMap:2
onNext:1
flatMap:3

预期的输出,因为我只调用了一次s.request(1)

flatMap:1
onNext:1

【问题讨论】:

    标签: rx-java rx-android rx-java2


    【解决方案1】:

    您的观察者只请求一项,但observeOn() 也会缓冲一项。 flatMap() 运算符本身将订阅连续的输入。

    1. 观察者订阅观察者链,并请求 1 项。
    2. observeOn() 为其缓冲区请求 1 个项目。
    3. range() 运算符发出 1。
    4. flatMap() 收到 1,并在内部订阅 flowable,导致第一条日志行。
    5. observeOn() 为其缓冲区获取一项,然后请求另一项。
    6. flatMap() 获取下一项,2。这被发射并传递给 observeOn() 缓冲区
    7. 观察者的onNext()被调用。
    8. flatMap() 获取下一项,3。

    如果您需要完美的锁定步骤,​​“请求一”->“流程一”,那么流程控制不是解决问题的方法。相反,您可能希望引入一个提供反馈循环的 observable,以便观察者告诉 observable 处理下一个。

    【讨论】:

    • 感谢您的详细步骤。我也认为flowable 做不到,所以我会尝试你的“反馈循环”方式。
    猜你喜欢
    • 1970-01-01
    • 2020-05-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-07-13
    • 1970-01-01
    • 2020-06-30
    • 2021-07-10
    相关资源
    最近更新 更多