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ş
Açıklaması şöyle. Bu metodu eski ismi flatMapLatest, en son sonuç nesnesini dışarı verir
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 Observable.flatMapSequential metodu - Sonuç Sırasını Korur

Giriş
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)); }
Çıktı şöyle
// 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
@Test
public 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)); }
Çıktı şöyle
// 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
AsyncSubject emits only the last value of the Observable and this only happens after the Observable completes.
Örnek 
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 onSubscribe
Observer 2 onSubscribe
Observer 1 onNext: 4
Observer 1 onComplete
Observer 2 onNext: 4
Observer 2 onComplete

RxJava UnicastSubject Sınıfı

Giriş
Açıklaması şöyle
UnicastSubject allows only a single subscriber and it emits all the items regardless of the time of subscription.
Örnek
Şö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: 1
onNext: 2
onNext: 3
onNext: 4
onNext: 5

RxJava BehaviorSubject Sınıfı - En Son Yayınlanan Nesneyi Yeni Aboneye Gönderir

Giriş
Açıklaması şöyle
BehaviorSubject emits the most recent item at the time of their subscription and all items after that. 
Örnek 
Elimizde şöyle bir kod olsun
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);
Çı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.
Observer 1 onSubscribe
Observer 1 onNext: 0
Observer 1 onNext: 1
Observer 1 onNext: 2
Observer 2 onSubscribe
Observer 2 onNext: 2
Observer 1 onNext: 3
Observer 2 onNext: 3
Observer 1 onNext: 4
Observer 2 onNext: 4

RxJava ReplaySubject Sınıfı

Giriş
Açıklaması şöyle
ReplaySubject emits all the items of the Observable, regardless of when the subscriber subscribes.
Örnek
Elimizde şöyle bir kod olsun
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);
Çı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.
Observer 1 onSubscribe
Observer 1 onNext: 0
Observer 1 onNext: 1
Observer 1 onNext: 2
Observer 2 onSubscribe
Observer 2 onNext: 0
Observer 2 onNext: 1
Observer 2 onNext: 2
Observer 1 onNext: 3
Observer 2 onNext: 3
Observer 1 onNext: 4
Observer 2 onNext: 4
Yani ReplaySubject pahalı işlemleri bir şekilde cache'lemek için kullanılabilir
Örnek
Elimizde şöyle bir kod olsun
Observable<Integer> observable = Observable.range(1, 5)
                .subscribeOn(Schedulers.io());

//Record everyting
ReplaySubject<Integer> subject = ReplaySubject.create();
observable.subscribe(subject);

//Replay everything
subject.subscribe(s -> System.out.println("subscriber one: " + s));

//Replay everything
subject.subscribe(s -> System.out.println("subscriber two: " + s));
Çıktı olarak şunu alırız
subscriber one: 1
subscriber one: 2
subscriber one: 3
subscriber one: 4
subscriber one: 5
subscriber two: 1
subscriber two: 2
subscriber two: 3
subscriber two: 4
subscriber two: 5

27 Mayıs 2021 Perşembe

RxJava Subjects

Giriş
Açıklaması şöyle
A 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.
Peki Subject neden Observable'a abone olup tekrar Observable hale getiriyor. Buradaki amaçlardan bir tanesi şu
Subjects 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.
Bir çok Subject gerçekleştirimi var. Bunlar şöyle

17 Mayıs 2021 Pazartesi

RxJava Observable.debounce metodu

Giriş
Eğer event'ler çok hızlı geliyorsa, seyreltmek için kullanılır.

Örnek
Şö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ş
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 buffer
2. Drop items

RxJava 3 Sınıfları
Sınıflar şöyle

Flowable
Açıklaması şöyle
A flow of 0..N items. It supports Reactive-Streams and backpressure.
Observable
Açıklaması şöyle
A flow of 0..N items. It doesn't support backpressure.
Single
Açıklaması şöyle
A flow of exactly: 1 item, or an error.
Maybe
Açıklaması şöyle
A flow with either: no items, exactly one item, or an error.
Completable
Açıklaması şöyle
A flow with no item but either: a completion, or an error signal.

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.
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 metodu
Single.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));
}