RxJava操作符--->组合/合并


引言

该篇文章主要是关于RxJava的组合/变换操作符使用的代码讲解。组合/变换操作符总共有四大类:

(1)组合多个被观察者

  • 按发送顺序:concat()、concatArray()
  • 按时间:merge()、mergeArray()
  • 错误处理:concatDelayError()、mergeDelayError()

(2)合并多个事件

  • 按数量:zip()
  • 按时间:combineLatest()、combineLatestDelayError()
  • 合并成一个事件发送:reduce()、collect()

(3)发送事件前追加发送事件

  • startWith()
  • startWithArray()

(4)统计发送事件数量

  • count()

1. concat()/concatArray()

组合多个被观察者一起发送数据,合并后按发送顺序串行执行。

二者区别:组合被观察者的数量,即concat()组合被观察者数量≤4个,而concatArray()则可>4个。

        Observable.concat(
                Observable.just(1,2,3),
                Observable.just(4,5,6),
                Observable.just(7,8,9),
                Observable.just(10,11,12))
                .subscribe(new Observer<Integer>() {
                    @Override
                    public void onSubscribe(Disposable d) {

                    }

                    @Override
                    public void onNext(Integer value) {
                        Log.d(Constant.TAG,"接收到了事件"+value);
                    }

                    @Override
                    public void onError(Throwable e) {
                        Log.d(Constant.TAG,"对Error事件做出响应");
                    }

                    @Override
                    public void onComplete() {
                        Log.d(Constant.TAG,"对Complete事件做出响应");
                    }
                });

        Observable.concatArray(Observable.just(1,2),
                Observable.just(3,4),
                Observable.just(5,6),
                Observable.just(7,8),
                Observable.just(9,10))
                .subscribe(new Observer<Integer>() {
                    @Override
                    public void onSubscribe(Disposable d) {

                    }

                    @Override
                    public void onNext(Integer value) {
                        Log.d(Constant.TAG,"接收到了事件"+value);
                    }

                    @Override
                    public void onError(Throwable e) {
                        Log.d(Constant.TAG,"对Error事件做出响应");
                    }

                    @Override
                    public void onComplete() {
                        Log.d(Constant.TAG,"对Complete事件做出响应");
                    }
                });

concat()的log信息:

06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件1
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件2
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件3
06-22 14:00:51.142 12967-12967/? D/RxJava: 接收到了事件4
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件5
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件6
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件7
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件8
06-22 14:00:51.143 12967-12967/? D/RxJava: 对Complete事件做出响应

concatArray()的log信息:

06-22 14:00:51.143 12967-12967/? D/RxJava: 对Complete事件做出响应
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件1
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件2
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件3
06-22 14:00:51.144 12967-12967/? D/RxJava: 接收到了事件4
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件5
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件6
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件7
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件8
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件9
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件10
06-22 14:00:51.143 12967-12967/? D/RxJava: 对Complete事件做出响应

2. merge()/mergeArray()

组合多个被观察者一起发送数据,合并后按时间线并行执行。

1.二者区别:和上述的concat和concatArray的一样;
2.区别上述concat操作符,同样是组合多个被观察者一起发送数据,但concat操作符合并后是按发送顺序串行执行。

        Observable.merge(
                Observable.intervalRange(0,3,1,1, TimeUnit.SECONDS),
                Observable.intervalRange(2,3,1,1, TimeUnit.SECONDS))
                .subscribe(new Observer<Long>() {
                    @Override
                    public void onSubscribe(Disposable d) {

                    }

                    @Override
                    public void onNext(Long value) {
                        Log.d(Constant.TAG,"接收到了事件"+value);
                    }

                    @Override
                    public void onError(Throwable e) {
                        Log.d(Constant.TAG,"对Error事件做出响应");
                    }

                    @Override
                    public void onComplete() {
                        Log.d(Constant.TAG,"对Complete事件做出响应");
                    }
                });

log信息:

06-22 14:23:11.358 14031-14082/com.gjj.frame D/RxJava: 接收到了事件0
06-22 14:23:11.366 14031-14083/com.gjj.frame D/RxJava: 接收到了事件2
06-22 14:23:12.357 14031-14082/com.gjj.frame D/RxJava: 接收到了事件1
06-22 14:23:12.358 14031-14082/com.gjj.frame D/RxJava: 接收到了事件3
06-22 14:23:13.358 14031-14082/com.gjj.frame D/RxJava: 接收到了事件2
06-22 14:23:13.359 14031-14082/com.gjj.frame D/RxJava: 接收到了事件4
06-22 14:23:13.362 14031-14082/com.gjj.frame D/RxJava: 对Complete事件做出响应

3. concatArrayDelayError()/mergeArrayDelayError()

