RxJava2源码分析(二) ---- subscribeOn

subscribeOn

    Observable.create((ObservableOnSubscribe<Integer>) e -> {
        System.out.println("observable : " + Thread.currentThread());
        e.onNext(1);
    })
            .subscribeOn(Schedulers.single())
            .subscribe(integer -> {
                System.out.println(integer);
                System.out.println("observer:  " + Thread.currentThread());
            });

Rxjava默认是在当前线程生发送事件, subscribeOn可以切换Observable发送事件所在的线程;
如果没有使用ObserveOn指定消费事件的线程, Observer将在Observable发送事件的的线程, 消费事件;

源码分析目的:

  1. Schduler 作用
  2. subscribeOn 做了什么

1. Schduler

Schduler不好直接用代码解释, 先说结论, 后面再去具体代码分析;

  1. 切换线程, 需要提供对应的Schduler;
  2. Schduler可以通过createWorker方法, 创建一个Worker类的实例;
  3. Worker有一个schedule方法, 提交runnable去运行; 切换线程, 就是把各个onNext的调用方法,封装成一个runnable 提交到指定线程去运行;
  4. 通过Worker.schedule提交runnable后, 会返回一个disposable对象, 用于取消或控制Observale的发射任务;
  5. Schduler本身是个管理类, 一般内部会创建具体的线程池, 同时通过统一的start shutdown等方法管理着线程池
  6. Schduler同时也管理着由createWorker创建的Worker; Worker一般都是持有Schduler中的线程池, 提交的runnable也是提交到该线程池

2. subscribeOn

2.1 Observable.subscribeOn
    public final Observable<T> subscribeOn(Scheduler scheduler) {
        ObjectHelper.requireNonNull(scheduler, "scheduler is null");
        return RxJavaPlugins.onAssembly(new ObservableSubscribeOn<T>(this, scheduler));
    }

1. 参数检测
2. 创建`ObservableOnSubscribe`对象, 并将当前Observable和Schduler传入;
3. RxJavaPlugins的hook; 这个前面说过, 用于hook, 默认传入什么 就返回什么;
2.2 ObservableSubscribeOn
  1. ObservableSubscribeOn是Observable的子类, 内部包含一个Observableschduler, 用于对原Obverable扩展, 是一个装饰模式;
  2. 上面说了, ObservableSubscribeOn是一个装饰模式, 继承于HasUpstreamObservableSource, 有一个source方法去获取被装饰的Observable对象;
  3. 上一篇说过, Observable.create方法创建的Observable, 实际是一个ObservableCreater对象, 现在ObservableScbscribeOn中包含的Observable即ObservableCreater;
  4. Observable的subscribe方法, 实际调用的是具体子类的subscribeActual方法;
2.3 ObservableSubscribeOn.subscribeActual

直接看 ObservableSubscribeOn.subscribeActual 的代码;

    @Override
    public void subscribeActual(final Observer<? super T> s) {
        final SubscribeOnObserver<T> parent = new SubscribeOnObserver<T>(s);

        s.onSubscribe(parent);

        parent.setDisposable(scheduler.scheduleDirect(new SubscribeTask(parent)));
    }
  1. 创建 SubscribeOnObserver, 并传入原始的Observer
  2. 调用Observer的onSubscribe方法
  3. 构建SubscribeTask, 并提交给Schduler去执行
2.4 SubscribeOnObserver

SubscribeOnObserverObservableSubscribeOn的静态内部类, 同时也是继承于Observer, 内部也包含一个原始的Observer, 也是一个装饰模式;

SubscribeOnObserver对被装饰类没有额外增加功能, 仅仅是一个封装, 在onNext, onError等方法中, 直接是调用的actual.onNext, actual.onError;

2.5 SubscribeTask

SubscribeTask 是一个runnable对象, 是ObservableSubscribeOn的内部类; 前面Schduler中说过, 切换线程, 就是将消息发送,包装成一个runnable, 提交给Worker去执行;
这个SubscribeTask将原先的发送事件代码 封装成的runnable, 然后送去对应的线程池执行;

直接看run方法

    public void run() {
        source.subscribe(parent);
    }

ObservableSubscribeOn是Observable的子类, 同时是装饰模式, 内部持有一个Observable, source是被包装的Observable, 在此处的代码中, source即是ObservableCreater, parent是SubscribeOnObserver, source.subscribe即和第一篇中的逻辑一样了;

2.6 scheduler.scheduleDirect
@NonNull
public Disposable scheduleDirect(@NonNull Runnable run) {
    return scheduleDirect(run, 0L, TimeUnit.NANOSECONDS);
}

@NonNull
public Disposable scheduleDirect(@NonNull Runnable run, long delay, @NonNull TimeUnit unit) {
    final Worker w = createWorker();

    final Runnable decoratedRun = RxJavaPlugins.onSchedule(run);

    DisposeTask task = new DisposeTask(decoratedRun, w);

    w.schedule(task, delay, unit);

    return task;
}
  1. 通过createWorker创建相应的Worker;
  2. hook处理相应的runnable, 默认没处理;
  3. 创建DisposeTask, 将需要运行的runnable对象, 封装成disposable对象, 用于执行取消操作;
  4. 将封装后的runnable提交给worker去运行;

此处的scheduler由Schedulers.single()生成, 实际是一个SingleScheduler;

2.6.1 Worker.schedule()

直接看 SingleScheduler的代码

####### Schedulers.createWorker 创建Worker; 获取公共的线程池, 创建Worker

    public Worker createWorker() {
        return new ScheduledWorker(executor.get());
    }

####### Worker.scheduler

    public Disposable schedule(@NonNull Runnable run, long delay, @NonNull TimeUnit unit) {
        if (disposed) {
            return EmptyDisposable.INSTANCE;
        }

        Runnable decoratedRun = RxJavaPlugins.onSchedule(run);

        ScheduledRunnable sr = new ScheduledRunnable(decoratedRun, tasks);
        tasks.add(sr);

        try {
            Future<?> f;
            if (delay <= 0L) {
                f = executor.submit((Callable<Object>)sr);
            } else {
                f = executor.schedule((Callable<Object>)sr, delay, unit);
            }

            sr.setFuture(f);
        } catch (RejectedExecutionException ex) {
            dispose();
            RxJavaPlugins.onError(ex);
            return EmptyDisposable.INSTANCE;
        }

        return sr;
    }
  1. 封装传入的runnable对象, 将其封装成ScheduledRunnable对象
  2. 提交给线程池运行, ScheduledRunnable本身是一个Callable对象, 可以用于取消执行

上述提交给线程池运行的流程, 最终封装的运行的run方法, 其实还是最先封装的SubscribeTask中的source.subscribe(parent);这一句代码;
SubscribeTask本身对应的runnable被一次次传递封装, 最后给线程池运行;
source.subscribe(parent);中, 上面说到是一个装饰模式, 运行的还是Observable.subscribeActual方法, 最后的运行逻辑和上一篇相同;
最后会调到ObservableSubscribeOn.onNext方法, 内部没做处理, 装饰模式,调用上一级的onNext方法

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