如何使用RxJava 2实现支持并发任务的异步事件总线?
RxJava 2 Event Bus: Getting Concurrent Execution for Long-Running Subscriber Tasks
我懂你现在的头疼点——用PublishSubject搭了事件总线,发射器能正常发事件,但订阅者处理长耗时任务时,想让任务一收到就立刻启动、不用等前一个做完,结果用了observeOn还是没达到并发的效果。这其实是PublishSubject的事件分发逻辑加上observeOn的调度特性导致的,咱们一步步来解决。
先搞清楚问题根源
首先得明白两个关键点:
- PublishSubject本身是串行推送事件给每个订阅者的——哪怕你用了
observeOn(Schedulers.io()),单个订阅者也是按事件顺序依次处理,前一个任务没跑完,后一个就不会在这个订阅者的线程里启动。 - 另外,
Schedulers.io()的线程池是动态扩容的,但如果你的任务都是同一个订阅者在处理,线程池再大也没用,因为这个订阅者的逻辑是串行接收事件的。
解决方案1:单个订阅者实现任务并发
如果你希望同一个订阅者收到每个事件后,都立刻启动独立的任务(互不等待),可以用flatMap把每个事件转换成一个独立的异步Observable,让每个任务都跑在自己的线程里:
原来的问题代码(串行处理)
eventBus.getSubject() .observeOn(Schedulers.io()) .subscribe(event -> { // 长耗时任务,比如模拟睡眠5秒 handleLongRunningTask(event); });
修改后的并发代码
eventBus.getSubject() .flatMap(event -> Observable.fromCallable(() -> { // 这里放你的长耗时任务 handleLongRunningTask(event); return event; }).subscribeOn(Schedulers.io())) .subscribe();
为什么这样能行?
flatMap会把每个事件拆成一个新的Observable实例,每个Observable都通过subscribeOn(Schedulers.io())获取独立的线程- 这样每个事件对应的任务都会立即启动,完全不需要等前一个任务完成,真正实现并发
解决方案2:多订阅者场景下的并发
如果是多个不同的订阅者,想让它们各自的任务同时执行,那要确保每个订阅者都通过observeOn指定调度器,甚至可以自定义线程池来控制并发数:
// 自定义固定大小的线程池,比如10个线程,适合批量并发任务 ExecutorService customExecutor = Executors.newFixedThreadPool(10); Scheduler customScheduler = Schedulers.from(customExecutor); // 订阅者1:处理任务A eventBus.getSubject() .observeOn(customScheduler) .subscribe(event -> handleTaskA(event)); // 订阅者2:处理任务B eventBus.getSubject() .observeOn(customScheduler) .subscribe(event -> handleTaskB(event));
这样两个订阅者的任务会在自定义线程池里同时执行,互相不阻塞。
额外提醒
- 记得在应用退出时关闭自定义线程池,避免内存泄漏:
customExecutor.shutdown() - 如果任务需要处理异常或者返回结果,记得在
flatMap里添加onErrorResumeNext之类的异常处理逻辑 - 要是你需要保持事件的处理顺序,但又想让任务并发启动,可以用
concatMapEager——它会同时启动所有任务,但按事件的原始顺序输出结果
内容的提问来源于stack exchange,提问作者olivierlemasle
相关产品推荐
相关产品推荐