使用concat和merge操作符时,若其中一个被观察者发出onError事件,则会马上终止其他被观察者继续发送事件,若希望onError事件推迟到其他被观察者发送事件结束后才处罚,就需要使用对应的concatDelayError或mergeDelayError()操作符。

(1)无使用concatArrayDelayError()的情况

        Observable.concat(Observable.create(new ObservableOnSubscribe<Integer>() {
            @Override
            public void subscribe(ObservableEmitter<Integer> e) throws Exception {
                e.onNext(1);
                e.onNext(2);
                e.onNext(3);
                //对error事件,因为无使用concatDelayError,所以第二个Observable将不会发送事件
                e.onError(new NullPointerException());
                e.onComplete();
            }
        }),Observable.just(4,5,6)).subscribe(new Observer<Integer>() {
            @Override
            public void onSubscribe(Disposable d) {

            }

            @Override
            public void onNext(Integer value) {
                Log.d(Constant.TAG,"接收到了事件"+value);
           }

            @Override
            public void onError(Throwable e) {
                Log.d(Constant.TAG,"对error事件做出响应");
            }

            @Override
            public void onComplete() {
                Log.d(Constant.TAG,"对Complete事件做出响应");
            }
        });

测试结果:第一个悲观者发送Error事件后,第2个被观察者则不会继续发送事件。
log信息:

06-25 11:03:06.905 21337-21337/com.gjj.frame D/RxJava: 接收到了事件1
06-25 11:03:06.905 21337-21337/com.gjj.frame D/RxJava: 接收到了事件2
06-25 11:03:06.905 21337-21337/com.gjj.frame D/RxJava: 接收到了事件3
06-25 11:03:06.906 21337-21337/com.gjj.frame D/RxJava: 对error事件做出响应

(2)使用concatArrayDelayError()的情况

        Observable.concatArrayDelayError(Observable.create(new ObservableOnSubscribe<Integer>() {
            @Override
            public void subscribe(ObservableEmitter<Integer> e) throws Exception {
                e.onNext(1);
                e.onNext(2);
                e.onNext(3);
                //对error事件,因为无使用concatDelayError,所以第二个Observable将不会发送事件
                e.onError(new NullPointerException());
                e.onComplete();
            }
        }),Observable.just(4,5,6)).subscribe(new Observer<Integer>() {
            @Override
            public void onSubscribe(Disposable d) {

            }

            @Override
            public void onNext(Integer value) {
                Log.d(Constant.TAG,"接收到了事件"+value);
           }

            @Override
            public void onError(Throwable e) {
                Log.d(Constant.TAG,"对error事件做出响应");
            }

            @Override
            public void onComplete() {
                Log.d(Constant.TAG,"对Complete事件做出响应");
            }
        });

测试结果:第1个被观察者的error事件将在第2个被观察者发送完事件后再继续发送。
log信息:

06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件1
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件2
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件3
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件4
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件5
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件6
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 对error事件做出响应

4. Zip()

合并多个被观察者(Observable)发送的事件,生成一个新的事件序列(即组合过后的事件序列),并最终发送。

        //创建第1个观察者
        Observable<Integer> observable1 = Observable.create(new ObservableOnSubscribe<Integer>() {
            @Override
            public void subscribe(ObservableEmitter<Integer> e) throws Exception {
                e.onNext(1);
                e.onNext(2);
                e.onNext(3);
            }
        }).subscribeOn(Schedulers.io());//设置被观察者1再工作线程1中工作

        //创建第2个观察者
        Observable<String> observable2 = Observable.create(new ObservableOnSubscribe<String>() {
            @Override
            public void subscribe(ObservableEmitter<String> e) throws Exception {
                e.onNext("A");
                e.onNext("B");
                e.onNext("C");
                e.onNext("D");
                e.onComplete();
            }
        }).subscribeOn(Schedulers.newThread());//设置被观察者2再工作线程2中工作

        Observable.zip(observable1, observable2, new BiFunction<Integer, String, String>() {
            @Override
            public String apply(Integer integer, String s) throws Exception {
                return integer+s;
            }
        }).subscribe(new Observer<String>() {
            @Override
            public void onSubscribe(Disposable d) {
            }

            @Override
            public void onNext(String value) {
                Log.d(Constant.TAG,"最终收到的事件 = "+value);
            }

            @Override
            public void onError(Throwable e) {
                Log.d(Constant.TAG,"onError");

            }

            @Override
            public void onComplete() {
                Log.d(Constant.TAG,"onComplete");

            }
        });

log信息:

06-26 16:30:02.147 29926-29985/com.gjj.frame D/RxJava: 最终收到的事件 = 1A
06-26 16:30:02.150 29926-29984/com.gjj.frame D/RxJava: 最终收到的事件 = 2B
06-26 16:30:02.151 29926-29984/com.gjj.frame D/RxJava: 最终收到的事件 = 3C

