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

RxJava中PublishSubject订阅无输出问题排查及多线程实现方案

问题分析与修复方案

首先,你遇到的核心问题是PublishSubject的订阅时机不匹配,以及手动使用Subject带来的事件丢失风险。我们一步步拆解问题并给出修复方案:

为什么PublishSubject没输出,ReplaySubject却可以?

当你调用repository.exec(inputNumber)时,方法会立即创建PublishSubject并启动populateList的异步操作。如果populateList的执行速度很快(比如输入值很小,生成集合的逻辑几乎不耗时),那么在Presenter完成订阅之前,populateList的回调就已经发送了所有onNext事件。而PublishSubject不会缓存已发送的事件,后续的订阅自然收不到内容。

ReplaySubject会缓存所有发送过的事件,所以不管什么时候订阅都能拿到结果,但这不符合你“实时接收计算完成结果”的需求,还可能带来内存占用问题。

修复步骤

1. 替换手动Subject为Observable操作符

手动创建Subject很容易出现订阅时机错误,改用RxJava的操作符构建冷Observable可以从根源解决问题——冷Observable只有在订阅发生时才会执行上游逻辑,保证事件都在订阅之后发射。

修改后的exec方法:

public Observable<BaseUnit> exec(int inputNumber) {
    // 更新现有unit的状态
    if (!unitList.isEmpty()) {
        for (BaseUnit unit : unitList) {
            unit.setInProgress();
        }
        unitList.clear(); // 根据需求决定是否清空,避免重复数据
    }

    return populateList(inputNumber)
            // 复用全局线程池,避免每次新建导致线程泄漏
            .subscribeOn(Schedulers.from(ThreadPool.getGlobalPool()))
            // 把集合中的每个List<Integer>单独发射
            .flatMapIterable(calculatedList -> calculatedList)
            // 为每个List<Integer>绑定所有ListOperationName
            .flatMap(elem -> Observable.fromArray(ListOperationName.values())
                    .map(op -> new AbstractMap.SimpleEntry<>(elem, op)))
            // 执行计算并生成Unit对象
            .map(entry -> {
                List<Integer> elem = entry.getKey();
                ListOperationName operationName = entry.getValue();
                ListUnit unit = new ListUnit(operationName, elem, 0);
                calculate(unit); // 同步耗时操作,已在后台线程执行
                unitList.add(unit);
                return unit;
            });
}

2. 复用线程池,避免资源浪费

你之前每次调用exec都新建FixedThreadPool,会导致线程泄漏和资源浪费。建议维护一个全局线程池:

public class ThreadPool {
    private static ExecutorService globalPool;

    public static ExecutorService getGlobalPool() {
        if (globalPool == null || globalPool.isShutdown()) {
            // 根据CPU核心数设置线程池大小,更合理
            globalPool = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());
        }
        return globalPool;
    }

    // 程序退出时记得关闭线程池
    public static void shutdown() {
        if (globalPool != null && !globalPool.isShutdown()) {
            globalPool.shutdown();
        }
    }
}

3. 简化Presenter的订阅逻辑

Presenter中多余的subscribeOn调用对Subject无效(Subject是热Observable),而且上游已经指定了执行线程,所以只保留observeOn切换到主线程即可:

// Presenter中添加订阅管理,避免内存泄漏
private CompositeDisposable disposables = new CompositeDisposable();

public void calculate(int inputNumber) {
    fragment.showAllProgressBars();
    Disposable disposable = repository.exec(inputNumber)
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(unit -> {
                Log.d("PRESENTER RESULT", unit.toString());
                fragment.setCellText(unit.getViewId(), unit.getTimeString());
            }, error -> Log.d("PRESENTER ERROR", error.toString()));
    disposables.add(disposable);
}

// 页面销毁时取消所有订阅
public void onDestroy() {
    disposables.dispose();
}

额外优化:异步计算的处理

如果calculate(unit)是异步操作,不要直接同步调用,应该把它包装成Observable,确保计算完成后再发射结果:

.flatMap(entry -> {
    List<Integer> elem = entry.getKey();
    ListOperationName operationName = entry.getValue();
    ListUnit unit = new ListUnit(operationName, elem, 0);
    // 假设calculateAsync返回Observable<Unit>,表示计算完成
    return calculateAsync(unit)
            .map(ignored -> unit);
})

内容的提问来源于stack exchange,提问作者white-imp

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 12:22:38