RxJava中Flowable.interval结合flatMap(Single)的背压问题及需求
解决RxJava Flowable.interval定时API调用的背压与并发问题
看起来你遇到了RxJava中定时API调用的经典痛点:用Flowable.interval按固定间隔触发API调用,但如果API耗时超过间隔时间,就会导致多个请求并发执行,进而引发背压问题。你想要实现的是仅当API调用未在进行时才发起新请求,下面给你两种适配不同场景的解决方案:
场景1:跳过忙碌时的间隔事件(严格按间隔检查,空闲才执行)
如果你的需求是:到了间隔时间点,若上一次API已经完成就发起新请求;若API还在执行,则跳过这次检查,等下一个间隔点再判断。可以通过onBackpressureLatest处理背压,结合flatMapSingle限制并发数来实现:
Flowable.interval(1, 1, TimeUnit.SECONDS) .onBackpressureLatest() // 下游忙碌时,丢弃旧的间隔事件,只保留最新的 .flatMapSingle(tick -> { System.out.println("发起API调用,当前刻度:" + tick); // 模拟耗时API调用(比如这里延迟2秒,超过1秒的间隔) return Single.just(1L) .doAfterSuccess(result -> System.out.println("API调用完成,结果:" + result)) .delay(2, TimeUnit.SECONDS); }, false, 1) // 设置最大并发数为1,确保同一时间只有一个API在执行 .subscribe( unused -> {}, error -> System.err.println("调用出错:" + error.getMessage()) );
代码解释:
onBackpressureLatest():解决背压的核心——当下游还在处理上一个API请求时,上游的间隔事件会被丢弃,只保留最新的那个,避免事件堆积导致内存问题。flatMapSingle(..., 1):通过maxConcurrency=1强制下游串行执行API调用,完美符合“仅当API空闲时才发起调用”的要求。- 实际效果:第一个刻度0发起请求,耗时2秒;期间刻度1、2的事件会被丢弃;当请求完成后,刻度3的事件会触发下一次API调用,以此类推。
场景2:API完成后再等待固定间隔(无并发,间隔从完成后计算)
如果你的需求是每次API调用完成后,再等待固定间隔发起下一次请求(而不是严格按固定时间点触发),可以用递归+concatWith的方式实现,完全避免并发和背压问题:
// 封装API调用逻辑 private Flowable<Long> executeApiCall() { return Single.just(1L) .doAfterSuccess(result -> System.out.println("API调用完成,结果:" + result)) .delay(2, TimeUnit.SECONDS) // 模拟耗时API .toFlowable() // 调用完成后,等待1秒再发起下一次调用 .concatWith(Flowable.timer(1, TimeUnit.SECONDS) .flatMap(tick -> executeApiCall())); } // 启动调用 executeApiCall().subscribe( unused -> {}, error -> System.err.println("调用出错:" + error.getMessage()) );
代码解释:
- 这种方式是串行递归触发:每一次API调用完成后,才会启动1秒的定时器,定时器结束后再发起下一次调用。
- 完全没有并发问题,也不存在背压,因为整个流程是严格串行的,一个请求完成后才会触发下一个环节。
关键知识点总结
- 背压问题的根源:
Flowable.interval是按固定速率发射事件,不管下游处理速度;而默认的flatMap允许高并发,导致事件堆积。 onBackpressureLatestvsonBackpressureDrop:前者保留最新事件,后者直接丢弃所有下游忙碌时的事件,根据你的业务场景选择。- 并发控制:
flatMapSingle的maxConcurrency参数是限制RxJava下游并发数的关键,设置为1就能实现串行执行。
内容的提问来源于stack exchange,提问作者Richard
相关产品推荐
相关产品推荐

