如何在flatMapSingle中指定maxConcurrency以限制并发?
实现flatMapSingle的并发限制:官方重载+手动实现思路
当然可以实现你说的并发限制逻辑!而且其实RxJava本身就给flatMapSingle提供了带maxConcurrency参数的重载方法,这是最直接省心的方案,先给你看怎么改你的代码:
方案1:直接使用flatMapSingle的官方重载
RxJava的flatMapSingle有一个重载签名是:
<R> Observable<R> flatMapSingle(Function<? super T, ? extends SingleSource<? extends R>> mapper, int maxConcurrency)
你只需要在调用时传入想要的最大并发数就行,比如设为3:
playersIdsObservable .distinct() .flatMapSingle(playerId -> dao.players().loadPlayerSRx(playerId), 3) // 限制同时最多3个请求
这个方法内部已经帮你实现了「保持maxConcurrency个并发任务,完成一个就补一个新任务」的逻辑,完全符合你描述的思路,不用自己手动管理。
方案2:手动实现预请求+补位逻辑(如果你想自定义控制)
如果你想自己手动实现这个逻辑,可以用信号量(Semaphore)来控制并发许可:
int maxConcurrency = 3; Semaphore concurrencySemaphore = new Semaphore(maxConcurrency); playersIdsObservable .distinct() .flatMapSingle(playerId -> // 先获取并发许可,没有则等待 Single.fromCallable(() -> { concurrencySemaphore.acquire(); return playerId; }) // 执行实际的请求 .flatMap(id -> dao.players().loadPlayerSRx(id)) // 无论成功失败,最后都释放许可,腾出位置给新任务 .doFinally(concurrencySemaphore::release) , Integer.MAX_VALUE) // 这里把flatMapSingle的默认并发设为最大,交给信号量控制
这个逻辑和你说的完全一致:一开始会同时获取maxConcurrency个许可,启动对应数量的请求,每当一个请求完成(无论成功失败),就释放一个许可,新的playerId就能获取许可并启动请求,保持并发数稳定在设定值。
关于buffer()操作符的区别
你提到buffer()和需求效果类似,但两者有本质区别:
buffer(n)是把上游的事件攒成大小为n的批次,一批处理完再处理下一批,比如buffer(3)就是等3个playerId到齐后,一次性处理这3个,然后等下一批3个,并发数是批次大小,但中间会有空窗期。- 而
flatMapSingle的maxConcurrency是持续保持最多n个并发任务,只要有一个任务完成就立刻启动新的,没有空窗期,吞吐量更高,更适合需要连续处理的场景。
内容的提问来源于stack exchange,提问作者Alex Kokorin
相关产品推荐
相关产品推荐

