RxJava Observalble create subscribe源码分析

先看一小段代码

Observable<String> observable = Observable.create(observer->{
           observer.onNext("处理的数字是"+Math.random()*100);
           observer.onComplete();
});
observable.subscribe(consumer->{
    System.out.println("我处理的元素是"+consumer);
});
observable.subscribe(consumer->{
    System.out.println("我处理的元素是"+consumer);
});

执行结果是

我处理的元素是处理的数字是19.702425673460567
我处理的元素是处理的数字是9.601318081392996

先看Observable.create方法

    @CheckReturnValue
    @NonNull
    @SchedulerSupport(SchedulerSupport.NONE)
    public static <T> Observable<T> create(ObservableOnSubscribe<T> source) {
        ObjectHelper.requireNonNull(source, "source is null");
        return RxJavaPlugins.onAssembly(new ObservableCreate<T>(source));
    }

参数是ObservableOnSubscribe

public interface ObservableOnSubscribe<T> {

    /**
     * Called for each Observer that subscribes.
     * @param emitter the safe emitter instance, never null
     * @throws Exception on error
     */
    void subscribe(@NonNull ObservableEmitter<T> emitter) throws Exception;
}

其实我们可以把我们最开始的例子改写成

ObservableOnSubscribe observableOnSubscribe = new ObservableOnSubscribe(){
        @Override
        public void subscribe(ObservableEmitter emitter) throws Exception {
               emitter.onNext("处理的数字是"+Math.random()*100);
               emitter.onComplete();
        }};
       
Observable<String> observable = Observable.create(observableOnSubscribe);

observable.subscribe(consumer->{
      System.out.println("我处理的元素是"+consumer);
});
observable.subscribe(consumer->{
      System.out.println("我处理的元素是"+consumer);
});

我们把create方法参数还原成1.8之前的写法,我们一眼就看出文章一开始的代码写的observer是影响我们理解代码的

observer->{
            observer.onNext("处理的数字是"+Math.random()*100);
            observer.onComplete();
        }

其实是emitter更为恰当

emitter->{
            emitter.onNext("处理的数字是"+Math.random()*100);
            emitter.onComplete();
        }

这个ObservableEmitter 又是个接口,也就是说下面这几行代码只是定义了一个模版,subscribe的时候,由ObservableEmitter的实现类还具体执行onNext和onComplete。那么实现类在哪里呢?

ObservableOnSubscribe observableOnSubscribe = new ObservableOnSubscribe(){
        @Override
        public void subscribe(ObservableEmitter emitter) throws Exception {
               emitter.onNext("处理的数字是"+Math.random()*100);
               emitter.onComplete();
        }};

我们再看Observable.create方法

    @CheckReturnValue
    @NonNull
    @SchedulerSupport(SchedulerSupport.NONE)
    public static <T> Observable<T> create(ObservableOnSubscribe<T> source) {
        ObjectHelper.requireNonNull(source, "source is null");
        return RxJavaPlugins.onAssembly(new ObservableCreate<T>(source));
    }

也就是说

Observable<String> observable = Observable.create(observableOnSubscribe);

observable等于ObservableCreate的一个实例。这个ObservableCreate留着待用。

我们再看observable.subscribe方法

 @CheckReturnValue
 @SchedulerSupport(SchedulerSupport.NONE)
 public final Disposable subscribe(Consumer<? super T> onNext) {
     return subscribe(onNext, Functions.ON_ERROR_MISSING, Functions.EMPTY_ACTION, Functions.emptyConsumer());
 }

可以看到除了onNext函数是往下传递的,剩下的参数都是默认值。
再放下跟subscribe方法

    @CheckReturnValue
    @SchedulerSupport(SchedulerSupport.NONE)
    public final Disposable subscribe(Consumer<? super T> onNext, Consumer<? super Throwable> onError,
            Action onComplete, Consumer<? super Disposable> onSubscribe) {
        ObjectHelper.requireNonNull(onNext, "onNext is null");
        ObjectHelper.requireNonNull(onError, "onError is null");
        ObjectHelper.requireNonNull(onComplete, "onComplete is null");
        ObjectHelper.requireNonNull(onSubscribe, "onSubscribe is null");
        LambdaObserver<T> ls = new LambdaObserver<T>(onNext, onError, onComplete, onSubscribe);
        subscribe(ls);
        return ls;
    }

注意这个LambdaObserver,传递进来的onNext函数,在这里包装成了一个observer对象。
继续进入subscribe(ls);

    @SchedulerSupport(SchedulerSupport.NONE)
    @Override
    public final void subscribe(Observer<? super T> observer) {
        ObjectHelper.requireNonNull(observer, "observer is null");
        try {
            observer = RxJavaPlugins.onSubscribe(this, observer);
            ObjectHelper.requireNonNull(observer, "The RxJavaPlugins.onSubscribe hook returned a null Observer. Please change the handler provided to RxJavaPlugins.setOnObservableSubscribe for invalid null returns. Further reading: https://github.com/ReactiveX/RxJava/wiki/Plugins");
            subscribeActual(observer);
        } catch (NullPointerException e) { // NOPMD
            throw e;
        } catch (Throwable e) {
            Exceptions.throwIfFatal(e);
            // can't call onError because no way to know if a Disposable has been set or not
            // can't call onSubscribe because the call might have set a Subscription already
            RxJavaPlugins.onError(e);
            NullPointerException npe = new NullPointerException("Actually not, but can't throw other exceptions due to RS");
            npe.initCause(e);
            throw npe;
        }
    }

没什么重要的,继续进入subscribeActual(observer);

protected abstract void subscribeActual(Observer<? super T> observer);

发现是个抽象方法,那么自然应该是刚刚待用的ObservableCreate的
subscribeActual方法

protected void subscribeActual(Observer<? super T> observer) {
        CreateEmitter<T> parent = new CreateEmitter<T>(observer);
        observer.onSubscribe(parent);

        try {
            source.subscribe(parent);
        } catch (Throwable ex) {
            Exceptions.throwIfFatal(ex);
            parent.onError(ex);
        }
    }

第一句把observer对象再包装成一个emitter对象
第二句什么也没执行,因为observer只有onNext是传进来一个lambda,其他三个参数都是默认的。记得是个emptyConsumer。
本文的重中之重就是下面这句

source.subscribe(parent);

source就是我们一开始定义的observableOnSubscribe
subscribe就是observableOnSubscribe的subscribe方法
参数parent就是刚刚的CreateEmitter

ObservableOnSubscribe observableOnSubscribe = new ObservableOnSubscribe(){
        @Override
        public void subscribe(ObservableEmitter emitter) throws Exception {
               emitter.onNext("处理的数字是"+Math.random()*100);
               emitter.onComplete();
        }};

至此所有逻辑拼接成功
先执行subscribe,然后再执行我们自己定义的onNext。
done

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 217,406评论 6 503
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,732评论 3 393
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 163,711评论 0 353
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,380评论 1 293
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,432评论 6 392
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,301评论 1 301
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,145评论 3 418
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 39,008评论 0 276
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,443评论 1 314
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,649评论 3 334
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,795评论 1 347
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,501评论 5 345
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,119评论 3 328
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,731评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,865评论 1 269
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,899评论 2 370
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,724评论 2 354