RxJava:如何在数组迭代的每次发射间设置对应时长的延迟?
问题分析
你的代码问题出在这两个关键地方:
Observable.interval(intervals[index.getAndIncrement()], TimeUnit.SECONDS)里的周期参数只会在Observable创建时计算一次,也就是只会用数组的第一个元素1L作为固定周期,后续的index.getAndIncrement()根本不会被执行——因为这个intervalObservable一旦初始化,周期就固定死了。zipWith会把fromArray同步发射的所有4个数组元素,和interval每隔1秒发射的元素配对,最终效果是每隔1秒打印一次时间,完全没有按照你预期的1→2→3→4秒依次等待。
正确实现方案
要实现“按数组中每个数字依次等待对应时长后执行操作”的需求,我们可以用concatMap操作符——它会依次处理每个上游发射的元素,前一个Observable完成后才会处理下一个,刚好符合“依次等待”的要求。
基础实现代码
Long[] intervals = {1L, 2L, 3L, 4L}; Observable.fromArray(intervals) .concatMap(waitTime -> // 等待指定时长后,发射一个信号触发后续操作 Observable.timer(waitTime, TimeUnit.SECONDS) ) .subscribe(__ -> { System.out.println(LocalDateTime.now()); });
代码解释
Observable.fromArray(intervals)会按顺序发射数组中的每个时长值。concatMap针对每个时长值,创建一个Observable.timer——这个Observable会在指定时长后发射一个元素,然后立即完成。- 得益于
concatMap的特性,只有当前一个timer执行完成(即等待完对应时长),才会处理下一个时长值,完美实现了“等待1秒打印→等待2秒打印→等待3秒打印→等待4秒打印”的效果。
额外优化(保留时长信息)
如果你需要在订阅时同时获取当前等待的时长,可以稍作修改:
Observable.fromArray(intervals) .concatMap(waitTime -> Observable.timer(waitTime, TimeUnit.SECONDS) .map(__ -> waitTime) // 将timer的发射项映射为当前等待的时长 ) .subscribe(waitTime -> { System.out.printf("等待了%d秒,当前时间:%s%n", waitTime, LocalDateTime.now()); });
内容的提问来源于stack exchange,提问作者Shvalb
相关产品推荐
相关产品推荐

