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

RxJava如何正确使用Observable与Observer?代码实现及内存泄漏疑问咨询

问题答复

基础疑问解答

  • onComplete触发后RxJava会自动终止当前订阅,对应的Disposable会被标记为已处置,不会残留订阅绑定关系,你这个单次发射的场景下该环节不会直接导致内存泄漏。
  • 每次点击按钮新建Observable的做法完全不必要,属于对RxJava使用场景的误用。你当前的写法本质上是每次点击执行一次单次事件回调,和直接调用PrimaryInfoModel的属性赋值方法没有区别,完全没有用到RxJava流式订阅的优势,还额外增加了对象创建开销。

现有代码隐藏问题

  1. 你在静态方法getPageObserver()返回的Observer实例中修改了非静态成员变量currentPage,如果PrimaryInfoModel不是单例实现,会出现变量更新不生效、引用错位的问题。
  2. 你指定了observeOn(Schedulers.io()),如果后续需要在事件回调中更新JavaFX UI,会直接抛出非UI线程操作异常,你当前仅更新成员变量所以暂时未触发问题。
  3. 快速连续点击按钮会生成大量临时Observable和订阅实例,虽然都是短生命周期,依然存在不必要的性能开销。

正确实现方案

你这个场景应该用Subject作为可观察数据源,仅初始化一次,每次点击直接发射事件即可,优化后代码如下:

1. PrimaryInfoModel调整
public class PrimaryInfoModel {
    private int currentPage;
    private int currentCategory;
    // 持有订阅实例方便后续销毁
    private Disposable pageDisposable;

    public Consumer<Integer> getPageUpdateConsumer() {
        return integer -> {
            System.out.println(integer);
            currentPage = integer;
        };
    }

    public void setPageDisposable(Disposable d) {
        this.pageDisposable = d;
    }

    // 页面销毁时调用,彻底清理订阅
    public void destroy() {
        if (pageDisposable != null && !pageDisposable.isDisposed()) {
            pageDisposable.dispose();
        }
    }
}
2. NavigationController调整
public class NavigationController implements Initializable {
    private int currentPage;
    // 用PublishSubject作为数据源,只初始化一次
    private PublishSubject<Integer> pageSubject = PublishSubject.create();
    private Disposable subscribeDisposable;

    @FXML
    private Button page1, page2, page3, page4, page5, prevPage, nextPage;

    @Override
    public void initialize(URL url, ResourceBundle resourceBundle) {
        // 初始化时就绑定订阅,仅绑定一次,这里替换为你实际获取PrimaryInfoModel实例的逻辑
        PrimaryInfoModel model = new PrimaryInfoModel(); 
        subscribeDisposable = pageSubject
                // 如果后续需要更新UI就替换为JavaFxScheduler.platform()
                .observeOn(Schedulers.io())
                .subscribeOn(Schedulers.io())
                .subscribe(model.getPageUpdateConsumer(),
                        e -> System.out.println(e.getMessage()),
                        () -> System.out.println("Complete"));
        model.setPageDisposable(subscribeDisposable);
    }

    @FXML
    private void setCurrentButton(MouseEvent event) {
        int oldPage = currentPage;
        cleanEffect();
        Object source = event.getSource();
        if (source == page1) {
            currentPage = 0;
        } else if (source == page2) {
            currentPage = 1;
        } else if (source == page3) {
            currentPage = 2;
        } else if (source == page4) {
            currentPage = 3;
        } else if (source == page5) {
            currentPage = 4;
        }
        if (oldPage != currentPage) {
            // 直接发射事件即可,不需要新建Observable
            pageSubject.onNext(currentPage);
        }
    }

    // 页面销毁时调用,清理Subject
    public void destroy() {
        pageSubject.onComplete();
        if (subscribeDisposable != null && !subscribeDisposable.isDisposed()) {
            subscribeDisposable.dispose();
        }
    }
}

内容的提问来源于stack exchange,提问作者MichaelHo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 21:39:04