RxJava2背压与并发处理:Flowable回调服务改造及疑问
咱们先逐个核对你的需求,看看现有实现是否达标,再逐一解答你的疑问,最后聊聊复用的方案。
一、现有实现是否符合需求?
你的MyBroadcaster实现基本满足所有要求,咱们拆解来看:
- 回调线程不阻塞:回调调用
packets::onNext时,PublishProcessor本身不会阻塞回调线程(它是无缓冲的,背压策略由onBackpressureLatest处理),再通过observeOn(scheduler)把事件处理移交到观察者指定的线程,回调线程只负责触发事件,不会等待下游处理,完全符合要求。 - 下游过慢不抛异常:
onBackpressureLatest会在下游跟不上生产者速度时,丢弃旧事件只保留最新的,不会抛出MissingBackpressureException,完美解决下游慢的问题。 - 支持多消费者独立订阅:
PublishProcessor本身就是多播组件,每个订阅者都会独立接收后续事件,加上onBackpressureLatest后,每个订阅者的背压处理也是独立的(各自保留自己的最新事件),多订阅场景完全没问题。 - 消费者自行缓存消息:你的生产者类
MyBroadcaster用的是PublishProcessor(无内置缓存),如果消费者需要缓存所有消息,只需要在订阅时自行添加cache()或replay()操作符,比如:
缓存逻辑在消费者的Flowable链中,不会污染生产者类,符合要求。broadcaster.observeAllPacketsOn(Schedulers.io()).cache()
关于你对onBackpressureLatest文档中observeOn说明的疑惑:
Note that due to the nature of how backpressure requests are propagated through subscribeOn/observeOn, requesting more than 1 from downstream doesn't guarantee a continuous delivery of onNext events
这段话的意思是:observeOn内部有一个默认大小为16的队列,当下游请求n>1的事件时,observeOn会从上游拉取n个事件放到队列里,但因为背压请求是异步传递的(observeOn在指定线程处理事件),上游不会一次性把所有请求的事件都发过来——当下游处理完一个事件,才会再向上游请求一个。所以即使下游请求了多个,也不会保证连续收到onNext,这是observeOn的正常行为,结合onBackpressureLatest的话,上游只会保留最新事件,这个特性不会影响你的需求。
二、你的疑问逐一解答
1. 调用onBackpressureLatest后是否会取消消息多播特性?
完全不会。PublishProcessor的多播能力是它的核心属性,onBackpressureLatest()只是在它的基础上添加了背压处理逻辑,返回的Flowable依然是多播的。每个订阅者都会独立应用这个背压策略,各自维护自己的最新事件,互相之间没有影响。
2. 如何测试上述需求是否满足?
可以分场景写测试用例,这里给你几个核心场景的测试思路:
- 测试回调线程不阻塞:
模拟回调线程发送事件,同时让下游做耗时操作,验证回调线程的执行耗时远小于下游处理时间,说明回调线程没有被阻塞。@Test public void testCallbackThreadNonBlocking() throws InterruptedException { CountDownLatch latch = new CountDownLatch(1); MyBroadcaster broadcaster = new MyBroadcaster(); // 下游做1秒的耗时操作 broadcaster.observeAllPacketsOn(Schedulers.io()) .subscribe(packet -> { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } latch.countDown(); }); // 在主线程模拟回调触发事件 long start = System.currentTimeMillis(); // 可以给MyBroadcaster加个测试用的方法触发onNext,或者临时把packets改成包可见 broadcaster.packets.onNext(new Packet()); long end = System.currentTimeMillis(); // 回调线程(主线程)耗时应该远小于1秒 assertTrue(end - start < 100); latch.await(); } - 测试下游过慢不抛异常:
上游持续发送大量事件,下游每个事件都sleep很长时间,观察是否没有抛出异常,且下游只收到最新的事件。 - 测试多消费者独立订阅:
创建两个观察者,一个快速处理,一个慢速处理,发送多个事件,验证快的观察者能收到更多(或所有)事件,慢的只收到最新的,两者互不干扰。 - 测试消费者自行缓存:
先发送几个事件,再让一个加了cache()的观察者订阅,验证它能收到之前发送的所有事件,而另一个没加缓存的观察者只能收到订阅后的事件。
3. 多个发布者如何实现模式复用?是否需要自定义Flowable?
不需要自定义Flowable,用抽象类就能完美复用这个模式,把通用的Flowable封装、背压、线程切换逻辑抽出来,子类只需要实现具体的回调注册逻辑即可:
// 通用的抽象广播器 abstract class AbstractBroadcaster<T> { private final PublishProcessor<T> processor = PublishProcessor.create(); private final Flowable<T> backpressuredFlowable = processor.onBackpressureLatest(); protected AbstractBroadcaster() { // 注册回调,把事件转发给processor registerCallback(this::onEventReceived); } // 子类实现具体的回调注册逻辑 protected abstract void registerCallback(Consumer<T> callback); // 接收回调事件,转发给processor private void onEventReceived(T event) { processor.onNext(event); } // 对外提供的订阅方法,指定调度器 public Flowable<T> observeOn(Scheduler scheduler) { return backpressuredFlowable.observeOn(scheduler); } } // 你的具体广播器,只需要实现回调注册 class MyBroadcaster extends AbstractBroadcaster<Packet> { @Override protected void registerCallback(Consumer<Packet> callback) { // 这里写你实际的回调注册逻辑,比如: // networkService.registerPacketCallback(callback); } }
如果之后需要其他类型的广播器,比如UserActivityBroadcaster,只需要继承AbstractBroadcaster<UserActivity>并实现registerCallback即可,完全复用了所有通用逻辑。
内容的提问来源于stack exchange,提问作者Aarjav

