手写简单Rxjava理解其内部实现(三)

上一篇我们实现了,操作符Map功能。本篇实现线程切换操作符subscribeOn及observeOn。

  • 创建抽象执行Runnable的Schedule
abstract class Scheduler {

    abstract fun createWorker(): Worker

    fun scheduleDirect(task: Runnable) {
        val worker = createWorker()
        worker.schedule(task)
    }

    interface Worker {
        fun schedule(runnable: Runnable)
    }
}
  • 创建主线程及子线程执行Schedule
class HandlerScheduler(var handler: Handler) : Scheduler() {

    override fun createWorker(): Worker {
        return HandlerWorker(handler)
    }

    class HandlerWorker(var handler: Handler) : Worker {

        override fun schedule(runnable: Runnable) {
            val message = Message.obtain(handler, runnable)
            message.obj = this
            handler.sendMessage(message)
        }

    }
}
class NewThreadScheduler : Scheduler() {

    override fun createWorker(): Worker {
        return NewThreadWorker()
    }

    class NewThreadWorker : Worker {
        var executorService: ExecutorService? = null

        init {
            executorService = Executors.newScheduledThreadPool(2)
        }

        override fun schedule(runnable: Runnable) {
            executorService?.execute(runnable)
        }
    }
}
  • 创建生产主线程、子线程的Schedulers
class Schedulers {
    companion object {
        private val MAIN_THREAD = HandlerScheduler(Handler(Looper.getMainLooper()))
        private val NEW_THREAD = NewThreadScheduler()

        fun mainThread(): Scheduler {
            return MAIN_THREAD
        }

        fun newThread(): Scheduler {
            return NEW_THREAD
        }
    }
}
  • 实现SubscribeOn的观察者及被观察者,同时创建一执行Runnable的任务
class ObservableSubscribeOn<T>(
    observableSource: ObservableSource<T>,
    private val scheduler: Scheduler
) : AbstractObservableWithUpStream<T, T>(observableSource) {

    override fun subscribeActual(observer: Observer<T>) {
        //将订阅逻辑抽离到一个Runnable里
        scheduler.scheduleDirect(SubscribeTask(observableSource, SubscribeOnObserver(observer)))
    }

    class SubscribeOnObserver<T>(downstream: Observer<T>) : BasicFuseabObserver<T, T>(downstream)

    class SubscribeTask<T>(
        private val observableSource: ObservableSource<T>,
        private val subscribeOnObserver: SubscribeOnObserver<T>
    ) : Runnable {
      //真正执行订阅逻辑的Runnable,运行线程决定了订阅线程
        override fun run() {
            observableSource.subscribe(subscribeOnObserver)
        }
    }
}
  • 实现ObserveOn的观察者及被观察者
class ObservableObserveOn<T>(
    observableSource: ObservableSource<T>,
    private val scheduler: Scheduler
) : AbstractObservableWithUpStream<T, T>(observableSource) {

    override fun subscribeActual(observer: Observer<T>) {
        val worker = scheduler.createWorker()
        observableSource.subscribe(ObserveOnObserver(observer, worker))
    }

    class ObserveOnObserver<T>(observer: Observer<T>, var worker: Scheduler.Worker) :
        BasicFuseabObserver<T, T>(observer), Runnable {

        @Volatile
        var done = false
        private var queue: ArrayDeque<T>? = null

        @Volatile
        var error: Throwable? = null

        @Volatile
        var over = false

        init {
            queue = ArrayDeque()
        }

        override fun onSubscribe() {
            oberver.onSubscribe()
            schedule()
        }

        override fun onNext(t: T) {
            if (done) {
                return
            }
            queue?.add(t)
            schedule()
        }

        override fun onComplete() {
            if (done) {
                return
            }
            done = true
            schedule()
        }

        override fun onError(t: Throwable) {
            if (done) {
                return
            }
            done = true
            error = t
            schedule()
        }

        override fun run() {
            drainNormal()
        }
        //执行了线程的切换
        private fun schedule() {
            worker.schedule(this)
        }
        //观察者的数据观察
        private fun drainNormal() {
            var arrayDeque = queue
            var a = oberver

            while (true) {
                var d = done
                var t = arrayDeque?.removeAt(0)
                val empty = t == null
                if (checkTerminated(d, empty, a)) {
                    return
                }
                if (t == null) {
                    break
                }
                a.onNext(t)
            }
        }

        /**
         * 判断是否终止
         */
        private fun checkTerminated(d: Boolean, empty: Boolean, a: Observer<T>): Boolean {
            if (over) {
                queue?.clear()
                return true
            }
            if (d) {
                var e = error
                if (e is Throwable) {
                    over = true
                    a.onError(e)
                    return true
                } else if (empty) {
                    over = true
                    a.onComplete()
                    return true
                }
            }
            return false
        }
    }
}
  • 通过以上实现,我们就应该能知道subscribeOn及observeOn对于线程切换时,真正起作用的逻辑。
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 214,588评论 6 496
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 91,456评论 3 389
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 160,146评论 0 350
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 57,387评论 1 288
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 66,481评论 6 386
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 50,510评论 1 293
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,522评论 3 414
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,296评论 0 270
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 44,745评论 1 307
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,039评论 2 330
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,202评论 1 343
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 34,901评论 5 338
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,538评论 3 322
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,165评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,415评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,081评论 2 365
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,085评论 2 352

推荐阅读更多精彩内容