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

