如需下載源碼,請訪問
https://github.com/fengchuanfang/Rxjava2Tutorial
本篇文章對應(yīng)的Demo類為:com.edward.edward.Rxjava2Demo.Rxjava2_2_Observable
文章原創(chuàng),轉(zhuǎn)載請注明出處:
Rxjava2入門教程二:Observable與Observer響應(yīng)式編程在Rxjava2中的典型實現(xiàn)
運行源碼缭裆,打開app炬灭,點擊Demo圖標(biāo)潜必,可進(jìn)入以下頁面
在RxJava中每币,函數(shù)響應(yīng)式編程具體表現(xiàn)為一個觀察者(Observer)訂閱一個可觀察對象(Observable)先壕,通過創(chuàng)建可觀察對象發(fā)射數(shù)據(jù)流翔脱,經(jīng)過一系列操作符(Operators)加工處理和線程調(diào)度器(Scheduler)在不同線程間的轉(zhuǎn)發(fā)奴拦,最后由觀察者接受并做出響應(yīng)的一個過程
ObservableSource與Observer是RxJava2中最典型的一組觀察者與可觀察對象的組合,其他四組可以看做是這一組的改進(jìn)版或者簡化版届吁。
Observable
抽象類Observable是接口ObservableSource下的一個抽象實現(xiàn)错妖,我們可以通過Observable創(chuàng)建一個可觀察對象發(fā)射數(shù)據(jù)流。
Observable<String> observable = Observable.create(new ObservableOnSubscribe<String>() {
@Override
public void subscribe(ObservableEmitter<String> emitter) throws Exception {
emitter.onNext("Hello World");
emitter.onComplete();
}
});
調(diào)用Observable.create方法疚沐,創(chuàng)建一個可觀察對象暂氯,并通過onNext發(fā)送一條數(shù)據(jù)“Hello World”,然后通過onComplete發(fā)送完成通知亮蛔。
Observer
創(chuàng)建一個觀察者Observer來接受并響應(yīng)可觀察對象發(fā)射的數(shù)據(jù)
Observer<String> observer = new Observer<String>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(String s) {
System.out.println(s);
}
@Override
public void onError(Throwable e) {
e.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("接受完成");
}
};
在onNext方法中接收到可觀察對象發(fā)射的數(shù)據(jù)"Hello World",并做出響應(yīng)——打印到控制臺痴施。
Observer訂閱Observable
observable.subscribe(observer);
通過subscribe方法,使Observer與Observable建立訂閱關(guān)系,Observer與Observable便成為了一個整體晾剖,Observer便可對Observable中的行為作出響應(yīng)锉矢。
Emitter/Observer
通過Observable.create創(chuàng)建可觀察對象時,我們可以發(fā)現(xiàn)具體執(zhí)行發(fā)射動作的是由接口ObservableEmitter的實例化對象完成的齿尽,而ObservableEmitter<T> 繼承自 接口Emitter<T>沽损,查看源碼接口Emitter的具體代碼如下:
public interface Emitter<T> {
//用來發(fā)送數(shù)據(jù),可多次調(diào)用循头,每調(diào)用一次發(fā)送一條數(shù)據(jù)
void onNext(@NonNull T value);
//用來發(fā)送異常通知绵估,只發(fā)送一次,若多次調(diào)用只發(fā)送第一條
void onError(@NonNull Throwable error);
//用來發(fā)送完成通知卡骂,只發(fā)送一次国裳,若多次調(diào)用只發(fā)送第一條
void onComplete();
}
onNext:用來發(fā)送數(shù)據(jù),可多次調(diào)用全跨,每調(diào)用一次發(fā)送一條數(shù)據(jù)
onError:用來發(fā)送異常通知缝左,只發(fā)送一次,若多次調(diào)用只發(fā)送第一條
onComplete:用來發(fā)送完成通知浓若,只發(fā)送一次渺杉,若多次調(diào)用只發(fā)送第一條
onError與onComplete互斥,兩個方法只能調(diào)用一個不能同時調(diào)用
數(shù)據(jù)在發(fā)送時挪钓,出現(xiàn)異呈窃剑可以調(diào)用onError發(fā)送異常通知也可以不調(diào)用,因為其所在的方法subscribe會拋出異常碌上,
若數(shù)據(jù)在全部發(fā)送完之后均正常倚评,可以調(diào)用onComplete發(fā)送一條完成通知
接口Observer中的三個方法(onNext,onError,onComplete)正好與Emitter中的三個方法相對應(yīng),對于Emitter中對應(yīng)方法發(fā)送的數(shù)據(jù)或通知進(jìn)行響應(yīng)馏予。
步驟簡化
去掉中間變量可以對之前的代碼簡化為以下形式:
public void demo2() {
Observable
.create(new ObservableOnSubscribe<String>() {
@Override
public void subscribe(ObservableEmitter<String> emitter) throws Exception {
emitter.onNext("Hello World");
emitter.onComplete();
}
})
.subscribe(new Observer<String>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(String s) {
System.out.println(s);
}
@Override
public void onError(Throwable e) {
e.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("接受完成");
}
});
}
再應(yīng)用Rxjava中強大的操作符天梧,可以將代碼簡化成以下形式:
public void demo3() {
Observable.just("Hello World")
.subscribe(new Consumer<String>() {
@Override
public void accept(@NonNull String s) throws Exception {
System.out.println(s);
}
});
}
再通過λ表達(dá)式,可進(jìn)一步簡化
public void demo3_1() {
Observable.just("Hello World").subscribe(System.out::println);
}
其中霞丧,just操作符可用來發(fā)送單條數(shù)據(jù)呢岗,數(shù)字,字符串蚯妇,數(shù)組,對象暂筝,集合都可以當(dāng)做單條數(shù)據(jù)發(fā)送箩言。
Consumer可以看做是對觀察者Observer功能單一化之后的產(chǎn)物——消費者,上例中的Consumer通過其函數(shù)accept只接收可觀察對象發(fā)射的數(shù)據(jù)焕襟,不接收異常信息或完成信息陨收。
如果想接收異常信息或完成信息可以用下面的代碼:
public void demo4() {
Observable.just("Hello World")
.subscribe(new Consumer<String>() {
@Override
public void accept(@NonNull String s) throws Exception {
System.out.println(s);
}
}, new Consumer<Throwable>() {
@Override
public void accept(@NonNull Throwable throwable) throws Exception {
throwable.printStackTrace();
}
}, new Action() {
@Override
public void run() throws Exception {
System.out.println("接受完成");
}
});
}
第二個參數(shù)Consumer規(guī)定泛型<Throwable>通過函數(shù)accept接收異常信息。
第三個參數(shù)Action也是對觀察者Observer功能單一化之后的產(chǎn)物--行動,通過函數(shù)run接收完成信息务漩,作出響應(yīng)行動拄衰。
發(fā)送數(shù)據(jù)序列
Observable可以發(fā)送單條數(shù)據(jù)也可以發(fā)送數(shù)據(jù)序列
通過最基礎(chǔ)的方法發(fā)送數(shù)據(jù)序列:
public void demo5(final List<String> list) {
Observable.create(new ObservableOnSubscribe<String>() {
@Override
public void subscribe(ObservableEmitter<String> emitter) throws Exception {
try {
for (String str : list) {
emitter.onNext(str);
}
emitter.onComplete();
} catch (Exception e) {
emitter.onError(e);
}
}
}).subscribe(new Observer<String>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(String s) {
System.out.println(s);
}
@Override
public void onError(Throwable e) {
e.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("接受完成");
}
});
}
在subscribe方法中,遍歷集合list中的String元素饵骨,通過emitter.onNext(str)逐一發(fā)送翘悉;發(fā)送完成后通過emitter.onComplete()發(fā)送完成通知;如果發(fā)送過程中遇到異常居触,通過emitter.onError(e)發(fā)送異常信息妖混。
Observer中通過onNext接收emitter發(fā)送的每一條信息并打印到控制臺(emitter發(fā)送幾次,Observer便接收幾次)轮洋,通過onError(Throwable e)接收異常信息制市,onComplete()接收完成信息。
同樣可以通過操作符對其進(jìn)行簡化弊予,如下;
public void demo6(final List<String> list) {
Observable
.fromIterable(list)
.subscribe(new Consumer<String>() {
@Override
public void accept(@NonNull String s) throws Exception {
System.out.println(s);
}
});
}
再用λ表達(dá)式祥楣,進(jìn)一步簡化
public void demo6_1(final List<String> list) {
Observable.fromIterable(list).subscribe(System.out::println);
}
其中fromIterable操作符,可用來將一個可迭代對象中的元素逐一發(fā)送
Disposable
在之前的例子中汉柒,可以看到Observer接口中還有一個方法
public void onSubscribe(Disposable d) {
}
是在觀察者Observer與可觀察對象Observable误褪,建立訂閱關(guān)系后,回調(diào)這個方法竭翠,并且傳過來一個Disposable類型的參數(shù)振坚,可通過Disposable來控制Observer與Observable之間的訂閱。
無論觀察者Observer以何種方式訂閱可觀察對象Observable斋扰,都會生成一個Disposable渡八,如下:
public void demo7(final List<String> list) {
Disposable disposable1 = Observable.just("Hello World")
.subscribe(new Consumer<String>() {
@Override
public void accept(@NonNull String s) throws Exception {
System.out.println(s);
}
});
Disposable disposable2 = Observable
.fromIterable(list)
.subscribe(new Consumer<String>() {
@Override
public void accept(@NonNull String s) throws Exception {
System.out.println(s);
}
});
}
查看Disposable接口的源碼,如下:
public interface Disposable {
void dispose();
boolean isDisposed();
}
其中isDisposed()方法用來判斷當(dāng)前訂閱是否失效传货,dispose()方法用來取消當(dāng)前訂閱屎鳍。
只有當(dāng)觀察者Observer與可觀察對象Observable之間建立訂閱關(guān)系,并且訂閱關(guān)系有效時问裕,Observer才能對Observable進(jìn)行響應(yīng)逮壁。如果Observer在響應(yīng)Observable的過程中,訂閱關(guān)系被取消粮宛,則Observer無法對取消訂閱關(guān)系之后Observable的行為作出響應(yīng)窥淆。
運行下面的代碼,當(dāng)Observable接收到第5條數(shù)據(jù)時巍杈,取消訂閱關(guān)系忧饭。
public void demo8() {
Observable.create(new ObservableOnSubscribe<Integer>() {
@Override
public void subscribe(ObservableEmitter<Integer> emitter) throws Exception {
for (int i = 0; i < 10; i++) {
System.out.println("發(fā)送" + i);
emitter.onNext(i);
}
emitter.onComplete();
}
}).subscribe(new Observer<Integer>() {
private Disposable disposable;
@Override
public void onSubscribe(Disposable d) {
disposable = d;
}
@Override
public void onNext(Integer integer) {
System.out.println("接收" + integer);
if (integer > 4) disposable.dispose();
}
@Override
public void onError(Throwable e) {
e.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("數(shù)據(jù)接受完成");
}
});
}
控制臺日志如下:
I/System.out: 發(fā)送0
I/System.out: 接收0
I/System.out: 發(fā)送1
I/System.out: 接收1
I/System.out: 發(fā)送2
I/System.out: 接收2
I/System.out: 發(fā)送3
I/System.out: 接收3
I/System.out: 發(fā)送4
I/System.out: 接收4
I/System.out: 發(fā)送5
I/System.out: 接收5
I/System.out: 發(fā)送6
I/System.out: 發(fā)送7
I/System.out: 發(fā)送8
I/System.out: 發(fā)送9
可以發(fā)現(xiàn)取消訂閱關(guān)系之前,Observable發(fā)送一條數(shù)據(jù)筷畦,Observe便接收一條词裤,但是取消訂閱關(guān)系之后刺洒,Observe將不再接收Observable發(fā)送的數(shù)據(jù)。
上一篇:Rxjava2入門教程一:函數(shù)響應(yīng)式編程及概述
下一篇吼砂;Rxjava2入門教程三:Operators操作符