欢迎您访问程序员文章站本站旨在为大家提供分享程序员计算机编程知识!
您现在的位置是: 首页

RxJava基本使用--合并型操作符

程序员文章站 2022-06-30 14:18:29
...

RxJava基本使用–合并型操作符

startWith,concatWith

先创建被观察者,然后再组合其他的被观察者,然后再订阅。

startWith

先执行startWith中的被观察者发送的数据,再执行另一个被观察者。

Observable.create(object : ObservableOnSubscribe<String>{
            override fun subscribe(emitter: ObservableEmitter<String>) {
                emitter.onNext("1")
                emitter.onNext("2")
                emitter.onNext("3")
                emitter.onComplete()
            }

        }).startWith(Observable.create(object : ObservableOnSubscribe<String>{
            override fun subscribe(emitter: ObservableEmitter<String>) {
                emitter.onNext("10000")
                emitter.onNext("20000")
                emitter.onNext("30000")
                emitter.onComplete()
            }

        })).subscribe(object : Consumer<String>{
            override fun accept(t: String?) {
                Log.e("RxJavaActivity","accept=$t")
            }
        })

必须添加emitter.onComplete(),否则无法显示第一个被观察者发送的对象。
RxJava基本使用--合并型操作符

concatWith

与startWith执行相反

Observable.create(object : ObservableOnSubscribe<String>{
            override fun subscribe(emitter: ObservableEmitter<String>) {
                emitter.onNext("1")
                emitter.onNext("2")
                emitter.onNext("3")
                emitter.onComplete()
            }

        }).concatWith(Observable.create(object : ObservableOnSubscribe<String>{
            override fun subscribe(emitter: ObservableEmitter<String>) {
                emitter.onNext("10000")
                emitter.onNext("20000")
                emitter.onNext("30000")
                emitter.onComplete()
            }

        })).subscribe(object : Consumer<String>{
            override fun accept(t: String?) {
                Log.e("RxJavaActivity","accept=$t")
            }
        })

RxJava基本使用--合并型操作符

concat()

RxJava基本使用--合并型操作符
最大能传入4个被观察者。
顺序执行4个被观察者对象。

Observable.concat(Observable.just("1"),
                Observable.just("2"),
                Observable.just("3"),
                Observable.just("4"))
                .subscribe(object : Consumer<String>{
                    override fun accept(t: String?) {
                        Log.e("RxJavaActivity","accept=$t")
                    }
                })

RxJava基本使用--合并型操作符

merge()

最多能合并4个,并发执行
为了体现并发特性,使用intervalRange实验:
参数:start开始累计 count累计多少个数量 initialDelay开始等待时间 period每隔多久执行 TimeUnit时间单位

        var observable1 = Observable.intervalRange(1, 5, 0, 0, TimeUnit.SECONDS)
        var observable2 = Observable.intervalRange(6, 5, 0, 0, TimeUnit.SECONDS)
        var observable3 = Observable.intervalRange(11, 5, 0, 0, TimeUnit.SECONDS)
        var observable4 = Observable.intervalRange(16, 5, 0, 0, TimeUnit.SECONDS)
        
        Observable.merge(observable1,observable2,observable3,observable4)
                .subscribe(object : Consumer<Long>{
                    override fun accept(t: Long?) {
                        Log.e("RxJavaActivity","accept=$t")
                    }
                })

RxJava基本使用--合并型操作符

zip()

组合两个被观察者对象,可以不同类型,但发射的数据必须一一对应,否则不对应的数据将会被忽略。

var observable1 = Observable.create(object : ObservableOnSubscribe<String>{
            override fun subscribe(emitter: ObservableEmitter<String>) {
                emitter.onNext("语文")
                emitter.onNext("数学")
                emitter.onNext("英语")
                emitter.onNext("政治") //被忽略
                emitter.onComplete()
            }
        })

        var observable2 = Observable.create(object : ObservableOnSubscribe<Int>{
            override fun subscribe(emitter: ObservableEmitter<Int>) {
                emitter.onNext(90)
                emitter.onNext(85)
                emitter.onNext(100)
                emitter.onComplete()
            }
        })

        Observable.zip(observable1,observable2,object : BiFunction<String, Int, StringBuffer> {
            override fun apply(t1: String, t2: Int): StringBuffer {
                return StringBuffer().append("课程:$t1,成绩:$t2")
            }
        }).subscribe(object : Consumer<StringBuffer>{
            override fun accept(t: StringBuffer?) {
                Log.e("RxJavaActivity","accept=$t")
            }
        })

RxJava基本使用--合并型操作符