日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當前位置: 首頁 > 编程语言 > java >内容正文

java

RxJava之PublishSubject、BehaviorSubject、ReplaySubject和AsyncSubject

發布時間:2023/12/10 java 30 豆豆
生活随笔 收集整理的這篇文章主要介紹了 RxJava之PublishSubject、BehaviorSubject、ReplaySubject和AsyncSubject 小編覺得挺不錯的,現在分享給大家,幫大家做個參考.
public class T2 {/*** subject 是一個神奇的對象,它可以是一個Observable同時也可以是一個Observer:它作為連接這兩個世界的一座橋梁。* 一個主題可以訂閱一個Observable,就像一個觀察者,并且它可以發射新的數據,或者傳遞它接受到的數據,就像一個Observable。* 很明顯,作為一個Observable,觀察者們或者其它主題都可以訂閱它。* 串行化如果你把 Subject 當作一個 Subscriber 使用,注意不要從多個線程中調用它的onNext方法(包括其它的on系列方法),* 這可能導致同時(非順序)調用,這會違反Observable協議,給Subject的結果增加了不確定性。* 要避免此類問題,你可以將 Subject 轉換為一個 SerializedSubject ,類似于這樣:* mySafeSubject = new SerializedSubject( myUnsafeSubject );*/public static void main(String[] args) {T2 t2 = new T2();System.out.println("===================testPublishSubject==========================");t2.testPublishSubject();System.out.println("===================testBehaviorSubject==========================");t2.testBehaviorSubject();System.out.println("===================testReplaySubject==========================");t2.testReplaySubject();System.out.println("===================testAsyncSubject==========================");t2.testAsyncSubject();}/*PublishSubject的觀察者接收到的是后續的消息輸出為:===================testPublishSubject==========================observer1 - A observer1 - B observer1 - C observer2 - C observer1 - D observer2 - D onCompletedonCompleted* */private void testPublishSubject() {Observer<String> observer1 = new Observer<String>() {@Overridepublic void onNext(String t) {System.out.print("observer1 - " + t + "\t");}@Overridepublic void onCompleted() {System.out.println("onCompleted");}@Overridepublic void onError(Throwable e) {System.out.println(e.getMessage());}};Observer<String> observer2 = new Observer<String>() {@Overridepublic void onNext(String t) {System.out.print("observer2 - " + t + "\t");}@Overridepublic void onCompleted() {System.out.println("onCompleted");}@Overridepublic void onError(Throwable e) {System.out.println(e.getMessage());}};PublishSubject<String> publishSubject = PublishSubject.create();publishSubject.subscribe(observer1);publishSubject.onNext("A");publishSubject.onNext("B");publishSubject.subscribe(observer2);publishSubject.onNext("C");publishSubject.onNext("D");publishSubject.onCompleted();System.out.println();}/** BehaviorSubject的觀察者接收到的永遠是最近的消息 和后續的消息* 輸出為===================testBehaviorSubject==========================* default A B C* B C D* onCompleted* error* */private void testBehaviorSubject() {Observer<String> observer = new Observer<String>() {@Overridepublic void onNext(String t) {System.out.print(t + "\t");}@Overridepublic void onCompleted() {System.out.println("onCompleted");}@Overridepublic void onError(Throwable e) {System.out.println(e.getMessage());}};//收到所有消息BehaviorSubject<String> subject1 = BehaviorSubject.create("default");subject1.subscribe(observer);subject1.onNext("A");subject1.onNext("B");subject1.onNext("C");System.out.println();//不能收到default、ABehaviorSubject<String> subject2 = BehaviorSubject.create("default");subject2.onNext("A");subject2.onNext("B");subject2.subscribe(observer);subject2.onNext("C");subject2.onNext("D");System.out.println();//只能收到onCompletedBehaviorSubject<String> subject3 = BehaviorSubject.create("default");subject3.onNext("A");subject3.onNext("B");subject3.onCompleted();subject3.subscribe(observer);System.out.println();// 只能收到errorBehaviorSubject<String> subject4 = BehaviorSubject.create("default");subject4.onNext("A");subject3.onNext("B");subject4.onError(new RuntimeException("error"));subject4.subscribe(observer);System.out.println();}/** ReplaySubject會緩存所有消息,所以觀察者都會收到所有消息* 輸出:===================testReplaySubject==========================* A B A B C C D D onCompleted* onCompleted* */private void testReplaySubject() {Observer<String> observer1 = new Observer<String>() {@Overridepublic void onNext(String t) {System.out.print(t + "\t");}@Overridepublic void onCompleted() {System.out.println("onCompleted");}@Overridepublic void onError(Throwable e) {System.out.println(e.getMessage());}};Observer<String> observer2 = new Observer<String>() {@Overridepublic void onNext(String t) {System.out.print(t + "\t");}@Overridepublic void onCompleted() {System.out.println("onCompleted");}@Overridepublic void onError(Throwable e) {System.out.println(e.getMessage());}};ReplaySubject<String> publishSubject = ReplaySubject.create();publishSubject.subscribe(observer1);publishSubject.onNext("A");publishSubject.onNext("B");publishSubject.subscribe(observer2);publishSubject.onNext("C");publishSubject.onNext("D");publishSubject.onCompleted();System.out.println();}/**當Observable完成時AsyncSubject只會發布最后一條消息給已經訂閱的每一個觀察者,* 如果沒有調用onCompleted則被觀察者不會發送任何消息給觀察者* 輸出===================testAsyncSubject==========================* C onCompleted* */private void testAsyncSubject() {Observer<String> observer = new Observer<String>() {@Overridepublic void onNext(String t) {System.out.print(t + "\t");}@Overridepublic void onCompleted() {System.out.println("onCompleted");}@Overridepublic void onError(Throwable e) {System.out.println(e.getMessage());}};AsyncSubject<String> publishSubject1 = AsyncSubject.create();publishSubject1.subscribe(observer);publishSubject1.onNext("A");publishSubject1.onNext("B");publishSubject1.onNext("C");AsyncSubject<String> publishSubject2 = AsyncSubject.create();publishSubject2.subscribe(observer);publishSubject2.onNext("A");publishSubject2.onNext("B");publishSubject2.onNext("C");publishSubject2.onCompleted();System.out.println();} }

總結

以上是生活随笔為你收集整理的RxJava之PublishSubject、BehaviorSubject、ReplaySubject和AsyncSubject的全部內容,希望文章能夠幫你解決所遇到的問題。

如果覺得生活随笔網站內容還不錯,歡迎將生活随笔推薦給好友。