Observable o1 = Observable.just(1, 2, 3);
Observable o2 = Observable.just(4, 5, 6);
Observable o3 = Observable.just(7, 8, 9);
Observable.combineLatest(o1, o2, o3, new Function3, Integer, Integer, Integer>() {
public Integer apply(@NonNull Integer integer, @NonNull Integer integer2, @NonNull Integer integer3) throws Exception {
return integer + integer2 + integer3;
}
}).subscribe(RxUtils.getObserver());
运行结果
onSubscribe
Thread:Thread[main,5,main]
onNext:16
Thread:Thread[main,5,main]
onNext:17
Thread:Thread[main,5,main]
onNext:18
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]
//产生0,5,10,15,20数列
Observable observable1 = Observable.interval(500,TimeUnit.MILLISECONDS)
// .delay(600,TimeUnit.MILLISECONDS)
.map(new Function, Long>() {
public Long apply(@NonNull Long aLong) throws Exception {
return aLong*5;
}
});
// 产生0,10,20,30,40数列
Observable observable2=Observable.interval(500,TimeUnit.MILLISECONDS)
// .delay(1000,TimeUnit.MILLISECONDS)
.map(new Function, Long>() {
public Long apply(@NonNull Long aLong) throws Exception {
return aLong*10;
}
});
// observable2.subscribe(RxUtils.getObserver());
observable1.join(observable2, new Function, ObservableSource>() {
public ObservableSource apply(@NonNull Long aLong) throws Exception {
return Observable.just(String.valueOf(aLong)).delay(600, TimeUnit.MILLISECONDS);
}
}, new Function, ObservableSource>() {
public ObservableSource apply(@NonNull Long aLong) throws Exception {
return Observable.just(String.valueOf(aLong));
}
}, new BiFunction, Long, String>() {
public String apply(@NonNull Long aLong, @NonNull Long aLong2) throws Exception {
return aLong + ":" + aLong2;
}
})
.subscribe(RxUtils.getObserver());
try {
Thread.sleep(Integer.MAX_VALUE);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
onSubscribe
Thread:Thread[main,5,main]
onNext:0:10
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:5:10
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:5:20
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:10:20
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:10:30
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:15:30
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:15:40
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:20:40
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:20:50
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:25:50
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:25:60
Thread:Thread[RxComputationThreadPool-2,5,main]
//产生0,5,10,15,20数列
Observable observable1 = Observable.interval(500, TimeUnit.MILLISECONDS)
// .delay(600,TimeUnit.MILLISECONDS)
.map(new Function, Long>() {
public Long apply(@NonNull Long aLong) throws Exception {
return aLong*5;
}
});
// 产生0,10,20,30,40数列
Observable observable2=Observable.interval(500,TimeUnit.MILLISECONDS)
.delay(1000,TimeUnit.MILLISECONDS)
.map(new Function, Long>() {
public Long apply(@NonNull Long aLong) throws Exception {
return aLong*10;
}
});
Observable.merge(observable1,observable2).subscribe(RxUtils.getObserver());
try {
Thread.sleep(Integer.MAX_VALUE);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
onSubscribe
Thread:Thread[main,5,main]
onNext:0
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:5
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:10
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:0
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:15
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:10
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:20
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:20
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:25
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:30
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:30
Observable.just(10, 20, 30).startWith(2).subscribe(RxUtils.getObserver());
onSubscribe
Thread:Thread[main,5,main]
onNext:2
Thread:Thread[main,5,main]
onNext:10
Thread:Thread[main,5,main]
onNext:20
Thread:Thread[main,5,main]
onNext:30
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]
// 每隔500毫秒产生一个observable
Observable> observable= Observable.interval(500, TimeUnit.MILLISECONDS).map(new Function, Observable>() {
public Observable apply(@NonNull Long aLong) throws Exception {
// 每隔250毫秒产生一组数据(0,10,20,30,40)
return Observable.interval(250,TimeUnit.MILLISECONDS).map(new Function,Long>() {
public Long apply(@NonNull Long l) throws Exception {
return l * 10;
}
}).take(5);
}
}).take(2);
// observable.subscribe(RxUtils.>getObserver());
Observable.switchOnNext(observable).subscribe(RxUtils.getObserver(Long.valueOf(1)));
try {
Thread.sleep(Integer.MAX_VALUE);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
onSubscribe:1
Thread:Thread[main,5,main]
onNext:0
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:0
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:10
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:20
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:30
Thread:Thread[RxComputationThreadPool-3,5,main]
onNext:40
Thread:Thread[RxComputationThreadPool-3,5,main]
onComplete
Thread:Thread[RxComputationThreadPool-3,5,main]
## zip
## 通过一个函数将多个Observables的发射物结合到一起,基于这个函数的结果为每个结合体发射单个数据项。
```java
Observable observable1 = Observable.just(10,20,30);
Observable observable2 = Observable.just(4, 8, 12, 16);
Observable.zip(observable1, observable2, new BiFunction, Integer, Integer>() {
public Integer apply(@NonNull Integer integer, @NonNull Integer integer2) throws Exception {
return integer+integer2;
}
}).subscribe(RxUtils.getObserver());
onSubscribe
Thread:Thread[main,5,main]
onNext:14
Thread:Thread[main,5,main]
onNext:28
Thread:Thread[main,5,main]
onNext:42
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]