Rx系列<第十五篇>:RxJava之组合/合并操作符

(1)concat和concatArray

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

    List<Observable<Integer>> list = new ArrayList<>();
    list.add(Observable.just(1,2));
    list.add(Observable.just(3, 4));
    list.add(Observable.just(5, 6));
    Observable.concat(list)
    .subscribe(new Consumer<Integer>() {
        @Override
        public void accept(Integer integer) throws Exception {
            Log.d("aaa", String.valueOf(integer));
        }
    });


    Observable.concatArray(Observable.just(1, 2), Observable.just(3, 4), Observable.just(5, 6))
    .subscribe(new Consumer<Integer>() {
        @Override
        public void accept(Integer integer) throws Exception {
            Log.d("aaa", String.valueOf(integer));
        }
    });

以上两种方式的合并,返回结果是:1 2 3 4 5 6

(2)merge和mergeArray

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

    List<Observable<Integer>> list = new ArrayList<>();
    list.add(Observable.just(1, 2).delay(2000, TimeUnit.MILLISECONDS));
    list.add(Observable.just(3, 4));
    list.add(Observable.just(5, 6));
    list.add(Observable.just(7, 8));
    Observable
            .merge(list)
            .subscribe(new Consumer<Integer>() {
                @Override
                public void accept(Integer integer) throws Exception {
                    Log.d("aaa", String.valueOf(integer));
                }
            });


    Observable
            .mergeArray(Observable.just(1, 2).delay(2000, TimeUnit.MILLISECONDS), Observable.just(3, 4), Observable.just(5, 6), Observable.just(7, 8))
            .subscribe(new Consumer<Integer>() {
                @Override
                public void accept(Integer integer) throws Exception {
                    Log.d("aaa", String.valueOf(integer));
                }
            });

以上两种方法的返回结果是:3 4 5 6 7 8 1 2

(3)concatDelayError和mergeDelayError

concatDelayError:多个Observable合并,并按顺序发射数据, 如果发生异常,则不会立即中断发射数据,异常将延迟发射。
mergeDelayError:多个Observable合并,并行发射数据, 如果发生异常,则不会立即中断发射数据,异常将延迟发射。

    List<Observable<Integer>> list = new ArrayList<>();
    list.add(Observable.create(new ObservableOnSubscribe<Integer>() {
        @Override
        public void subscribe(ObservableEmitter<Integer> e) throws Exception {
            e.onNext(1);
            e.onNext(2);
            e.onError(new NullPointerException("exception"));
            e.onComplete();
        }
    }).delay(2000, TimeUnit.MILLISECONDS));

    list.add(Observable.create(new ObservableOnSubscribe<Integer>() {
        @Override
        public void subscribe(ObservableEmitter<Integer> e) throws Exception {
            e.onNext(3);
            e.onNext(4);
            e.onError(new NullPointerException("exception"));
            e.onComplete();
        }
    }));

    list.add(Observable.create(new ObservableOnSubscribe<Integer>() {
        @Override
        public void subscribe(ObservableEmitter<Integer> e) throws Exception {
            e.onNext(5);
            e.onNext(6);
            e.onError(new NullPointerException("exception"));
            e.onComplete();
        }
    }));

    Observable
            .concatDelayError(list)
            .subscribe(new Consumer<Integer>() {
        @Override
        public void accept(Integer integer) throws Exception {
            Log.d("aaa", String.valueOf(integer));
        }
    }, new Consumer<Throwable>() {
        @Override
        public void accept(Throwable throwable) throws Exception {
            Log.d("aaa", "发生了异常");
        }
    });

返回结果是:

图片.png
(4)zip

合并多个被观察者的数据流, 然后发送(Emit)最终合并的数据。(数据和数据之间是一对一的关系)

    Observable observable1=Observable.create(new ObservableOnSubscribe<Integer>() {
        @Override
        public void subscribe(ObservableEmitter<Integer> e) throws Exception {
            e.onNext(1);
            SystemClock.sleep(1000);
            e.onNext(2);
            SystemClock.sleep(1000);
            e.onNext(3);
            SystemClock.sleep(1000);
            e.onNext(4);
            SystemClock.sleep(1000);
            e.onComplete();
        }
    }).subscribeOn(Schedulers.io());

    Observable observable2=Observable.create(new ObservableOnSubscribe<String>() {
        @Override
        public void subscribe(ObservableEmitter<String> e) throws Exception {
            e.onNext("A");
            SystemClock.sleep(1000);
            e.onNext("B");
            SystemClock.sleep(1000);
            e.onNext("C");
            SystemClock.sleep(1000);
            e.onComplete();
        }
    }).subscribeOn(Schedulers.io());

    Observable.zip(observable1, observable2, new BiFunction<Integer,String,String>() {
        @Override
        public String apply(Integer a,String b) throws Exception {
            return a+b;
        }
    }).subscribe(new Consumer<String>() {
        @Override
        public void accept(String s) throws Exception {
            Log.d("aaa", s);
        }
    });