注意:最终合并的事件数量是多个被观察者中最少的数量,多余的事件将不会发送。

5. combineLatest()

当两个Observable中的任何一个发送了数据后,将先发送了数据的Observables的最新(最后)一个数据与另外一个Observable发送的每一个数据结合,最终基于该函数的结果发送数据。

        Observable.combineLatest(Observable.just(1L, 2L, 3L), Observable.intervalRange(0, 3, 1, 1, TimeUnit.SECONDS), new BiFunction<Long, Long, Long>() {
            @Override
            public Long apply(Long aLong, Long aLong2) throws Exception {
                Log.d(Constant.TAG,"合并的数据是:"+aLong+" "+aLong2);
                return aLong+aLong2;
            }
        }).subscribe(new Consumer<Long>() {
            @Override
            public void accept(Long aLong) throws Exception {
                Log.d(Constant.TAG,"合并的结果是:"+aLong);
            }
        });

log信息:

06-26 16:48:37.010 30604-30634/com.gjj.frame D/RxJava: 合并的数据是:3 0
06-26 16:48:37.011 30604-30634/com.gjj.frame D/RxJava: 合并的结果是:3
06-26 16:48:38.010 30604-30634/com.gjj.frame D/RxJava: 合并的数据是:3 1
06-26 16:48:38.011 30604-30634/com.gjj.frame D/RxJava: 合并的结果是:4
06-26 16:48:39.012 30604-30634/com.gjj.frame D/RxJava: 合并的数据是:3 2
06-26 16:48:39.013 30604-30634/com.gjj.frame D/RxJava: 合并的结果是:5

6. combineLatestDelayError()

作用类似于concatArrayDelayError()。

7. reduce()

把被观察者需要发送的事件聚合成一个事件&发送

        Observable.just(1,2,3,4)
                .reduce(new BiFunction<Integer, Integer, Integer>() {
                    @Override
                    public Integer apply(Integer integer, Integer integer2) throws Exception {
                        Log.d(Constant.TAG,"本次计算的数据是:"+integer+"乘"+integer2);
                        return integer * integer2;
                    }
                }).subscribe(new Consumer<Integer>() {
            @Override
            public void accept(Integer integer) throws Exception {
                Log.d(Constant.TAG,"最终计算的结果是:"+integer);
            }
        });

log信息:

06-26 16:59:56.401 31613-31613/com.gjj.frame D/RxJava: 本次计算的数据是:1乘2
06-26 16:59:56.402 31613-31613/com.gjj.frame D/RxJava: 本次计算的数据是:2乘3
06-26 16:59:56.402 31613-31613/com.gjj.frame D/RxJava: 本次计算的数据是:6乘4
06-26 16:59:56.402 31613-31613/com.gjj.frame D/RxJava: 最终计算的结果是:24

8. collect()

将被观察者Observable发送的数据事件收集到一个数据结构里

        Observable.just(1,2,3,4,5,6)
                .collect(new Callable<ArrayList<Integer>>() {
                    @Override
                    public ArrayList<Integer> call() throws Exception {
                        return new ArrayList<>();
                    }
                }, new BiConsumer<ArrayList<Integer>, Integer>() {
                    @Override
                    public void accept(ArrayList<Integer> list, Integer integer) throws Exception {
                        list.add(integer);
                    }
                }).subscribe(new Consumer<ArrayList<Integer>>() {
            @Override
            public void accept(ArrayList<Integer> integers) throws Exception {
                Log.d(Constant.TAG,"本次发送的数据是:"+integers);
            }
        });

log信息:

06-26 17:04:40.264 31785-31785/com.gjj.frame D/RxJava: 本次发送的数据是:[1, 2, 3, 4, 5, 6]

9. startWith()/startWithArray()

在一个被观察者发送事件钱,追加发送一些数据/一个新的被观察者

        Observable.just(3,4)
                .startWith(0)//追加单个数据
                .startWithArray(1,2)//追加多个数据
                .subscribe(new Consumer<Integer>() {
                    @Override
                    public void accept(Integer integer) throws Exception {
                        Log.d(Constant.TAG,"接收到了事件"+integer);
                    }
                });

log信息:

06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件1
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件2
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件0
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件3
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件4

10.count()

统计被观察者发送事件的数量。

        Observable.just(1,2,3,4)
                .count()
                .subscribe(new Consumer<Long>() {
                    @Override
                    public void accept(Long integer) throws Exception {
                        Log.d(Constant.TAG,"发送的事件数量 = "+integer);
                    }
                });

log信息:

06-26 17:42:20.639 4750-4750/com.gjj.frame D/RxJava: 发送的事件数量 = 4

参考文章:
Android RxJava:组合 / 合并操作符 详细教程

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

推荐阅读更多精彩内容