实现简单的 RxKotlin (中)

线程切换的操作在 Rx 里面非常常用,主要有 subscribeOn observeOn 他们都需要一个 Scheduler 参数,很明显这个是接口 可以实现各种 调度器。根据这个我们可以这样写...

//Scheduler.kt
abstract class Scheduler {

    abstract fun createWorker(): Worker

    abstract class Worker {
        abstract fun schedule(action: () -> Unit)
    }
}

Worker 是实际工作的,有个抽象方法 schedule 需要我们去实现
Executor 是 Java 里面线程池的执行器, Executors 里面有各种已经实现的线程池,至此 我们可以来实现一个ExecutorScheduler

//ExecutorScheduler.kt
class ExecutorScheduler(private var executor: Executor) : Scheduler() {

    override fun createWorker(): Worker {
        return WorkerImpl(executor)
    }

    private class WorkerImpl(private var executor: Executor) : Worker() {

        override fun schedule(action: () -> Unit) {
            executor.execute {
                action()
            }
        }
    }
}

参考 rxjava 建立一个 Schedulers 存放 Scheduler,这样就直接可以调用 Schedulers.io()

class Schedulers {
    companion object {

        private val io = ExecutorScheduler(Executors.newSingleThreadExecutor())
       
        fun io(): Scheduler {
            return io
        }
    }
}

实现 Android 特有 的 AndroidScheduler,主要利用 Looper

//LooperScheduler.kt
class LooperScheduler(looper: Looper) : Scheduler() {

    private  var handler : Handler = Handler(looper)

    override fun createWorker(): Worker = LooperSchedulerWorker(handler)

    private class LooperSchedulerWorker(private var handler : Handler) : Worker() {
        override fun schedule(action: () -> Unit) {
            handler.post {
                action()
            }
        }
    }

}
//AndroidSchedulers.kt
class AndroidSchedulers {

    companion object {

        private val mainThread = LooperScheduler(Looper.getMainLooper())

        fun mainThread(): Scheduler {
            return mainThread
        }
    }
}

实现 subscribeOn observeOn

//Observable.kt
fun <R> lift(operator: Operator<R, T>): Observable<R> {
    return create(OnSubscribeLift(onSubscribe!!, operator))
}

fun subscribeOn(scheduler: Scheduler): Observable<T> {
    return create(OperatorSubscribeOn(this, scheduler))
}

fun observeOn(scheduler: Scheduler): Observable<T> {
    return lift(OperatorObserveOn(scheduler))
}

interface Operator<T, R> {
    fun call(subscriber: Subscriber<T>) : Subscriber<R>
}

//OnSubscribeLift.kt
class OnSubscribeLift<T, R>(private var parent: Observable.OnSubscribe<T>, private var operator: Observable.Operator<R, T>) : Observable.OnSubscribe<R>{

    override fun call(subscriber: Subscriber<R>) {
        try {
            val st = operator.call(subscriber)
            st.onStart()
            parent.call(st)
        }catch (e: Exception) {
            subscriber.onError(e)
        }
    }
}

在 rxjava 里面 subscribeOn 只有调用的第一次起作用,而 observeOn 则是看最近的那次调用。
原因在于 subscribeOn 调度的是执行 OnSubscribe 因此多次调用最上游的还是 在 第一次的线程那里执行的,除非更改下游 观察者所在的线程 也就是 observeOn。如果不明白的,可以多看这部分代码或者 rxjava 的源码加深理解。

//OperatorSubscribeOn.kt
class OperatorSubscribeOn<T>(private var source: Observable<T> , private var scheduler: Scheduler) : Observable.OnSubscribe<T>{

    override fun call(subscriber: Subscriber<T>) {
        val worker = scheduler.createWorker()
        worker.schedule {
            source.subscribe(SubscribeOnSubscriber(subscriber))
        }
    }

    class SubscribeOnSubscriber<T>(private var actual: Subscriber<T>): Subscriber<T>() {
        override fun onCompleted() {
            actual.onCompleted()
        }

        override fun onError(t: Throwable) {
            actual.onError(t)
        }

        override fun onNext(t: T) {
            actual.onNext(t)
        }
    }
}

//OperatorObserveOn.kt
class  OperatorObserveOn<T>(private var scheduler: Scheduler): Observable.Operator<T, T> {

    override fun call(subscriber: Subscriber<T>): Subscriber<T> {
        val worker = scheduler.createWorker()
        return object : Subscriber<T>(){
            override fun onCompleted() {
                worker.schedule {
                    subscriber.onCompleted()
                }
            }

            override fun onError(t: Throwable) {
                worker.schedule {
                    subscriber.onError(t)
                }
            }

            override fun onNext(t: T) {
                worker.schedule {
                    subscriber.onNext(t)
                }
            }
        }
    }
}
//Test.kt
Observable.just("1", "2", "3")
                .subscribeOn(Schedulers.io())
                .map {
                    it.toInt() + 1
                }
                .filter {
                    it != 1
                }
                .map {
                    it
                }
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe(object : Subscriber<Int>() {
                    override fun onCompleted() {

                    }

                    override fun onError(t: Throwable) {
                    }

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

推荐阅读更多精彩内容

  • 前言我从去年开始使用 RxJava ,到现在一年多了。今年加入了 Flipboard 后,看到 Flipboard...
    占导zqq阅读 9,159评论 6 151
  • 我从去年开始使用 RxJava ,到现在一年多了。今年加入了 Flipboard 后,看到 Flipboard 的...
    Jason_andy阅读 5,460评论 7 62
  • 引入依赖: implementation 'io.reactivex.rxjava2:rxandroid:2.0....
    为梦想战斗阅读 1,300评论 0 0
  • 最近项目里面有用到Rxjava框架,感觉很强大的巨作,所以在网上搜了很多相关文章,发现一片文章很不错,今天把这篇文...
    Scus阅读 6,868评论 2 50
  • 面试官:“如何制作一个菱形的Button,比如现在button的背景图是个菱形,如何实现点击图片中的菱形内有响应而...
    keyser_fayee阅读 1,958评论 0 0