【问题标题】:Is unsubscribe thread safe in RxJava?在 RxJava 中取消订阅线程安全吗?
【发布时间】:2015-09-22 17:36:03
【问题描述】:

假设我有以下 RxJava 代码(它访问数据库,但确切的用例无关紧要):

public Observable<List<DbPlaceDto>> getPlaceByStringId(final List<String> stringIds) {
    return Observable.create(new Observable.OnSubscribe<List<DbPlaceDto>>() {
        @Override
        public void call(Subscriber<? super List<DbPlaceDto>> subscriber) {
            try {
                Cursor c = getPlacseDb(stringIds);

                List<DbPlaceDto> dbPlaceDtoList = new ArrayList<>();
                while (c.moveToNext()) {
                    dbPlaceDtoList.add(getDbPlaceDto(c));
                }
                c.close();

                if (!subscriber.isUnsubscribed()) {
                    subscriber.onNext(dbPlaceDtoList);
                    subscriber.onCompleted();
                }
            } catch (Exception e) {
                if (!subscriber.isUnsubscribed()) {
                    subscriber.onError(e);
                }
            }
        }
    });
}

鉴于此代码,我有以下问题:

  1. 如果有人取消订阅从此方法返回的 observable(在之前的订阅之后),那么该操作是线程安全的吗?那么无论安排如何,我的“isUnsubscribed()”检查在这个意义上是否正确?

  2. 有没有比我在这里使用的更简洁的方法来检查未订阅状态的样板代码?我在框架中找不到任何东西。我以为 SafeSubscriber 解决了订阅者退订时不转发事件的问题,但显然它没有。

【问题讨论】:

    标签: java reactive-programming rx-java


    【解决方案1】:

    该操作是线程安全的吗?

    是的。您收到一个 rx.Subscriber,它(最终)检查一个在订阅者的订阅被取消订阅时设置为 true 的 volatile 布尔值。

    用更少的样板代码更简洁地检查未订阅状态

    SyncOnSubscribeAsyncOnSubscribe(在 1.0.15 版中作为 @Experimental api 可用)是为此用例创建的。它们可以作为调用Observable.create 的安全替代方法。这是同步情况的(人为的)示例。

    public static class FooState {
        public Integer next() {
            return 1;
        }
        public void shutdown() {
    
        }
        public FooState nextState() {
            return new FooState();
        }
    }
    public static void main(String[] args) {
        OnSubscribe<Integer> sos = SyncOnSubscribe.createStateful(FooState::new, 
                (state, o) -> {
                    o.onNext(state.next());
                    return state.nextState();
                }, 
                state -> state.shutdown() );
        Observable<Integer> obs = Observable.create(sos); 
    }
    

    请注意,SyncOnSubscribe next 函数不允许在每次迭代中多次调用observer.onNext,也不能同时调用该观察者。以下是1.x 分支头部的SyncOnSubscribe implementationtests 的几个链接。它的主要用途是简化编写同步迭代或解析数据的可观察对象以及下游的 onNext,但在支持背压并检查是否取消订阅的框架中这样做。本质上,您将创建一个next 函数,每次下游操作员需要一个新的数据元素onNexted 时都会调用该函数。您的 next 函数可以调用 onNext 0 次或 1 次。

    AsyncOnSubscribe 旨在很好地处理异步操作的可观察源(例如开箱即用调用)的背压。下一个函数的参数包括请求计数,并且您提供的 observable 应该提供一个 observable 来满足请求数量的数据。这种行为的一个例子是来自外部数据源的分页查询。

    以前,将OnSubscribe 转换为Iterable 并使用Observable.from(Iterable) 是一种安全的做法。此实现获取一个迭代器并为您检查subscriber.isUnsubscribed()

    【讨论】:

    • 谢谢,您实际上回答了我的另一个问题,即创建具有适当背压支持的“自定义”可观察对象!我检查了 SyncSubscriber,它看起来非常好。在许多情况下,我会看到将操作转换为 Iterable 在语义上有点尴尬,但很高兴知道我们可以通过这种方式轻松获得背压支持。我看到虽然这个类仍然被标记为@Experimental,但你认为它什么时候可以粗略地被认为是生产就绪的?
    • 很高兴听到!下一步是将其发布到版本中(应该很快会在 1.0.15 中发布)。之后,一般会升级为@Beta状态或当我们获得信心后直接进入公共状态。
    猜你喜欢
    • 2015-07-23
    • 2019-03-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多