List list=new ArrayList();
Observable.range(1,5).buffer(2).subscribe(RxUtils.getObserver(list));
运行结果
onSubscribe:[]
Thread:Thread[main,5,main]
onNext:[1, 2]
Thread:Thread[main,5,main]
onNext:[3, 4]
Thread:Thread[main,5,main]
onNext:[5]
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]
Observable.just("大","小").flatMap(new Function, ObservableSource>() {
public ObservableSource apply(@NonNull String s) throws Exception {
return Observable.just(s+"sb");
}
}).subscribe(RxUtils.getObserver("0"));
onSubscribe:0
Thread:Thread[main,5,main]
onNext:大sb
Thread:Thread[main,5,main]
onNext:小sb
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]
Observable.just("大","小").map(new Function, String>() {
public String apply(@NonNull String s) throws Exception {
return s+"sb";
}
}).subscribe(RxUtils.getObserver("0"));
onSubscribe:0
Thread:Thread[main,5,main]
onNext:大sb
Thread:Thread[main,5,main]
onNext:小sb
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]
Observable,Long>> observable= Observable.interval(1, TimeUnit.SECONDS).take(10).groupBy(new Function, Long>() {
public Long apply(@NonNull Long aLong) throws Exception {
return aLong%3;
}
});
observable.subscribe(new Observer, Long>>() {
public void onSubscribe(@NonNull Disposable d) {
System.out.println("onSubscribe:");
}
public void onNext(@NonNull final GroupedObservable, Long> longLongGroupedObservable) {
longLongGroupedObservable.subscribe(new Consumer() {
public void accept(@NonNull Long aLong) throws Exception {
System.out.println("key:" + longLongGroupedObservable.getKey() +", value:" + aLong);
}
});
}
public void onError(@NonNull Throwable e) {
System.out.println("onNext:"+e);
}
public void onComplete() {
}
});
try {
Thread.sleep(Integer.MAX_VALUE);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
运行结果
onSubscribe:
key:0, value:0
key:1, value:1
key:2, value:2
key:0, value:3
key:1, value:4
key:2, value:5
key:0, value:6
key:1, value:7
key:2, value:8
key:0, value:9
可能代码看的有点晕,看下官方给的图就会豁然开朗
不明白?这里走 再看看 官网的文档
Observable.just(1,2,3,4,5).scan(new BiFunction, Integer, Integer>() {
public Integer apply(@NonNull Integer integer, @NonNull Integer integer2) throws Exception {
return integer+integer2;
}
}).subscribe(RxUtils.getObserver());
以上代码实现了从1到5的累加
onSubscribe
Thread:Thread[main,5,main]
onNext:1
Thread:Thread[main,5,main]
onNext:3
Thread:Thread[main,5,main]
onNext:6
Thread:Thread[main,5,main]
onNext:10
Thread:Thread[main,5,main]
onNext:15
Thread:Thread[main,5,main]
onComplete
Thread:Thread[main,5,main]
Observable.interval(1, TimeUnit.SECONDS).take(12)
.window(3, TimeUnit.SECONDS)
.subscribe(new Observer>() {
public void onSubscribe(@NonNull Disposable d) {
}
public void onNext(@NonNull Observable longObservable) {
longObservable.subscribe(RxUtils.getObserver());
}
public void onError(@NonNull Throwable e) {
}
public void onComplete() {
}
});
try {
Thread.sleep(Integer.MAX_VALUE);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
运行结果
onSubscribe
Thread:Thread[main,5,main]
onNext:0
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:1
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:2
Thread:Thread[RxComputationThreadPool-2,5,main]
onComplete
Thread:Thread[RxComputationThreadPool-2,5,main]
onSubscribe
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:3
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:4
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:5
Thread:Thread[RxComputationThreadPool-2,5,main]
onComplete
Thread:Thread[RxComputationThreadPool-2,5,main]
onSubscribe
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:6
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:7
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:8
Thread:Thread[RxComputationThreadPool-2,5,main]
onComplete
Thread:Thread[RxComputationThreadPool-2,5,main]
onSubscribe
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:9
Thread:Thread[RxComputationThreadPool-2,5,main]
onNext:10
Thread:Thread[RxComputationThreadPool-2,5,main]
onComplete
Thread:Thread[RxComputationThreadPool-1,5,main]
onSubscribe
Thread:Thread[RxComputationThreadPool-1,5,main]
onNext:11
Thread:Thread[RxComputationThreadPool-2,5,main]
onComplete
Thread:Thread[RxComputationThreadPool-2,5,main]
``
注意可能你的结果可能和我的可能不同,由于两个事件源不是不一个线程,不同时间运行会有些时间差
总结:以上是RxJava的变换操作符做常用的莫过于map flatmap这两个,一定要理解两者的区别