当Publisher与Subscriber同线程时,RxJava背压的作用是什么?
RxJava单线程场景下的背压疑问解答
测试代码
public class IntervalObserverTest { @Test public void test() throws InterruptedException { CountDownLatch latch = new CountDownLatch(10); IntervalPublisher o = new IntervalPublisher(); o.getPublisher().subscribe(new Subscriber<String>() { volatile Subscription s; @Override public void onSubscribe(Subscription s) { this.s = s; s.request(10); } @Override public void onNext(String s) { System.out.println(Thread.currentThread().getName()); latch.countDown(); } @Override public void onError(Throwable t) { s.cancel(); } @Override public void onComplete() { } }); if(!latch.await(5, TimeUnit.SECONDS)) { throw new RuntimeException(); } } } class IntervalPublisher { public Publisher<String> getPublisher() { return Flowable.range(1, 100) .map(String::valueOf) .delay(50, TimeUnit.MILLISECONDS) .delaySubscription(100, TimeUnit.MILLISECONDS); } }
问题背景
运行上述代码时,onNext()方法打印的线程名称表明,它和Publisher(Flowable)运行在同一个线程(如RxComputationThreadPool-2),即Publisher与Subscriber处于同一线程。由此产生以下疑问:
- 在此单线程场景下,RxJava的背压设置有何作用?
- 单线程下Subscriber不可能过载,这种说法是否正确?
- 背压设置的通用作用是什么?认为仅当Publisher与Subscriber处于不同线程时,背压才用于解决经典有界队列的生产者消费者问题,同线程下无意义,这种观点是否正确?
问题解答
1. 单线程场景下背压的作用
在单线程场景中,背压并非毫无用处,它的核心作用包括:
- 主动控制数据接收量:通过
request(n),Subscriber可以指定每次接收的数据条数,避免一次性接收全部数据。比如示例代码中request(10),Flowable只会先发射10条数据,剩余数据需要再次调用request才会继续发送,让Subscriber拥有数据接收的主动权。 - 协调中间操作符的缓存:RxJava的很多操作符(如
delay)内部会维护队列缓存数据,背压机制可以控制上游发送到这些队列的数据量,防止队列无限膨胀导致内存溢出。 - 统一API逻辑:背压是Flowable的标准特性,无论单线程还是多线程场景,API设计保持一致,开发者无需为不同线程模型修改代码逻辑,降低了学习和维护成本。
2. 单线程下Subscriber不可能过载的说法是否正确?
这种说法错误。
即使在单线程环境中,Subscriber依然可能出现过载:
- 如果
onNext()的处理逻辑耗时较长,而上游Publisher发射数据的速度远快于处理速度,数据会积压在中间操作符的队列中。比如在示例代码的onNext()中加入Thread.sleep(1000)模拟耗时操作,上游delay(50ms)每50ms产生一条数据,单线程下数据会持续堆积在delay的内部队列中,最终引发内存溢出。
3. 背压的通用作用及同线程下观点的判断
背压的通用作用是让下游Subscriber向上游Publisher反馈自身处理能力,调节数据发射速度,避免数据积压和内存溢出,是一套通用的流量控制机制。
认为“同线程下背压无意义”的观点不正确:
- 单线程场景下依然存在数据积压的风险,背压可以通过
request(n)匹配上下游的处理速度,避免中间队列溢出。 - 背压不仅仅解决跨线程的生产者-消费者问题,它适用于所有需要协调数据流速的场景,包括单线程下中间操作符的异步缓存场景。
- 很多看似单线程的场景,内部操作符可能已经引入了异步逻辑(如
delay依赖调度器),背压依然需要发挥作用来控制这些异步环节的数据流动。
内容的提问来源于stack exchange,提问作者ng.newbie
相关产品推荐
相关产品推荐

