RxJS:将Observable拆分为两个并等待第一个完成后执行第二个
我来帮你梳理一下这个问题的解决方案~其实核心难点在于:在源Observable发射元素的过程中,我们没法提前知道哪个是最后一个元素,必须等源完全完成后才能确定拆分点。下面针对不同场景给你两种可行的实现方案:
方案一:保持原发射时序(适合冷Observable)
如果你的源是冷Observable(比如可以重复订阅、没有副作用的序列),同时希望前n-1个元素保持原有的发射间隔,推荐用这种方式:
// 先定义你的源Observable(示例模拟了你的时序) Observable<String> originalSource = Observable.create(emitter -> { emitter.onNext("a"); Thread.sleep(1000); emitter.onNext("b"); Thread.sleep(500); emitter.onNext("c"); Thread.sleep(1500); emitter.onNext("d"); Thread.sleep(1000); emitter.onNext("e"); Thread.sleep(2000); emitter.onNext("f"); emitter.onComplete(); }); // 缓存源,避免重复执行(如果源有副作用比如网络请求,这一步很关键) Observable<String> cachedSource = originalSource.cache(); // 第一步:先获取所有元素的列表,确定拆分位置 Single<List<String>> allElements = cachedSource.toList(); // 第一个Observable:发射除最后一个之外的所有元素,保持原时序 Observable<String> firstPart = allElements.flatMapObservable(list -> { if (list.size() <= 1) { return Observable.empty(); } // take(n-1)会按原间隔发射前n-1个元素 return cachedSource.take(list.size() - 1); }); // 第二个Observable:仅发射最后一个元素 Observable<String> secondPart = allElements.flatMapObservable(list -> { if (list.isEmpty()) { return Observable.empty(); } // skip(n-1)会跳过前n-1个元素,只发射最后一个 return cachedSource.skip(list.size() - 1); }); // 用concat严格保证第一个Observable完成后,再执行第二个 Observable.concat(firstPart, secondPart) .subscribe( item -> System.out.println("收到元素:" + item), error -> System.err.println("出错:" + error), () -> System.out.println("所有序列执行完成") );
方案说明
cache()让源只执行一次,所有订阅者共享同一个发射序列,避免重复触发副作用;toList()会等待源完全完成,收集所有元素到列表,这样我们就能明确知道元素总数和最后一个元素的位置;take和skip基于索引拆分,不会因为元素重复导致误判;concat操作符的特性就是必须等前一个Observable完全完成(包括内部的所有逻辑和onComplete信号),才会订阅并执行下一个Observable,完美匹配你的需求。
方案二:简化实现(无需保持原发射时序)
如果不要求前n-1个元素保持原有的发射间隔,能接受一次性发射,可以用更简洁的写法:
Observable<String> originalSource = ...; Single<List<String>> allElements = originalSource.toList(); Observable<String> firstPart = allElements.flatMapObservable(list -> { if (list.size() <= 1) { return Observable.empty(); } // 直接从列表中取前n-1个元素一次性发射 return Observable.fromIterable(list.subList(0, list.size() - 1)); }); Observable<String> secondPart = allElements.flatMapObservable(list -> { if (list.isEmpty()) { return Observable.empty(); } // 直接发射列表最后一个元素 return Observable.just(list.get(list.size() - 1)); }); Observable.concat(firstPart, secondPart).subscribe(...);
注意事项
- 如果你的源是热Observable(比如实时的UI事件流),
cache()可能不适用,建议用方案二先收集所有元素再拆分; - 不要用
takeWhile(t -> !t.equals(lastElement))的方式拆分,因为如果序列中有重复元素,会导致提前停止发射。
内容的提问来源于stack exchange,提问作者feerlay
相关产品推荐
相关产品推荐

