RxJava 2 Maybe无输出问题:多生产者单消费者批量处理排查
问题分析与解决方案
你的代码存在几个关键问题,导致Maybe始终没有输出,咱们一步步拆解并解决:
1. reduce操作符的使用场景错误
你在bus.window(...).reduce(...)里用的reduce是作用在窗口Observable的序列上的——它会收集所有窗口Observable,直到上游的window序列完全结束(也就是bus调用onComplete)才会发射最终合并后的结果。但你的bus是PublishSubject,从来没有触发onComplete,所以reduce返回的Maybe永远不会产生任何输出。
你真正需要的是每个窗口内部的元素批量合并,而非把多个窗口合并成一个Observable。
2. 主线程被无限流阻塞
main方法里的Stream.iterate(0, i -> i+1).forEach(...)是同步执行的无限流,它会一直占用主线程往bus里发数据,导致后面的Thread.sleep(100000L)根本执行不到,窗口的时间触发条件自然也无法生效。
3. 窗口合并逻辑不符合预期
你在reduce里用zipWith把两个窗口的元素按位置拼接,这会把第一个窗口的第n个元素和第二个窗口的第n个元素合并,而不是把单个窗口里的所有元素拼接成一个完整字符串,这和你的预期输出完全不符。
修正后的代码
下面是调整后的实现,完美解决了上述问题:
import io.reactivex.Observable; import io.reactivex.schedulers.Schedulers; import io.reactivex.subjects.PublishSubject; import java.util.concurrent.TimeUnit; import java.util.stream.IntStream; public class Rxtest { private PublishSubject<String> bus = PublishSubject.create(); public Rxtest() { bus // 满足10个元素或10秒任一条件就触发窗口 .window(10, TimeUnit.SECONDS, 10) // 逐个处理每个窗口:把窗口内的所有元素拼接成一个字符串 .concatMap(windowObs -> windowObs .reduce("", (acc, item) -> acc + item) // 窗口为空时返回空字符串(可选,根据业务需求调整) .defaultIfEmpty("") ) // 订阅处理结果,这里可以替换成你的耗时业务操作 .subscribe( batchStr -> System.out.println(batchStr), Throwable::printStackTrace, () -> System.out.println("Done here") ); } public static void main(String[] args) throws InterruptedException { Rxtest test = new Rxtest(); // 用后台线程生产数据,避免阻塞主线程 Schedulers.io().scheduleDirect(() -> IntStream.range(0, 100).forEach(i -> { test.bus.onNext(String.valueOf(i + 1)); // 从1开始,匹配你的预期输出格式 try { // 模拟生产数据的间隔(可选,根据实际场景调整) TimeUnit.MILLISECONDS.sleep(50); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }) ); // 等待足够时间让窗口触发,最后调用onComplete结束整个流 Thread.sleep(15000L); test.bus.onComplete(); Thread.sleep(1000L); } }
代码说明
- 窗口处理逻辑:用
concatMap替换外层的reduce,对每个窗口Observable单独处理——在窗口内部用reduce拼接所有元素,得到该窗口的批量字符串。 - 后台线程生产数据:用
Schedulers.io().scheduleDirect把生产数据的逻辑放到后台线程,避免阻塞主线程,保证窗口的时间触发条件能正常生效。 - 正确触发流结束:在主线程sleep足够时间后,调用
bus.onComplete(),这样整个流会正常结束,触发onComplete回调输出"Done here"。 - 匹配预期输出:把生产的数字从1开始,每个窗口会输出10个连续数字拼接的字符串,比如
12345678910、11121314151617181920等。
内容的提问来源于stack exchange,提问作者nvalada
相关产品推荐
相关产品推荐

