The switchMap operator must not be confused with concatMap, as it looks very similar at the first glance. However, it works by cancelling the previous inner subscriber whenever the outer publisher emits an item.
RxJava etiketine sahip kayıtlar gösteriliyor. Tüm kayıtları göster
RxJava etiketine sahip kayıtlar gösteriliyor. Tüm kayıtları göster
21 Kasım 2023 Salı
RxJava Observable.switchMap metodu - Yani flatMapLatest
Giriş
RxJava Observable.flatMapSequential metodu - Sonuç Sırasını Korur
GirişÇıktı şöyle
Açıklaması şöyle
This operator eagerly subscribes to its inner publishers like flatMap, but queues up elements from later inner publishers to match the actual natural ordering and thereby prevents interleaving like concatMap.
Örnek
Şöyle yaparız. Burada flatMapSequential ile çalıştırma paralel olsa bile çıktı sırası korunuyor
@Test void test_flatMap() { Flux.just(1, 2, 3) .flatMap(this::doSomethingAsync) //.flatMapSequential(this::doSomethingAsync) //.concatMap(this::doSomethingAsync) .doOnNext(n -> log.info("Done {}", n)) .blockLast(); } private Mono<Integer> doSomethingAsync(Integer number) { //add some delay for the second item... return number == 2 ? Mono.just(number).doOnNext(n -> log.info("Executing {}", n)) .delayElement(Duration.ofSeconds(1)) : Mono.just(number).doOnNext(n -> log.info("Executing {}", n)); }
// flatMap does not preserve original ordering, and has subscribed to all three elements // eagerly. Also, notice that element 3 has proceeded before element 2. 2022-04-22 19:38:49,164 INFO main - Executing 1 2022-04-22 19:38:49,168 INFO main - Done 1 2022-04-22 19:38:49,198 INFO main - Executing 2 2022-04-22 19:38:49,200 INFO main - Executing 3 2022-04-22 19:38:49,200 INFO main - Done 3 2022-04-22 19:38:50,200 INFO parallel-1 - Done 2 // flatMapSequential has subscribed to all three elements eagerly like flatMap // but preserves the order by queuing elements received out of order. 2022-04-22 19:53:40,229 INFO main - Executing 1 2022-04-22 19:53:40,232 INFO main - Done 1 2022-04-22 19:53:40,261 INFO main - Executing 2 2022-04-22 19:53:40,263 INFO main - Executing 3 2022-04-22 19:53:41,263 INFO parallel-1 - Done 2 2022-04-22 19:53:41,264 INFO parallel-1 - Done 3 //concatMap naturally preserves the same order as the source elements. 2022-04-22 19:59:31,817 INFO main - Executing 1 2022-04-22 19:59:31,820 INFO main - Done 1 2022-04-22 19:59:31,853 INFO main - Executing 2 2022-04-22 19:59:32,857 INFO parallel-1 - Done 2 2022-04-22 19:59:32,857 INFO parallel-1 - Executing 3 2022-04-22 19:59:32,857 INFO parallel-1 - Done 3
RxJava Observable.concatMap metodu
Giriş
Açıklaması şöyle. Yani concatMap dışarıdaki Observable ile gelen girdinin sırasını değiştirmez. Ayrıca bir nesnenin işi bitmeden bir sonrakine geçmez.
The concatMap operator is actually quite similar to the flatMap, except that the operator waits for its inner publishers to complete before subscribing to the next one.
Örnek
Şöyle yaparız. Burada sıra korunuyor
@Testpublic void flatMapVsConcatMap() throws Exception {System.out.println("******** Using flatMap() *********");Observable.range(1, 15).flatMap(item -> Observable.just(item).delay(1, TimeUnit.MILLISECONDS)).subscribe(x -> System.out.print(x + " "));Thread.sleep(100);System.out.println("\n******** Using concatMap() *********");Observable.range(1, 15).concatMap(item -> Observable.just(item).delay(1, TimeUnit.MILLISECONDS)).subscribe(x -> System.out.print(x + " "));Thread.sleep(100);}
Çıktı şöyle
******** Using flatMap() ********* 1 2 3 4 5 6 7 9 8 11 13 15 10 12 14 ******** Using concatMap() ********* 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
Örnek
Şöyle yaparız. Burada sıra korunuyor
@Test void test_flatMap() { Flux.just(1, 2, 3) .flatMap(this::doSomethingAsync) //.flatMapSequential(this::doSomethingAsync) //.concatMap(this::doSomethingAsync) .doOnNext(n -> log.info("Done {}", n)) .blockLast(); } private Mono<Integer> doSomethingAsync(Integer number) { //add some delay for the second item... return number == 2 ? Mono.just(number).doOnNext(n -> log.info("Executing {}", n)) .delayElement(Duration.ofSeconds(1)) : Mono.just(number).doOnNext(n -> log.info("Executing {}", n)); }
// flatMap does not preserve original ordering, and has subscribed to all three elements // eagerly. Also, notice that element 3 has proceeded before element 2. 2022-04-22 19:38:49,164 INFO main - Executing 1 2022-04-22 19:38:49,168 INFO main - Done 1 2022-04-22 19:38:49,198 INFO main - Executing 2 2022-04-22 19:38:49,200 INFO main - Executing 3 2022-04-22 19:38:49,200 INFO main - Done 3 2022-04-22 19:38:50,200 INFO parallel-1 - Done 2 // flatMapSequential has subscribed to all three elements eagerly like flatMap // but preserves the order by queuing elements received out of order. 2022-04-22 19:53:40,229 INFO main - Executing 1 2022-04-22 19:53:40,232 INFO main - Done 1 2022-04-22 19:53:40,261 INFO main - Executing 2 2022-04-22 19:53:40,263 INFO main - Executing 3 2022-04-22 19:53:41,263 INFO parallel-1 - Done 2 2022-04-22 19:53:41,264 INFO parallel-1 - Done 3 //concatMap naturally preserves the same order as the source elements. 2022-04-22 19:59:31,817 INFO main - Executing 1 2022-04-22 19:59:31,820 INFO main - Done 1 2022-04-22 19:59:31,853 INFO main - Executing 2 2022-04-22 19:59:32,857 INFO parallel-1 - Done 2 2022-04-22 19:59:32,857 INFO parallel-1 - Executing 3 2022-04-22 19:59:32,857 INFO parallel-1 - Done 3
4 Haziran 2021 Cuma
RxJava Observable.create metodu - Vertx İle Birlikte Kullanılabilir
Giriş
ObservableOnSubscribe nesnesi alır. Bu bir emitter'dır. İş bitince emitter nesnesinin onNext() veya onError() metodu çağrılır. İşin bittiğini belirtmek için de onComplete() çağrılır.
Örnek
Şöyle yaparız.
Observable.create(new ObservableOnSubscribe<Integer>() {
@Override
public void subscribe(ObservableEmitter<Integer> observableEmitter) throws Exception {
if (!observableEmitter.isDisposed())
observableEmitter.onComplete();
}
}).subscribe (...);
Örnek
Şöyle yaparız
Vertx vertx = Vertx.vertx();
WebClient webClient = WebClient.create(vertx);
Observable<Object> google = hitURL("www.google.com", webClient);
Observable<Object> yahoo = hitURL("www.yahoo.com", webClient);
for (int i = 0; i < 100; i++) {
google.repeat(100).subscribe(timeTaken -> {
if ((Long) timeTaken > 10000) {
System.out.println(timeTaken);
}
}, error -> {System.out.println(error.getMessage());});
yahoo.repeat(100).subscribe(timeTaken -> {
if ((Long) timeTaken > 10000) {
System.out.println(timeTaken);
}
}, error -> {System.out.println(error.getMessage());});
}
}
public static Observable<Object> hitURL(String url, WebClient webClient) {
return Observable.create(emitter -> {
Long l1 = System.currentTimeMillis();
webClient.get(80, url, "").send(ar -> {
if (ar.succeeded()) {
Long elapsedTime = (System.currentTimeMillis() - l1);
emitter.onNext(elapsedTime);
} else {
emitter.onError(ar.cause());
}
emitter.onComplete();
});
});
}
28 Mayıs 2021 Cuma
RxJava AsyncSubject Sınıfı
Giriş
Açıklaması şöyle
ÖrnekAsyncSubject emits only the last value of the Observable and this only happens after the Observable completes.
Elimizde şöyle bir kod olsun
AsyncSubject<Integer> pSubject = AsyncSubject.create();pSubject.onNext(0);pSubject.subscribe(it -> System.out.println("Observer 1 onNext: " + it),(Throwable error) -> { }, () -> System.out.println("Observer 1 onComplete"),on1 -> System.out.println("Observer 1 onSubscribe"));pSubject.onNext(1);pSubject.onNext(2);pSubject.subscribe(it -> System.out.println("Observer 2 onNext: " + it),(Throwable error) -> { }, () -> System.out.println("Observer 2 onComplete"),on1 -> System.out.println("Observer 2 onSubscribe"));pSubject.onNext(3);pSubject.onNext(4);/* This is very important in AsyncSubject */pSubject.onComplete();
Çıktı olarak şunu alırız
Observer 1 onSubscribeObserver 2 onSubscribeObserver 1 onNext: 4Observer 1 onCompleteObserver 2 onNext: 4Observer 2 onComplete
RxJava UnicastSubject Sınıfı
Giriş
Açıklaması şöyle
ÖrnekUnicastSubject allows only a single subscriber and it emits all the items regardless of the time of subscription.
Şöyle yaparız
Observable<Integer> observable = Observable.range(1, 5).subscribeOn(Schedulers.io());UnicastSubject<Integer> pSubject = UnicastSubject.create();observable.subscribe(pSubject);pSubject.subscribe(it -> System.out.println("onNext: " + it));
Çıktı olarak şunu alırız
onNext: 1onNext: 2onNext: 3onNext: 4onNext: 5
RxJava BehaviorSubject Sınıfı - En Son Yayınlanan Nesneyi Yeni Aboneye Gönderir
Giriş
Açıklaması şöyle
ÖrnekBehaviorSubject emits the most recent item at the time of their subscription and all items after that.
Elimizde şöyle bir kod olsun
Çıktı olarak şunu alırız. 2 numaralı katılımcı abone olduğunda 2 değeri zaten yayınlanmıştı, ancak yine de duyabilir.BehaviorSubject<Integer> pSubject = BehaviorSubject.create();pSubject.onNext(0);pSubject.subscribe(it -> System.out.println("Observer 1 onNext: " + it),(Throwable error) -> { }, () -> {},on1 -> System.out.println("Observer 1 onSubscribe"));pSubject.onNext(1);pSubject.onNext(2);pSubject.subscribe(it -> System.out.println("Observer 2 onNext: " + it),(Throwable error) -> { }, () -> {},on1 -> System.out.println("Observer 2 onSubscribe"));pSubject.onNext(3);pSubject.onNext(4);
Observer 1 onSubscribeObserver 1 onNext: 0Observer 1 onNext: 1Observer 1 onNext: 2Observer 2 onSubscribeObserver 2 onNext: 2Observer 1 onNext: 3Observer 2 onNext: 3Observer 1 onNext: 4Observer 2 onNext: 4
RxJava ReplaySubject Sınıfı
Giriş
Açıklaması şöyle
ÖrnekReplaySubject emits all the items of the Observable, regardless of when the subscriber subscribes.
Elimizde şöyle bir kod olsun
Çıktı olarak şunu alırız. Aslında bir anlamda Cold Observale'ın tüm çıktısı kaydediliyor ve kaydedilen şey tekrar en baştan oynatılıyor.ReplaySubject<Integer> pSubject = ReplaySubject.create();pSubject.onNext(0);pSubject.subscribe(it -> System.out.println("Observer 1 onNext: " + it),(Throwable error) -> { }, () -> {},on1 -> System.out.println("Observer 1 onSubscribe"));pSubject.onNext(1);pSubject.onNext(2);pSubject.subscribe(it -> System.out.println("Observer 2 onNext: " + it),(Throwable error) -> { }, () -> {},on1 -> System.out.println("Observer 2 onSubscribe"));pSubject.onNext(3);pSubject.onNext(4);
Yani ReplaySubject pahalı işlemleri bir şekilde cache'lemek için kullanılabilirObserver 1 onSubscribeObserver 1 onNext: 0Observer 1 onNext: 1Observer 1 onNext: 2Observer 2 onSubscribeObserver 2 onNext: 0Observer 2 onNext: 1Observer 2 onNext: 2Observer 1 onNext: 3Observer 2 onNext: 3Observer 1 onNext: 4Observer 2 onNext: 4
Örnek
Elimizde şöyle bir kod olsun
Çıktı olarak şunu alırızObservable<Integer> observable = Observable.range(1, 5).subscribeOn(Schedulers.io());//Record everytingReplaySubject<Integer> subject = ReplaySubject.create();observable.subscribe(subject);//Replay everythingsubject.subscribe(s -> System.out.println("subscriber one: " + s));//Replay everythingsubject.subscribe(s -> System.out.println("subscriber two: " + s));
subscriber one: 1subscriber one: 2subscriber one: 3subscriber one: 4subscriber one: 5subscriber two: 1subscriber two: 2subscriber two: 3subscriber two: 4subscriber two: 5
27 Mayıs 2021 Perşembe
RxJava Subjects
Giriş
Açıklaması şöyle
Peki Subject neden Observable'a abone olup tekrar Observable hale getiriyor. Buradaki amaçlardan bir tanesi şuA Subject extends an Observable and implements Observer at the same time. It acts as an Observable to clients and registers to multiple events taking place in the app. It acts as an Observer by broadcasting the event to multiple subscribers.
Bir çok Subject gerçekleştirimi var. Bunlar şöyleSubjects are considered as HOT Observables.A HOT Observable, such as Subjects, emits items only once regardless of number of subscribers and its subscribers receive items only from the point of their subscription. Subjects convert cold observable into hot observable.
17 Mayıs 2021 Pazartesi
RxJava Observable.debounce metodu
Giriş
Eğer event'ler çok hızlı geliyorsa, seyreltmek için kullanılır.
Şöyle yaparız
inputObservable.debounce(1, TimeUnit.SECONDS)
.subscribe(new Action1<String>() {
@Override
public void call(String s) {
...
}
});RxJava Observable.zip metodu
Giriş
İşleri paralel çalıştırır.
Örnek
Şöyle yaparız
Observable.zip(api.getUserDetails2(userId), api.getUserPhoto2(userId),
(details, photo) -> Pair.of(details, photo))
.subscribe(p -> {
// Do your task.
});19 Mart 2021 Cuma
RxJava Backpressure - Geri tepme
Giriş
RxJava 3 Sınıfları
Açıklaması şöyle. Yani hızlı bir Producer ve yavaş bir Consumer varsa, Producer bu durum karşısında ezilmez.
... considering a fast data producer and a slow data consumer, backpressure is the mechanism that 'pushes back' on the producer not to be overwhelmed by data.
Şeklen şöyle
Açıklaması şöyle. Producer'ın ezilmemesinin sebebi, tüketen tarafın Subscription.request() metodu ile veriyi çekmesi.
A Subscriber MUST signal demand via Subscription.request(long n) to receive onNext signals.The intent of this rule is to establish that it is the responsibility of the Subscriber to decide when and how many elements it is able and willing to receive. To avoid signal reordering caused by reentrant Subscription methods, it is strongly RECOMMENDED for synchronous Subscriber implementations to invoke Subscription methods at the very end of any signal processing. It is RECOMMENDED that Subscribers request the upper limit of what they are able to process, as requesting only one element at a time results in an inherently inefficient "stop-and-wait" protocol.-- Reactive Streams specifications for the JVM
Açıklaması şöyle
To cope with that, RxJava offers two main strategies to handle 'overproduced' items:1. Store items in a buffer2. Drop items
Sınıflar şöyle
Flowable
Açıklaması şöyle
A flow of 0..N items. It supports Reactive-Streams and backpressure.Observable
20 Temmuz 2020 Pazartesi
RxJava Single Sınıfı - Bir Observable
Giriş
Single nesnesi yaratıldıktan sonra subscribeOn metodu çağrılır.
create metodu
Şöyle yaparız.
Örnek
Şöyle yaparız.
Single.error() yerine direkt exception da fırlatılabilir.
Örnek
Şöyle yaparız
Şöyle yaparız.
Örnek
Şöyle yaparız
Şöyle yaparız.
Örnek
Şöyle yaparız
Single nesnesi yaratıldıktan sonra subscribeOn metodu çağrılır.
create metodu
Şöyle yaparız.
Single<String> saveBookToRepository(AddBookRequest addBookRequest) {
return Single.create(singleSubscriber -> {
if (notFound) {
singleSubscriber.onError(new EntityNotFoundException());
else {
singleSubscriber.onSuccess(addedBookId);
}
});
}
defer metoduÖrnek
Şöyle yaparız.
public Single<List<String>> getHealthChecks(JsonArray endpoints) {
return Single.defer(() -> {
List<Single<String>> healthChecks = endpoints
.stream()
.map(endpoint -> getHealthStatus(client, endpoint.toString()))
.collect(Collectors.toList());
return consumeHealthChecks(healthChecks);
});
}
private Single<List<String>> consumeHealthChecks(List<Single<String>> healthChecks) {
return Single.merge(healthChecks)
.timeout(1500, TimeUnit.MILLISECONDS)
.toList();
}
error metoduSingle.error() yerine direkt exception da fırlatılabilir.
Örnek
Şöyle yaparız
public static Single<Integer> externalMethod(int x){
int result = 0;
try{
/* Some database / time consuming logic */
result = x % 0;
}
catch (Exception e){
throw new CustomException(e.getMessage()); // --> APPROACH 1
OR
return Single.error(new CustomException(e.getMessage())); // --> APPROACH 2
}
return Single.just(result);
}
fromCallable metoduŞöyle yaparız.
Single.fromCallable(() -> {
...
return true;
}
)
.observeOn(Schedulers.io())
.subscribeOn(AndroidSchedulers.mainThread())
.subscribe(ignored -> {
// finished
}, error -> {
// error
});
subscribe metoduÖrnek
Şöyle yaparız
Single.just(1)
.flatMap(x -> externalMethod(x))
.subscribe(
s -> System.out.println("Success : " + s),
e -> System.out.println("Error : " + e)
);
public static Single<Integer> externalMethod(int x){
return Single.just(...);
}
subscribeOn metoduŞöyle yaparız.
return bookService.addBook(toAddBookRequest(addBookWebRequest)) //Single nesnesi döner
.subscribeOn(Schedulers.io())
.map(s -> ResponseEntity.created(...).body(...));
zip metodu
Örnek
Şöyle yaparız
public static <A, B> Single<B> zipOver(Single<Function<A, B>> applicativeFunctor,
Single<A> applicativeValue) {
return Single.zip(
applicativeFunctor,
applicativeValue,
(Function<A, B> f, A a) -> f.apply(a));
}
Kaydol:
Kayıtlar (Atom)
