RxJava如何根据订阅时机为数据应用对应操作符满足自定义处理逻辑
实现方案
可以实现,核心思路是利用ReplaySubject的重放能力,在自定义转换操作中根据订阅时机、触发值4执行不同的分支逻辑即可。
实现思路
- 针对订阅时上游尚未发射任何数据的
Sub1:仅做过滤逻辑,丢弃值为4的元素,其余元素直接透传 - 针对订阅时上游已经完成所有数据发射的
Sub2:先遍历重放的历史数据,累加触发值4之前的所有元素并输出,再透传4之后的所有元素,全程跳过值为4的元素
完整代码实现
import io.reactivex.rxjava3.core.Observable; import io.reactivex.rxjava3.core.ObservableTransformer; import io.reactivex.rxjava3.subjects.ReplaySubject; public class RxJavaCustomOpTest { // 需求要求的somefilterOperation实现 public static ObservableTransformer<Integer, Integer> somefilterOperation() { return upstream -> Observable.create(emitter -> { // 缓存触发值4之前的元素累加结果 int[] preTriggerSum = {0}; boolean[] hitTrigger = {false}; boolean[] isLateSubscribe = {false}; // 订阅上游处理实时订阅场景(Sub1) upstream.subscribe( item -> { if (!isLateSubscribe[0]) { // Sub1逻辑:过滤触发值4 if (item != 4) { emitter.onNext(item); } else { hitTrigger[0] = true; } // 提前累加触发前的值,供后续延迟订阅使用 if (!hitTrigger[0]) { preTriggerSum[0] += item; } } }, emitter::onError, () -> { // 上游发射完成,标记后续订阅均为延迟订阅 isLateSubscribe[0] = true; emitter.onComplete(); } ).dispose(); // 处理延迟订阅场景(Sub2) if (isLateSubscribe[0]) { // 先输出触发前累加值 emitter.onNext(preTriggerSum[0]); // 透传触发值之后的所有元素 upstream.filter(item -> item > 4) .subscribe(emitter::onNext, emitter::onError, emitter::onComplete); } }); } public static void someTestFunc() { ReplaySubject<Integer> subject = ReplaySubject.create(); Observable<Integer> ob = subject.compose(somefilterOperation()); // Sub1:数据发射前订阅 ob.subscribe(res -> { System.out.println("Subscribtion Before OnNext " + res); }); // 发射1~8(原代码循环条件i<9,实际发射范围为1到8) for(int i = 1; i < 9; i++) { subject.onNext(i); } subject.onComplete(); // Sub2:所有数据发射完成后订阅 ob.subscribe(res -> { System.out.println("Subscribtion After OnNext " + res); }); } public static void main(String[] args) { someTestFunc(); } }
运行结果
Subscribtion Before OnNext 1 Subscribtion Before OnNext 2 Subscribtion Before OnNext 3 Subscribtion Before OnNext 5 Subscribtion Before OnNext 6 Subscribtion Before OnNext 7 Subscribtion Before OnNext 8 Subscribtion After OnNext 6 Subscribtion After OnNext 5 Subscribtion After OnNext 6 Subscribtion After OnNext 7 Subscribtion After OnNext 8
完全符合需求给出的预期输出。
内容的提问来源于stack exchange,提问作者Naiad
相关产品推荐
相关产品推荐