返回结果:

图片.png
(5)combineLatest

按照同一时间线来进行合并。

    Observable observable1=Observable.create(new ObservableOnSubscribe<Integer>() {
        @Override
        public void subscribe(ObservableEmitter<Integer> e) throws Exception {
            SystemClock.sleep(700);
            e.onNext(1);
            SystemClock.sleep(700);
            e.onNext(2);
            SystemClock.sleep(700);
            e.onNext(3);
            SystemClock.sleep(700);
            e.onNext(4);
            SystemClock.sleep(700);
            e.onComplete();
        }
    }).subscribeOn(Schedulers.io());

    Observable observable2=Observable.create(new ObservableOnSubscribe<String>() {
        @Override
        public void subscribe(ObservableEmitter<String> e) throws Exception {
            e.onNext("A");
            SystemClock.sleep(600);
            e.onNext("B");
            SystemClock.sleep(600);
            e.onNext("C");
            SystemClock.sleep(600);
            e.onComplete();
        }
    }).subscribeOn(Schedulers.io());

    Observable.combineLatest(observable1, observable2, new BiFunction<Integer,String,String>() {
        @Override
        public String apply(Integer a,String b) throws Exception {
            return a+b;
        }
    }).subscribe(new Consumer<String>() {
        @Override
        public void accept(String s) throws Exception {
            Log.d("aaa", s);
        }
    });

接下来,根据代码,我们来画一张图。

画两个时间线observable1和observable2,根据代码中指定的时间画时间线,最后观察两个被观察者时间线重合的地方。

图片.png

实际上代码输出的结果也是:

图片.png
(6)combineLatestDelayError

作用类似于concatDelayError() / mergeDelayError() ,即错误处理,上面已经介绍过类似的了。

(7)reduce

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

    Observable.just(1,2,3,4,5)
            .reduce(new BiFunction<Integer, Integer, Integer>() {
                @Override
                public Integer apply(Integer integer, Integer integer2) throws Exception {
                    return integer + integer2;
                }
            })
            .subscribe(new Consumer<Integer>() {
                @Override
                public void accept(Integer integer) throws Exception {
                    Log.d("aaa", String.valueOf(integer));
                }
            });

日志如下:

图片.png

相当于做了一个这样的计算1+2+3+4+5 = 15,再将15发射出去。

(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> list) throws Exception {
            for (int result : list){
                Log.d("aaa", String.valueOf(result));
            }
        }
    });

执行结果:

图片.png
(9)startWith和startWithArray

startWith: 在已有数据流之前追加一个或一组数据流。

图片.png

startWith可以传递的参数是:一个数据,一个数据列表,一个Observable。

    Observable.just(1, 2, 3)
            .startWith(4)
            .subscribe(new Consumer<Integer>() {
                @Override
                public void accept(Integer integer) throws Exception {
                    Log.d("aaa", String.valueOf(integer));
                }
            });

返回结果: 4 1 2 3

startWithArray:在已有数据流之前追加一组数据流。

图片.png
    Observable.just(1, 2, 3)
            .startWithArray(4, 5)
            .subscribe(new Consumer<Integer>() {
                @Override
                public void accept(Integer integer) throws Exception {
                    Log.d("aaa", String.valueOf(integer));
                }
            });

返回结果:4 5 1 2 3

(10)count

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

    Observable.just(1, 2, 3, 4)
            .count()
            .subscribe(new Consumer<Long>() {
                @Override
                public void accept(Long aLong) throws Exception {
                    Log.d("aaa", String.valueOf(aLong));
                }
            });

返回结果:4

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

推荐阅读更多精彩内容

  • 一、RxJava操作符概述 RxJava中的操作符就是为了提供函数式的特性,函数式最大的好处就是使得数据处理简洁易...
    BrotherChen阅读 1,606评论 0 10
  • 一、RxJava操作符概述 RxJava中的操作符就是为了提供函数式的特性,函数式最大的好处就是使得数据处理简洁易...
    无求_95dd阅读 3,045评论 0 21
  • 创建操作 用于创建Observable的操作符Create通过调用观察者的方法从头创建一个ObservableEm...
    rkua阅读 1,817评论 0 1
  • 注:只包含标准包中的操作符,用于个人学习及备忘参考博客:http://blog.csdn.net/maplejaw...
    小白要超神阅读 2,190评论 2 8
  • 记录RxJava操作符,方便查询(2.2.2版本) 英文文档地址:http://reactivex.io/docu...
    凌云飞鱼阅读 819评论 0 0