【发布时间】:2017-06-15 07:56:40
【问题描述】:
我想知道如何以类似于 switchMap 的方式转换可观察对象,但不是限制为单个活动流,而是具有多个(有限)流。
目的是让多个任务同时工作,达到一定的任务计数限制,并允许新任务以先进先出队列策略启动,这意味着任何新任务到达将立即启动,队列中最旧的任务将被取消。
switchMap 将为源的每次发射创建 Observable,并在创建新的 Observable 流后取消之前运行的 Observable 流,我想实现类似但允许某种级别的并发(如 flatMap),这意味着允许创建多个 Observable对于每次发射,并发运行到某个并发限制,当达到并发限制时,最旧的 observable 将被取消,新的 observable 将启动。
其实这也和 flatMap 的 maxConcurrent 类似,只是当达到 maxConcurrent 时,新的 Observables 不会在队列中等待,而是取消旧的 Observables 并立即进入新的 Observables。
【问题讨论】:
-
我不相信 rxJava(或我熟悉的任何 RX 实现)中存在这样的运算符。如果您有一个好的用例,您可以继续提交功能请求here。
-
嗯,这么想,我会尝试自己想出一些东西,并分享它,看看我是否在正确的方向
-
这是用于 rxjava 1 还是 2?
-
它是 rxjava 1 的,已经看过 rxjava-extras Dave,也许我错过了?
标签: concurrency rx-java