You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.28 10:15:05