您的位置:首页 > 编程语言 > Java开发

Rxjava结合操作符—merge、 Join

2017-08-23 15:00 369 查看

1、merge

merge可以合并多个发射物

Javadoc: merge(Iterable)

Javadoc: merge(Iterable,int)

Javadoc: merge(Observable[])

Javadoc: merge(Observable,Observable) (接受二到九个Observable)



两个Obserable合并成一个Observable

Observable<Integer> odds=Observable.just(1,3,5,7);
Observable<Integer> events=Observable.just(2,4,6,8);
Observable.merge(odds,events).subscribe(i->Log.d("TAG","merge->"+i));


运行结果:

merge->1
merge->3
merge->5
merge->7
merge->2
merge->4
merge->6
merge->8


2、mergeWith



这个图有个叉,表示如果某个阶段出错了,后续的数据就会停止发射了。

Observable<Integer> odds=Observable.just(1,3,5,7);
Observable<Integer> events=Observable.just(2,4,6,8);
odds.mergeWith(events).subscribe(i->Log.d("TAG","merge->"+i));


运行结果同上

3、mergeDelayError



Observable<Integer> odds=Observable.just(1,3,5,7);
Observable<Integer> events=Observable.just(2,4,6,8);
Observable.mergeDelayError(odds,events).subscribe(i->Log.d("TAG","merge->"+i));


运行结果同上

merge和mergeWith相当于一个是静态方法,一个是实例方法。merge如果以onError错误终止的话,数据也会相应的停止发射。如果想让它继续发射数据,然后才报告错误,可以使用mergeDelayError

注: 合并后的新Observables发射时,可以认为任何一个Observable发生错误,都将会打断合并。如果想避免这种情况发生,可以调用mergeDelayError()方法,表明它能从一个Observable中继续发射数据即便是其中有一个抛出了错误。当所有的Observables都完成时,mergeDelayError()将会发射onError()。

4、join

任何时候,只要在另一个Observable发射的数据定义的时间窗口内,这个Observable发射了一条数据,就结合两个Observable发射的数据



Join操作符结合两个Observable发射的数据,基于时间窗口(你定义的针对每条数据特定的原则)选择待集合的数据项。你将这些时间窗口实现为一些Observables,它们的生命周期从任何一条Observable发射的每一条数据开始。当这个定义时间窗口的Observable发射了一条数据或者完成时,与这条数据关联的窗口也会关闭。只要这条数据的窗口是打开的,它将继续结合其它Observable发射的任何数据项。你定义一个用于结合数据的函数。

join表示形式:

join(Observable, Func1, Func1, Func2)

join有四个参数分别表示:

Observable:源Observable需要组合的Observable,这里我们姑且称之为目标Observable;

Func1:接收从源Observable发射来的数据,并返回一个Observable,这个Observable的声明周期决定了源Obsrvable发射出来的数据的有效期;

Func1:接收目标Observable发射来的数据,并返回一个Observable,这个Observable的声明周期决定了目标Obsrvable发射出来的数据的有效期;

Func2:接收从源Observable和目标Observable发射出来的数据,并将这两个数据组合后返回。

//输出[0,1,2,3]序列
Observable<Integer> ob =Observable.create(subscriber -> {
for (int i = 0; i < 4; i++) {
subscriber.onNext(i);
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
});

Observable.just("hello")
.join(ob, s -> {
Log.d("TAG",s);
//使Observable延迟3000毫秒执行
return Observable.timer(3000, TimeUnit.MILLISECONDS);
}, integer -> {
Log.d("TAG",integer+"");
//使Observable延迟2000毫秒执行
return Observable.timer(2000, TimeUnit.MILLISECONDS);
//结合上面发射的数据
(s, integer) -> s+integer)
.subscribe(o -> Log.d("TAG",o));


运行结果:

0
hello
hello0
1
hello
hello1
2
hello
hello2
3
内容来自用户分享和网络整理,不保证内容的准确性,如有侵权内容,可联系管理员处理 点击这里给我发消息
标签: