Skip to content

7. RxJava kütüphanesi

RxJava kütüphanesi şu konsepte dayanır: T türündeki Observable<T> öğelerinden oluşan bir akış, bir veya daha fazla Subscriber<T> abonesi (aboneler, gözlemciler, tüketiciler) tarafından gözlemlenir. RxJava kütüphanesi, geliştiricinin bu iş parçacıklarının yaşam döngüsünü yönetmekle veya iş parçacıkları arasında veri paylaşımı ve genel bir görevi yürütmek için iş parçacıklarının senkronizasyonu gibi doğal olarak zor sorunlarlabu iş parçacıklarının yaşam döngüsünü yönetmekle veya iş parçacıkları arasında veri paylaşımı ve genel bir görevi yürütmek için iş parçacıklarının senkronizasyonu gibi doğal olarak zor sorunlarla uğraşmasına gerek kalmadan. Böylece asenkron programlamayı kolaylaştırır.

Bir Observable<T> akışı, T türündeki elemanları üretir ve bu elemanlar üretildikçe gözlemlenebilir hale gelir. Gözlemci ve gözlemlenebilir (konuşma dilinde Observable<T> türünü ifade etmek için kullanılır) aynı iş parçacığında bulunuyorsa, gözlemlenebilir, gözlemci i öğesini tükettiğinde ancak (i+1) öğesini üretebilir. Bu mimarinin faydalı olduğu durumlar çok azdır. Gözlemci ve gözlemlenebilir aynı iş parçacığında bulunmuyorsa, gözlemlenebilir ve gözlemcisi birbirinden bağımsız davranır: gözlemlenebilir kendi hızında üretir, gözlemci ise kendi hızında tüketir. Kütüphanenin önemi işte burada yatmaktadır. Şimdiye kadar hep tek bir gözlemciden bahsettik. Oysa gerçekte, bir gözlemlenebilirin istediği kadar gözlemcisi olabilir.

RxJava kütüphanesi, keşif bölümünün 2. paragrafında ele alınan ve burada tekrar hatırlatılan mimariye özellikle uygundur:

Image

  • [1]'te, bir hizmet katmanı hizmetler sunar; bu hizmetlerin bazılarının elde edilmesi uzun sürer (örneğin ağ istekleri);
  • bu hizmet katmanı, bir grafik kullanıcı arayüzü ([1]; Swing, Android, JavaFx) tarafından çağrılır. Hizmet katmanı, onu kullanan [swing] yöntemiyle aynı iş parçacığında çalıştırılırsa, hizmetin sonucunu beklerken grafik arayüz donar (tepki vermez);
  • [2]'te, RxJava ile uygulanan ince bir uyarlama katmanı, grafik katmanına aynı hizmetin asenkron bir uygulamasını sunmayı sağlar: bu hizmet, onu çağıran grafik katmanı yönteminden farklı bir iş parçacığında çalıştırılabilir. Bu durumda, [3] grafik arayüzü duyarlı kalır: kullanıcı arayüzle etkileşime devam edebilir; örneğin, ilk isteğin yanı sıra paralel olarak yeni bir ağ isteği başlatabilir ve en önemlisi, kullanıcıya çok uzun süren işlemleri iptal etme olanağı sunulabilir; bu, grafik arayüzün donmuş olması durumunda imkansızdır;
  • [4] çağrısı senkron iken, [5-6] çağrısı asenkron;

Bu mimaride, [2] katmanı, [3] grafik katmanının yöntemlerinin abone olabileceği Observable<T> türlerini döndüren hizmetler sunar. Böylece, [2] katmanındaki bir hizmet sonuçlarını tek tek sunar ve [3] katmanı, örneğin grafik arayüzün bir veya daha fazla bileşenini güncelleyerek her bir sonuca tepki verebilir.

Observable<T> sınıfı onlarca yönteme sahiptir. Bu, kütüphanenin zorluklarından biridir: çok zengin bir yapıya sahiptir ve tüm olanaklarını kavramak zordur. Bunlardan bazılarını tanıtacağız. Diğer yöntemleri zamanla öğreneceksiniz.

7.1. Gözlemlenebilir nesneler oluşturma ve bunlara abone olma

7.1.1. Örnek-01: [Observable.from] yöntemi

  

Aşağıdaki kodu ele alalım:


package dvp.rxjava.observables;

import rx.Observable;
import rx.functions.Action0;
import rx.functions.Action1;

import java.util.Arrays;

public class Exemple01 {
  public static void main(String[] args) {
    // tamsayı gözlemlenebilirleri
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
    });
  }
}
  • 12. satır: Bir tamsayı listesinden Observable<Integer> türü oluşturulur.

Observable<T> sınıfı, T türündeki öğelerden oluşan bir akıştır; bu öğeler üretildikçe, tercihen asenkron olarak (ancak zorunlu değildir) gözlemlenebilir. Tanımı şöyledir:

 

Daha önce de belirtildiği gibi, Observable<T> sınıfı onlarca yönteme sahiptir. Bunlardan bazıları, 5. paragrafta incelenen Stream<T> sınıfındakilere benzemektedir. RxJava belgeleri, bu yöntemlerin işleyişini gösteren [2] "marble diyagramları" içerir:

  • 3. satır, gözlemlenebilirin zaman içindeki emisyonlarını gösterir;
  • [4] yöntemi, gözlemlenebilir tarafından yayılan parçacıklara uygulanır. Bu yöntem genellikle yeni bir gözlemlenebilir üretir;
  • 5. satır, elde edilen yeni gözlemleneni gösterir;

[Observable.from] yönteminin imzası şöyledir:

 

Statik [Observable.from] yöntemi, T türündeki öğelerden oluşan bir koleksiyondan bir Observable<T> oluşturmaya olanak tanır. Bu, gözlemlenebilirlerle çalışmaya başlamak için çok basit bir yoldur. Şu satır:


    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));

satırı üç öğe yayınlayacaktır. Bu öğeleri hemen yayınlamaz. Her gözlemci kendini bildirdiğinde öğeleri tamamen yayınlayacaktır. Buna "soğuk gözlemlenebilir" denir. Gözlemlenebilir, her yeni abone için öğelerini yeniden yayınlar.

Yukarıdaki komutu, gözlemlenebilirin yapılandırma eylemi olarak düşünebiliriz. Gözlemlenebilir bir kez yapılandırılır ve n abone olduğunda n kez yürütülür.

Nasıl abone olunur?

Bunu yapmanın bir yolu, burada aşağıdaki gibi tanımlanan [Observable.subscribe] yöntemini kullanmaktır:

 
  • Yöntemin ilk parametresi [Action1<T> onNext] (bkz. paragraf 6.2), gözlemlenebilir birim yeni bir T öğesi yayınladığında yürütülecek yöntemdir;
  • yöntemin ikinci parametresi [Action1<Throwable> onError], gözlemlenebilir bir istisna attığında çalıştırılacak yöntemdir;
  • yöntemin üçüncü parametresi [Action0 onComplete] (bkz. paragraf 6.1), gözlemlenebilir bir istisna yayınladığında yürütülecek yöntemdir;
  • yöntem, [Subscription] türünde bir değer döndürür;

[Subscription] türü, gözlemlenebilir nesneye bir aboneliği temsil eder. Tanımı şöyledir:

 

Bu [1] arayüzünün önemi, bir aboneliği iptal etmeyi sağlayan [2] yönteminde yatmaktadır.

Örneğimizde, gözlemlenebilir öğeye abonelik kodu şöyledir:


    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
});
  • 1. satır: [Subscription] türündeki sonuç göz ardı edilir;
  • satır 1-15: üç parametre de anonim sınıfların örnekleridir. Ayrıca lambda ifadeleri de kullanacağız. Anonim sınıfların avantajı, bu sınıfların tek yönteminin beklediği veri türlerini açıkça görebilmemizdir;
  • satır 2-5: [Action1<Integer>] türündeki ilk parametrenin uygulaması;
  • satır 6-10: [Action1<Throwable>] türündeki ikinci parametrenin uygulaması;
  • satır 11-15: [Action0] türündeki üçüncü parametrenin uygulanması;

Kodun tamamı şu şekildedir:


package dvp.rxjava.observables;

import rx.Observable;
import rx.functions.Action0;
import rx.functions.Action1;

import java.util.Arrays;

public class Exemple01 {
  public static void main(String[] args) {
    // tamsayı gözlemlenebilirleri
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // abonelik
    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
    });
  }
}
  1. satırdaki gözlemlenebilir, 14. satırda [subscribe] yöntemi çağrılır çağrılmaz 3 öğesini yayınlamaya başlar. Bu andan itibaren:
  • her eleman yayınlandığında, 15-18. satırlar yürütülür.
  • 3 öğe bittiğinde, 24-29. satırlar yürütülür;
  • 19-24. satırlar hiçbir zaman çalıştırılmayacaktır çünkü gözlemlenebilir burada bir istisna yaymaz;

Varsayılan olarak, gözlemlenebilir ve gözlemci aynı iş parçacığında çalışır. Ana iş parçacığından (burada main yönteminin iş parçacığı) farklı bir iş parçacığında çalışan birkaç önceden tanımlanmış gözlemlenebilir vardır, ancak çoğu için durum böyle değildir. Dolayısıyla burada her şey [main] yönteminin iş parçacığında gerçekleşir:

  • gözlemlenebilir, 1 öğesini yayınlar;
  • 15-18. satırlar çalıştırılır ve bu eleman görüntülenir;
  • gözlemlenebilir, 2 numaralı öğeyi yayınlar;
  • 15-18. satırlar çalışır ve bu öğeyi görüntüler;
  • gözlemlenebilir, 3 numaralı öğeyi yayınlar;
  • 15-18. satırlar çalışır ve bu öğeyi görüntüler;
  • gözlemlenebilir, [completed] bildirimini yayınlar;
  • 24-29. satırlar çalışır;

Elde edilen sonuçlar şunu göstermektedir:

1
2
3
4
next : 1
next : 2
next : 3
completed

[Exemple02] sınıfı, bu kez [Observable.subscribe] yönteminin parametreleri olarak lambda işlevleri kullanarak [Exemple01] sınıfını devralır:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;

public class Exemple02 {
  public static void main(String[] args) {
    // tamsayı gözlemlenebilirleri
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // abonelik
    obs1.subscribe(
      (integer) -> System.out.printf("next : %s%n", integer),
      (th) -> System.out.println(th),
      () -> System.out.println("completed"));
  }
}

7.1.2. Örnek-03: Observer sınıfı

  

Bir gözlemlenene abone olmayı sağlayan [Observable.subscribe] yöntemi, aşağıdakiler de dahil olmak üzere çeşitli sürümleri vardır:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;
import rx.Observer;

public class Exemple03 {
    public static void main(String[] args) {
        // tamsayı gözlemlenebilirleri
        Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
        // abonelik
        obs1.subscribe(new Observer<Integer>() {
            @Override
            public void onCompleted() {
                System.out.println("completed");
            }

            @Override
            public void onError(Throwable th) {
                System.out.printf("throwable %s", th);
            }

            @Override
            public void onNext(Integer integer) {
                System.out.printf("next : %s%n", integer);
            }
        });
    };
}
  1. satırda, [subscribe] yöntemine üç parametre aktarmak yerine, aşağıdaki [Observer] türü aktarılır:
 

[Observer] türü, üç yönteme sahip bir arayüzdür:

  • [onNext(T t)], gözlemlenebilir nesne bir t öğesi yayınladığında her seferinde çağrılır;
  • [onError(Throwable th)], gözlemlenebilir bir th istisnası yayınladığında çağrılır;
  • [onCompleted], gözlemlenebilir nesne yayınlamayı bitirdiğini bildirdiğinde çağrılır;

Kodun işleyişi, daha önce açıklananla benzerdir. Aşağıdaki sonuçlar elde edilir:

1
2
3
4
next : 1
next : 2
next : 3
completed

7.1.3. Örnek-04: [Observable.create] yöntemi

  

Observable.create statik yöntemi şu şekilde tanımlanmıştır:

 
  • [create] yöntemi, Observable<T> türünde bir değer döndürür;
  • [create] yönteminin parametresi, aşağıdaki şekilde tanımlanmış [Observable.OnSubscribe<T>] türünde bir işlevdir:
 

[Observable.OnSubscribe<T>] türü, [Action1<Subscriber<? super T>>] işlevsel arayüzünü genişleten bir işlevsel arayüzdür. Bu arayüzün [call] yöntemi, aşağıdaki şekilde tanımlanmış bir [Subscriber] türü (abone, kayıtlı kullanıcı, gözlemci) bekler:

 

[1]'te, [Subscriber<T>] sınıfının 7.1.2. paragrafında sunulan [Observer<T>] arayüzünü uyguladığı görülmektedir.

Sonuç olarak, [<T> Observable.create] yöntemi:

  • parametre olarak, tek yöntem imzası void call(Subscriber<T> s) olan [Observable.OnSubscribe<T>] türünde bir örnek bekler. [Subscriber<T>] türü, [Observer<T>] türünü genişletir ve bu nedenle onNext, onError, onCompleted yöntemlerine sahiptir;
  • bir Observable<T> türü döndürür;

[<T> Observable.create] yöntemi, yapılandırılmış bir gözlemlenebilir döndürür. Henüz hiçbir öğe yayınlanmamıştır. Bir [Subscriber<T> s] abonesi bu gözlemlenebilir nesneye abone olduğunda, [<T> Observable.create] yöntemine parametre olarak geçirilen fonksiyonun [void call(s)] yöntemi çağrılır. Bu yöntemin görevi, T türündeki t öğelerini yayınlamak ve her yayınlamada gözlemcinin [s.onNext(t)] yöntemini çağırmaktır. Bu yöntem tamamlandığında, gözlemcinin [s.onCompleted(t)] yöntemi çağrılmalı ve [call] yöntemi sonlandırılmalıdır. [call] yöntemi bir th istisnasına rastlarsa, gözlemcinin [s.onError(th)] yöntemi çağrılmalı ve [call] yöntemi sonlandırılmalıdır;

Bu karmaşık işleyişi açıklamak için aşağıdaki [Exemple04] kodunu kullanacağız:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.util.Random;

public class Exemple04 {
    public static void main(String[] args) {
        // reel sayılar için gözlemlenebilir yapılandırma
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                for (int i = 0; i < 3; i++) {
                    // i elemanı yayını
                    subscriber.onNext(new Random((i + 1)).nextDouble());
                }
                // yayın sonu
                subscriber.onCompleted();
            }
        });
        // abonelik ve dolayısıyla yayın
        obs1.subscribe((d) -> System.out.printf("onNext %s%n", d), (th) -> System.out.printf("onError %s%n", th),
                () -> System.out.println("onCompleted"));
    }
}
  • 11. satır: Double türlerini yayan bir gözlemlenebilir oluşturulur;
  • 11-21. satırlar: [create] yönteminin parametresi, 12-20. satırlarda yer alan tek yöntem [call]'e sahip anonim bir sınıfla örneklenir. 11. satırda oluşturulan gözlemlenebilir, yayın yapmaya hazırdır ancak yalnızca bir gözlemci geldiğinde yayın yapacaktır;
  • 13-21. satırlar: [call] yöntemi, bir gözlemcinin referansını alır;
  • 14-17. satırlar: gözlemciye 3 öğe gönderilir;
  • 19. satır: gözlemciye yayın sonu bildirimi;
  • satır 23-24: satır 11'deki gözlemlenene abone olunur. [subscribe] yönteminin üç [onNext, onError, onCompleted] parametresi, üç lambda ile uygulanır. Bu abonelik, 13. satırdaki [call] yöntemine aktarılacak olan [Subscriber<Double>] abonesini oluşturacaktır. Ardından öğelerin yayını başlayacaktır;
  • her şey aynı iş parçacığında gerçekleşir: gözlemlenebilir ve gözlemci;

Aşağıdaki sonuçlar elde edilir:

1
2
3
4
onNext 0.7308781907032909
onNext 0.7311469360199058
onNext 0.731057369148862
onCompleted

[Observable.create] yöntemi, herhangi bir olaydan bir gözlemlenebilir oluşturmaya olanak tanır. Keşif bölümünün 2. paragrafında, senkron bir arayüzü asenkron bir arayüze dönüştürmek için kullandığımız yöntem budur.

7.1.4. Örnek-05: [Exemple-04]'in yeniden yapılandırılması

  

Aşağıdaki örnek, [Observable.subscribe] statik yönteminin yeni bir sürümünü göstermektedir:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;

public class Exemple05 {
    public static void main(String[] args) {
        // reel sayılardan oluşan bir gözlemlenebilirin yapılandırılması
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // bekleme
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // hata
                        subscriber.onError(e);
                    }
                    // eylem
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // bitti
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // bir abone
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // abonelik
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}
  • 56. satır: [Observable.subscribe] statik yönteminin yeni sürümü, bir önceki paragrafta tanıttığımız [Subscriber] türünü parametre olarak kabul eder;
  • 37-52. satırlar: abone (gözlemci). Observer arayüzünü, onNext, onError ve onCompleted adlı üç yöntemiyle uygular;
  • satır 61-64: bundan sonra, gözlemlenebilir nesnenin ve gözlemcisinin yürütüldüğü iş parçacıklarına odaklanacağız;
  • 62. satır: iş parçacığının adı;
  • satır 63: saniye ve milisaniye cinsinden ifade edilen geçerli zaman. Bu, gözlemlenebilirin öğeleri ne zaman yayınladığını ve gözlemcinin bunları ne zaman işlediğini zaman içinde görmemizi sağlayacaktır;
  • Bu kod, önceki kodla aynı işlevselliğe sahiptir. Sadece önceki kodu yeniden düzenledik;

Elde edilen sonuçlar şunlardır:

avant souscription ------Thread[main] ---- Time[31:685]
Observable.call start ------Thread[main] ---- Time[31:691]
Observable.call onNext(80.39999999999999) ------Thread[main] ---- Time[32:194]
Subscriber.onNext (80.39999999999999) ------Thread[main] ---- Time[32:195]
Observable.call onNext(73.2) ------Thread[main] ---- Time[32:595]
Subscriber.onNext (73.2) ------Thread[main] ---- Time[32:595]
Observable.call onNext(106.8) ------Thread[main] ---- Time[32:897]
Subscriber.onNext (106.8) ------Thread[main] ---- Time[32:897]
Observable.call onCompleted ------Thread[main] ---- Time[32:898]
Subscriber.onCompleted ------Thread[main] ---- Time[32:898]
après souscription ------Thread[main] ---- Time[32:899]
  • Sonuçların 1. satırı: Kodun 56. satırından önce henüz hiçbir şey gerçekleşmemiştir. Gözlemlenebilir sadece yapılandırılmıştır;
  • sonuçların 2. satırı: kodun 56. satırı, 15. satırdaki [call] yönteminin çağrılmasına neden olur. 3. satırda, 80,39 sayısal değeri gözlemciye gönderilir;
  • 4. satır: gözlemci, gönderilen sayıyı alır;
  • 5-8. satırlar: önceki işlem 2 kez tekrarlanır;
  • 9. satır: gözlemlenebilir, gönderim sonu bildirimini gönderir;
  • 10. satır: gözlemci bunu alır;
  • satır 11: kodun 57. satırı tarafından görüntülenir;

Görüldüğü gibi, yalnızca 56. satırdaki abonelik, sonuçların 2-10. satırlarının görüntülenmesine neden olmuştur. RxJava kütüphanesini kullanmaya başladığımızda, olayların birbiriyle nasıl bağlantılı olduğu ve özellikle gözlemci ile gözlemleneni birbirine bağlayan ilişkiler nasıl işlediği merak edilir. Burada görüldüğü gibi, 56. satır, yani gözlemlenene abonelik,

  • gözlemlenenin tüm öğelerinin gönderilmesine neden olmuştur;
  • gözlemlenenin ve gözlemcinin aynı iş parçacığında çalıştığını;
  • bu nedenle şu sırayı gözlemliyoruz: i öğesinin gönderilmesi, i öğesinin gözlemlenmesi, (i+1) öğesinin gönderilmesi, (i+1) öğesinin gözlemlenmesi, ...

Yayınlayıcının, öğelerini yayınlamadan önce beklediğini hatırlayalım:


                    // bekleme
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // hata
                        subscriber.onError(e);
}

burada 3. satırdaki i, gönderim numarasını temsil eder (0<=i<3). Gözlemlenebilirin elemanlarının gönderim zamanlarına bakarsak:

  • 2. ve 3. satırlar: 0 numaralı eleman, aboneliğin başlamasından yaklaşık 500 ms sonra gönderilmiştir;
  • 3. ve 5. satırlar: 1 numaralı eleman, 0 numaralı elemandan yaklaşık 400 ms sonra gönderilmiştir;
  • satır 5, 7: eleman 2, eleman 1'den yaklaşık 300 ms sonra yayınlanmıştır;

7.2. Çalıştırma iş parçacığı, gözlem iş parçacığı

7.2.1. Örnek-06: [main] dışındaki bir iş parçacığında gözlemlenebilir ve gözlemci

  

Önceki örneği şu şekilde yeniden düzenliyoruz: [Exemple06]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;
import java.util.concurrent.CountDownLatch;

public class Exemple06 {
    public static void main(String[] args) {

        // bariyer koruyucu
        CountDownLatch latch = new CountDownLatch(1);

        // reel sayılardan oluşan bir gözlemlenebilirin yapılandırılması
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // bekleme
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // hata
                        subscriber.onError(e);
                    }
                    // eylem
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // bitti
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // bir abone
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // bariyer indiriliyor
                latch.countDown();
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // gözlemlenebilir yapılandırma devam ediyor
        obs1 = obs1.subscribeOn(Schedulers.computation());
        // abonelik
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // bariyer önünde bekleme
        try {
            showInfos("début attente barrière");
            latch.await();
            showInfos("fin attente barrière");
        } catch (InterruptedException e1) {
            System.out.println(e1);
        }
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}
  • 16. satır: [CountDownLatch] türünde bir nesne ile bir bariyer (semafor) oluşturulur. Bu nesne, iş parçacıkları arasında senkronizasyon sağlamak için kullanılır. Burada, bariyerin (veya semaforun) değeri olarak adlandıracağımız 1 değeriyle başlatılır. Bir iş parçacığı, aşağıdaki işlemle bariyerin serbest kalmasını bekler:

latch.await();

Bariyer değeri >0 ise iş parçacığı bloke edilir. Bir iş parçacığı, bariyerin iç değerini artırabilir veya azaltabilir. 48. satırda, bariyer değeri 1 azaltılır.

  • 63. satır: Gözlemlenebilir, [Schedulers.computation()] zamanlayıcısı tarafından sağlanan bir iş parçacığı üzerinde çalışacak şekilde yapılandırılmıştır. Bu zamanlayıcı, yürütme makinesindeki çekirdek sayısı kadar iş parçacığı sağlayabilir. Örnek uygulamanın ilgili bölümünde diğer zamanlayıcıların kullanımı gösterilmiştir (bkz. bölüm 2.8);

Kodun çalışma prensibi şöyledir:

  • [main] yöntemi ana iş parçacığında (main) çalışır;
  • 66. satır: gözlemlenebilir nesnenin elemanlarının yayınlanmasını başlatır. Bu elemanlar, ana iş parçacığından farklı bir iş parçacığında yayınlanacaktır;
  • 70. satır: bariyer değeri 1 olduğu için (bkz. 16. satır) ana iş parçacığı bloke edilir. Bu değer 0'a düştüğünde devam edebilir. Bu, 48. satırda gerçekleşir. Gözlemlenebilirin yayınlarını tamamladığına dair bildirimi aldığında bariyeri indiren gözlemcidir;

Çalıştırma sonucu şu sonuçları verir:

avant souscription ------Thread[main] ---- Time[09:268]
Observable.call start ------Thread[RxComputationThreadPool-1] ---- Time[09:278]
début attente barrière ------Thread[main] ---- Time[09:278]
Observable.call onNext(44.4) ------Thread[RxComputationThreadPool-1] ---- Time[09:783]
Subscriber.onNext (44.4) ------Thread[RxComputationThreadPool-1] ---- Time[09:783]
Observable.call onNext(18.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:183]
Subscriber.onNext (18.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:184]
Observable.call onNext(54.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:486]
Subscriber.onNext (54.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:488]
Observable.call onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[10:489]
Subscriber.onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[10:490]
fin attente barrière ------Thread[main] ---- Time[10:491]
après souscription ------Thread[main] ---- Time[10:493]
  • 1. satır: abonelik gerçekleşecek;
  • 2. satır: Bu, [call] yönteminin [RxComputationThreadPool-1] iş parçacığı üzerinde yürütülmesini tetikler. Artık iki iş parçacığıyla paralel bir yürütme söz konusudur;
  • satır 3: nedeni bilinmeyen bir sebepten dolayı, [RxComputationThreadPool-1] iş parçacığı kontrolü devretti. Bunun üzerine [main] iş parçacığı kontrolü devralır ve bariyer tarafından engellenir (kodun 70. satırı). Bu andan itibaren yalnızca [RxComputationThreadPool-1] iş parçacığı işlem yapabilir;
  • 4-11. satırlar: Gözlemlenebilir ile gözlemcisi arasında daha önce gözlemlenen davranış tekrarlanır, ancak artık her şey [RxComputationThreadPool-1] iş parçacığı üzerinden gerçekleşir;
  • satır 12-13: gözlemci bariyeri indirdi (kodun 48. satırı) ve [RxComputationThreadPool-1] iş parçacığı sona erdi. [main] iş parçacığı devreye girer ve iki mesaj görüntüler;

7.2.2. Örnek-07: İki farklı iş parçacığında gözlemlenebilir ve gözlemci

  

Önceki örneği şu şekilde değiştiriyoruz:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;
import java.util.concurrent.CountDownLatch;

public class Exemple07 {
    public static void main(String[] args) {

        // bariyer bekçisi
        CountDownLatch latch = new CountDownLatch(1);

        // gerçek sayılardan oluşan bir gözlemlenebilirin yapılandırılması
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // bekleme
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // hata
                        subscriber.onError(e);
                    }
                    // eylem
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // bitti
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // bir abone
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // bariyer indiriliyor
                latch.countDown();
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // gözlemlenebilir yapılandırma devam ediyor
        obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
        // abonelik
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // bariyerin yükselmesini bekleme
        try {
            showInfos("début attente barrière");
            latch.await();
            showInfos("fin attente barrière");
        } catch (InterruptedException e1) {
            System.out.println(e1);
        }
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}

Kod, 63. satır hariç önceki örnektekiyle aynıdır:


obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());

Bu satır, gözlemleneni (subscribeOn) ve gözlemciyi (observeOn), [Schedulers.computation()] zamanlayıcısı tarafından sağlanan iş parçacıklarından birinde çalıştırılacak şekilde yapılandırır.

Elde edilen sonuçlar şunlardır:

avant souscription ------Thread[main] ---- Time[09:643]
début attente barrière ------Thread[main] ---- Time[09:656]
Observable.call start ------Thread[RxComputationThreadPool-4] ---- Time[09:656]
Observable.call onNext(39.6) ------Thread[RxComputationThreadPool-4] ---- Time[10:162]
Subscriber.onNext (39.6) ------Thread[RxComputationThreadPool-3] ---- Time[10:163]
Observable.call onNext(98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[10:562]
Subscriber.onNext (98.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[10:564]
Observable.call onNext(46.8) ------Thread[RxComputationThreadPool-4] ---- Time[10:864]
Observable.call onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[10:866]
Subscriber.onNext (46.8) ------Thread[RxComputationThreadPool-3] ---- Time[10:866]
Subscriber.onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[10:868]
fin attente barrière ------Thread[main] ---- Time[10:869]
après souscription ------Thread[main] ---- Time[10:870]

Aşağıdaki noktalar dikkat çekmektedir:

  • gözlemlenebilir, [RxComputationThreadPool-4] iş parçacığında çalışır (satır 3-4, 6, 8-9);
  • gözlemci, [RxComputationThreadPool-3] iş parçacığında çalışır (5., 7., 10-11. satırlar);
  • her ikisi de bağımsız olarak çalışır. Dolayısıyla 8-9. satırlarda, gözlemlenebilir, gözlemcinin [onNext] bildirimini almadan önce 2 bildirim (onNext, onCompleted) gönderir (10. satır);

RxJava kütüphanesi, gözlemlenebilirin iş parçacığından gözlemcinin iş parçacığına veri aktarımını (yayınları) üstlenir. Geliştiricinin bununla ilgilenmesi gerekmez.

Gözlemlenebilirlerin (Observable.from, Observable.create) nasıl oluşturulduğunu gördük. Şimdi RxJava kütüphanesindeki önceden tanımlanmış gözlemlenebilirlere bakalım.

7.3. Önceden tanımlanmış gözlemlenebilirler

7.3.1. Örnek-08: [Observable.range] yöntemi

 

Bundan sonra, gözlemlenen süreçler ve bunların gözlemcileri için özel sınıflar kullanacağız. Buradaki amaç, zaman içinde bu süreçleri takip edebilmek için isimlerini, yürütme iş parçacıklarını ve yürütme saatlerini günlüğe kaydedebilmektir.

[Process] sınıfı, basitçe isim verilebilen bir Observable olacaktır. Bu sınıf, aşağıdaki [IProcess] arayüzünü uygulayacaktır:


package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcess<T> {

    // gözlemlenebilirin adı
    public String getName();

    // gözlemlenebilir değişken
    public Observable<T> getObservable();

}

Bu arayüz, aşağıdaki [Process<T>] sınıfı tarafından uygulanabilir:


package dvp.rxjava.observables.utils;

import rx.Observable;
import rx.Scheduler;

public class Process<T> implements IProcess<T>{

    // gözlemlenebilirin adı
    protected String name;
    // gözlemlenen süreç
    protected Observable<T> observable;

    // oluşturucular
    public Process(String name, Observable<T> observable) {
        // yerel başlatmalar
        this.name = name;
        this.observable = observable;
    }

    // alıcı ve ayarlayıcılar
    public String getName() {
        return name;
    }

    public Observable<T> getObservable() {
        return observable;
    }

}
  • 9. satır: işlemin adı;
  • 11. satır: gözlemlenen değişken;
  • 14-18. satırlar: oluşturucu;

Gözlemci ise aşağıdaki [Observateur] sınıfı ile tanımlanacaktır:


package dvp.rxjava.observables.utils;

import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import rx.Subscriber;

public class Observateur<T> extends Subscriber<T> {

...
}
  • 11. satırda, Observateur<T> sınıfı, 7.1.3. paragrafında kısaca tanıttığımız Subscriber<T> sınıfını genişletir. Bunu, [Observable.subscribe] yönteminin argümanı olarak kullanacağız:

// gözlemlenebilir yürütme (gözlem)
obs1.subscribe(observateur);

Yukarıdaki 2. satırda kullanılan [Observable.subscribe] yöntemi şu şekilde tanımlanmıştır:

 

[Subscriber]'in rolü, esas olarak [Observer] arayüzünün yöntemleri aracılığıyla abone olduğu gözlemlenebilir tarafından gönderilen öğeleri yönetmektir: onNext, onError, onCompleted. [Subscriber] sınıfı aşağıdaki yöntemlere sahiptir:

 

[Observateur] sınıfının kodunda, abonenin aboneliğinin iptal edilip edilmediğini öğrenmek için [1] ve isUnsubscribed yöntemlerini kullanacağız. [Observateur<T>] sınıfının tam kodu aşağıdaki gibidir:


package dvp.rxjava.observables.utils;

import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import rx.Subscriber;

public class Observateur<T> extends Subscriber<T> {

    // bir bariyer (semafor)
    private CountDownLatch latch;
    // bir görüntüleme yöntemi
    private Consumer<String> showInfos;
    // gözlemcinin adı
    private String observerName;
    // gözlemlenen sürecin adı
    private String processName;

    // yapıcılar
    public Observateur() {

    }

    public Observateur(String name, CountDownLatch latch, Consumer<String> showInfos, String observedName) {
        this.observerName = name;
        this.latch = latch;
        this.showInfos = showInfos;
        this.processName = observedName;
    }

    // --------------------------- Observer<T> arayüzünün uygulaması
    @Override
    public void onCompleted() {
        // yayınların sonu
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s,%s].onCompleted", observerName, processName));
        }
        // ana iş parçacığının kilitlenmesi sona erdi
        latch.countDown();
    }

    @Override
    public void onError(Throwable e) {
        // yayın hatası
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s, %s].onError (%s)", observerName, processName, e));
        }
    }

    @Override
    public void onNext(T value) {
        // ek bir yayın
        if (!isUnsubscribed()) {
            try {
                showInfos.accept(String.format("Subscriber [%s,%s] : onNext (%s)", observerName, processName,
                        new ObjectMapper().writeValueAsString(value)));
            } catch (JsonProcessingException e) {
                showInfos.accept(String.format("Subscriber [%s,%s].onNext (%s)", observerName, processName, e));
            }
        }
    }
}
  • Bir Subscriber'in özelliklerine ek olarak, Observateur gözlemcisi aşağıdaki bilgileri de beraberinde taşıyacaktır:
    • 14. satır: Gözlemci, gözlemlenebilir tarafından gönderilen tüm öğeleri alana kadar ana iş parçacığını bloke etmek için kullanılacak bir bariyer veya semafor. Bu işlem, gözlemci gözlemlenebilirden gönderim sonu bildirimini aldığında kodun 36. satırında gerçekleşecektir;
    • 16. satır: Konsolda bir mesaj görüntülemek için kullanılacak bir Consumer<String> örneği;
    • 18. satır: Birden fazla gözlemci olduğunda bunları birbirinden ayırt etmek için kullanılan gözlemcinin adı;
    • 20. satır: gözlemlenen işlemin adı;
  • satır 36, 46, 54: [Subscriber<T>] soyut sınıfı tarafından uygulanan [Observer<T>] arayüzünün [onCompleted, onError, onNext] yöntemleri. Bu sınıf bunları uygulamaz. Bu nedenle, bu işlemler alt sınıflarda gerçekleştirilmelidir. Bu yöntemlerde herhangi bir işlem yapmadan önce, gözlemcinin gözlemlediği gözlemlenebilir öğeden aboneliğinin iptal edilip edilmediğine bakılır;
  • 59. satır: Gözlemcinin [onNext] yöntemi, alınan öğenin jSON dizesini yazar. Bu, çeşitli öğe türlerini görüntülememizi sağlayacaktır;

Bunu belirttikten sonra, Observable sınıfının yeni bir yöntemi olan [range] yöntemini inceleyelim:

 

Observable.range(n,m) gözlemlenebilir öğesi, n ile n+m-1 arasında değişen (m) tamsayı üretir. Bunu aşağıdaki [Exemple08] koduyla inceleyeceğiz:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.Observateur;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple08 {
    public static void main(String[] args) throws InterruptedException {

        // gözlemci sayısı
        final int nbObservateurs = 2;

        // semafor
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // gözlemlenebilir yapılandırma
        Observable<Integer> obs1 = Observable.range(15, 3).subscribeOn(Schedulers.computation());
        // gözlemlenebilir yürütme (gözlem)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur[%d]", i), latch, showInfos,"obs1"));
        }
        // bekleme
        showInfos.accept("main : attente fin observation");
        latch.await();
        // son
        showInfos.accept("main : fin observation");
    }

    // görüntülemeler
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • 16. satır: iki gözlemci kullanacağız;
  • satır 19: her gözlemciyi farklı bir iş parçacığına yerleştireceğimiz için bariyer (semafor) iki olarak başlatılır. Dolayısıyla ana iş parçacığı, iki gözlem iş parçacığının tamamlanmasını beklemek zorunda kalacaktır;
  • 22. satır: gözlemleneni, [Schedulers.computation()] zamanlayıcısının bir iş parçacığı üzerinde çalışacak şekilde yapılandırıyoruz. Gözlemci, gözlemleneni ile aynı iş parçacığı üzerinde olacaktır;
  • 25-27. satırlar: Gözlemlenene iki gözlemci abone oluyor. Bu, her bir gözlemci için gözlemlenenin tamamen yürütülmesini tetikleyecektir: 15, 16 ve 17 tamsayıları gönderilecektir;
  • 30. satır: ana iş parçacığı, gözlemcilerin işlerinin bitmesini bekler;

Elde edilen sonuçlar şunlardır:

main : début observation ------Thread[main] ---- Time[27:875]
main : attente fin observation ------Thread[main] ---- Time[27:893]
Subscriber[observateur[1],obs1] : onNext (15) ------Thread[RxComputationThreadPool-2] ---- Time[28:245]
Subscriber[observateur[0],obs1] : onNext (15) ------Thread[RxComputationThreadPool-1] ---- Time[28:245]
Subscriber[observateur[1],obs1] : onNext (16) ------Thread[RxComputationThreadPool-2] ---- Time[28:247]
Subscriber[observateur[0],obs1] : onNext (16) ------Thread[RxComputationThreadPool-1] ---- Time[28:248]
Subscriber[observateur[1],obs1] : onNext (17) ------Thread[RxComputationThreadPool-2] ---- Time[28:249]
Subscriber[observateur[1],obs1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[28:250]
Subscriber[observateur[0],obs1] : onNext (17) ------Thread[RxComputationThreadPool-1] ---- Time[28:251]
Subscriber[observateur[0],obs1].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[28:252]
main : fin observation ------Thread[main] ---- Time[28:252]
  • 2. satır: ana iş parçacığı, 2 gözlemcinin tamamlanmasını beklerken bloke olur;
  • satır 3-4: gözlemci 0'ın [RxComputationThreadPool-1] iş parçacığında, gözlemci 1'in ise [RxComputationThreadPool-2] iş parçacığında olduğu görülmektedir;
  • satır 3-10: her iki gözlemcinin de tam olarak aynı elemanları aldığı görülüyor;

Diğer gözlemlenebilir türlerin davranışını açıklamak için bu şekilde tanımlanan Observateur sınıfını kullanacağız.

7.3.2. Örnek-09: Observable.[interval, take, doNext] yöntemleri

  
 

Bu örnek, düzenli zaman aralıklarında uzun tamsayılar üreten Observable.interval (uzun aralık, TimeUnit birimi) gözlemlenebilirinin kullanımını göstermektedir. [1] noktasına dikkat edilmelidir: varsayılan olarak, [Observable.interval] gözlemlenebilir nesnesi, [Schedulers.computation] zamanlayıcısının iş parçacıklarından birinde çalışır.

Kod şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.Observateur;
import rx.Observable;

public class Exemple09 {
    public static void main(String[] args) throws InterruptedException {

        // gözlemci sayısı
        final int nbObservateurs = 2;

        // semafor
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // gözlemlenebilir yapılandırma
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // gözlemlenebilir yürütme (gözlem)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur [%d]", i), latch, showInfos,
                    "obs1"));
        }
        // bekleme
        showInfos.accept("main : attente fin observation");
        latch.await();
        // bitiş
        showInfos.accept("main : fin observation");
    }

    // görüntülemeler
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • 22. satır: gözlemlenebilir, her 500 milisaniyede bir uzun tamsayılar üretir. Seri, 0 sayısıyla başlar;
  • 22. satır: Bu gözlemlenebilir, sonsuz sayıda değer üretir. [Observable.take(n)] yöntemi, üretilen ilk n öğeyi saklayan yeni bir gözlemlenebilir oluşturur;
 

Gözlemlenebilirin koduna geri dönelim:


Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
.doOnNext(l -> showInfos.accept(l.toString()));

2. satırda, [Observable.doOnNext] yöntemi, gözlemlenebilir her yeni bir eleman yayınladığında çalışır. Bu yöntem genellikle bilgileri günlüğe kaydetmek için kullanılır. Burada, 500 milisaniyelik aralığın doğru bir şekilde kontrol edilip edilmediğini doğrulamak için elemanların yayınlanma tarihini günlüğe kaydetmek istiyoruz. [Observable.doOnNext] yöntemi, uygulandığı gözlemlenebilir nesneyi değiştirmez. Tanımı şöyledir:

 

Yürütme işlemi şu sonuçları verir:

main : début observation ------Thread[main] ---- Time[55:892]
main : attente fin observation ------Thread[main] ---- Time[55:911]
0 ------Thread[RxComputationThreadPool-1] ---- Time[56:412]
0 ------Thread[RxComputationThreadPool-2] ---- Time[56:413]
Subscriber[observateur [1],obs1] : onNext (0) ------Thread[RxComputationThreadPool-2] ---- Time[56:723]
Subscriber[observateur [0],obs1] : onNext (0) ------Thread[RxComputationThreadPool-1] ---- Time[56:723]
1 ------Thread[RxComputationThreadPool-1] ---- Time[56:906]
Subscriber[observateur [0],obs1] : onNext (1) ------Thread[RxComputationThreadPool-1] ---- Time[56:908]
1 ------Thread[RxComputationThreadPool-2] ---- Time[56:912]
Subscriber[observateur [1],obs1] : onNext (1) ------Thread[RxComputationThreadPool-2] ---- Time[56:914]
2 ------Thread[RxComputationThreadPool-1] ---- Time[57:405]
Subscriber[observateur [0],obs1] : onNext (2) ------Thread[RxComputationThreadPool-1] ---- Time[57:407]
Subscriber[observateur [0],obs1].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[57:408]
2 ------Thread[RxComputationThreadPool-2] ---- Time[57:412]
Subscriber[observateur [1],obs1] : onNext (2) ------Thread[RxComputationThreadPool-2] ---- Time[57:414]
Subscriber[observateur [1],obs1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[57:415]
main : fin observation ------Thread[main] ---- Time[57:416]
  • 3., 7. ve 11. satırlar: Gönderim aralığının yaklaşık olarak 500 ms'ye yakın olduğu görülmektedir;
  • Her iki gözlemci de elbette farklı iş parçacıklarında yer alıyor; oysa gözlemlenebilir nesne belirli bir zamanlayıcıyla çalışacak şekilde yapılandırılmamıştı. Burada gördüğümüz, [Observable.interval] gözlemlenebilir nesnesinin varsayılan çalışma şeklidir;

7.3.3. Örnekler-10/12: Observable.[error, empty, never] yöntemleri

 

Bundan sonra, [Observable] sınıfının yöntemlerini örneklerimizi daha kısa ve öz bir şekilde sunacağız. Önceki kod şöyleydi:


package dvp.rxjava.observables;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;

import rx.Observable;

public class Exemple09 {
    public static void main(String[] args) throws InterruptedException {

        // gözlemci sayısı
        final int nbObservateurs = 2;

        // semafor
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // gözlemlenebilir yapılandırma
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // gözlemlenebilir yürütme (gözlem)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur [%d]", i), latch, showInfos,
                    "obs1"));
        }
        // bekleme
        showInfos.accept("main : attente fin observation");
        latch.await();
        // son
        showInfos.accept("main : fin observation");
    }

    // görüntüler
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}

Bu kod, önceki örnekte de kullanılmıştı. Yalnızca 21-22. satırlar değişmişti. Bu nedenle, bu kodun büyük bir kısmını aşağıdaki [ProcessUtils] sınıfında faktörlere ayıracağız:


package dvp.rxjava.observables.utils;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import rx.Observable;

public class ProcessUtils {

    @SafeVarargs
    public static void subscribe(int nbObservateurs, IProcess<?>... processes) throws InterruptedException {

        // semafor
        CountDownLatch latch = new CountDownLatch(nbObservateurs * processes.length);

        // gözlemlenebilir yürütme (gözlem)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            for (IProcess<?> process : processes) {
                Observable<?> obs = process.getObservable();
                obs.subscribe(new Observateur<>(String.format("observateur[%d]", i), latch, showInfos, process.getName()));
            }
        }
        // bekleme
        showInfos.accept("main : attente fin observation");
        latch.await();
        // son
        showInfos.accept("main : fin observation");
    }

    // görüntülemeler
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • 13. satır: Yöntem iki parametre kabul eder:
    • nbObservateurs: ikinci parametre olarak geçirilen süreçlerin gözlemci sayısı;
    • processes: gözlemlenecek süreçler (adlandırılmış gözlemlenebilirler). [IProcess<?>] notasyonu sayesinde, süreçler farklı türlerdeki elemanları gönderebilecek;
  • 16. satır: Tüm gözlemciler tüm gözlemlerini tamamladığında semafor yeşile geçmelidir. Semaforun başlangıç değeri, gözlemci sayısı ile gözlem sayısının çarpımıdır;
  • satır 20-25: her gözlemci, gözlemlenmesi gereken tüm süreçlere abone edilir;
  • 23. satır: süreçten gözlemlenebilir değer alınır (bkz. paragraf 7.3.1);
  • satır 23: gözlemciye bu gözlemlenebilir değere abone olunur. Gözlemciye 4 bilgi aktarılır:
    • adı;
    • gözlemlediği gözlemlenebilirin yayın sonu bildirimini aldığında azaltması gereken semafor;
    • konsola bilgi kaydetmek istediğinde kullanacağı yöntem;
    • gözlemleyeceği işlemin adı;

Bu sınıflar tanımlandıktan sonra, örnek 10 şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple10 {
    public static void main(String[] args) throws InterruptedException {
        // gözlemlenebilir yapılandırma
        Observable<?> obs = Observable.error(new RuntimeException("Erreur !!!")).subscribeOn(Schedulers.computation());
        // gözlemlenebilir yürütme (gözlem)
        ProcessUtils.subscribe(2,new Process<>("process1", obs));
    }
}

11. satırda, [Observable.error] statik yöntemi şu şekilde tanımlanmıştır:

 

Dolayısıyla 8. satır, abonelerine [onError] yöntemine yönelik bir istisna atmakla yetinen bir gözlemleneni yapılandırır. Yürütme işlemi şu sonuçları verir:


main : début observation ------Thread[main] ---- Time[22:618]
main : attente fin observation ------Thread[main] ---- Time[22:636]
Subscriber[observateur[1], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-2] ---- Time[22:638]
Subscriber[observateur[0], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-1] ---- Time[22:638]

3. ve 4. satırlarda, her iki abonenin [onError] yöntemi, gözlemlenebilir tarafından atılan istisnayı almıştır.

Bu çalıştırmanın bir özelliği vardır: her iki gözlemcinin [onCompleted] yöntemleri çağrılmamıştır. Dolayısıyla, bariyer indirilmemiştir ve ana iş parçacığı, aşağıdaki 3. satırdaki [ProcessUtils.subscribe] statik yönteminde bloke kalmıştır:


// bekleme
showInfos.accept("main : attente fin observation");
latch.await();
// son
showInfos.accept("main : fin observation");

Burada, gözlemlenebilirde bir hata olması durumunda abonelerin [onCompleted] yönteminin çağrılmadığını görüyoruz. Bunun üzerine [Observateur.onError] yöntemini şu şekilde değiştiriyoruz:


    @Override
    public void onError(Throwable e) {
        // gönderim hatası
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber[%s, %s].onError (%s)", observerName, processName, e));
        }
        // ana iş parçacığı kilitlenmesi sonu
        latch.countDown();
}

Gözlemlenebilir değişkende hata olması durumunda engeli kaldırmak için 7-8. satırları ekliyoruz. Bu yeni kodla, çalıştırma aşağıdaki sonuçları veriyor:


main : début observation ------Thread[main] ---- Time[40:750]
main : attente fin observation ------Thread[main] ---- Time[40:764]
Subscriber[observateur[0], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-1] ---- Time[40:766]
Subscriber[observateur[1], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-2] ---- Time[40:766]
main : fin observation ------Thread[main] ---- Time[40:767]

Daha önce elde edemediğimiz 5. satırı elde ediyoruz.

Örnek 11 şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple11 {
    public static void main(String[] args) throws InterruptedException {
        // gözlemlenebilir yapılandırma
        Observable<?> obs1 = Observable.empty();
        // gözlemlenebilir yürütme (gözlem)
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

10. satırda, statik yöntem [Observable.empty], hiçbir eleman yaymayan bir gözlemlenebilir oluşturur. Yalnızca yayın sonu bildirimini yayar;

 

Yukarıdaki örnekteki kodun çalıştırılması şu sonuçları verir:

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[37:073]
Subscriber[observateur[0],process1].onCompleted ------Thread[main] ---- Time[37:086]
Subscriber[observateur[1],process1].onCompleted ------Thread[main] ---- Time[37:086]
main : attente fin observation ------Thread[main] ---- Time[37:087]
main : fin observation ------Thread[main] ---- Time[37:087]
  • 2. ve 3. satırlar: Her iki gözlemcinin de daha önce herhangi bir öğe almadan yayın sonu bildirimini aldıkları görülmektedir.

Bu yöntemin ne işe yaradığı merak edilebilir. Bu yöntem, başlangıçta boş olan ve daha sonra elemanların biriktirildiği bir koleksiyona benzer şekilde kullanılabilir:

1
2
3
4
Observable obs=Observable.empty() ;
for(Observable o : observables){
    obs=obs.mergeWith(o) ;
}

3. satırda, başlangıçtaki gözlemlenebilir obs (1. satır) diğer gözlemlenebilirlerle birleştirilir.

Örnek 12, statik [Observable.never] yöntemini göstermektedir:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple12 {
    public static void main(String[] args) throws InterruptedException {
        // gözlemlenebilir yapılandırma
        Observable<?> obs1 = Observable.never();
        // yürütme (gözlem) gözlemlenebilir
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

Statik [Observable.never] yöntemi, hiçbir zaman sinyal yaymayan bir gözlemlenebilir oluşturur:

 

Örneğin çalıştırılması şu sonuçları verir:

main : début observation ------Thread[main] ---- Time[27:018]
main : attente fin observation ------Thread[main] ---- Time[27:030]

2. satırda, ana iş parçacığı süresiz olarak bekler. Zira, semaforu (bariyer) yeşile çevirmeye (bariyeri indirme) olanak tanıyan [onCompleted] bildirimini gönderen hiçbir gözlemlenebilir yoktur.

7.4. Multi-threading

7.4.1. Örnek-13: eylem iş parçacığı, gözlem iş parçacığı

7.1.3 numaralı paragrafta, [Observable.create] statik yöntemi ile bir gözlemlenebilir oluşturduk:

 
  • [create] yöntemi, Observable<T> türünde bir değer döndürür;
  • [create] yönteminin parametresi, aşağıdaki şekilde tanımlanmış [Observable.OnSubscribe<T>] türünde bir işlevdir:
 

[Observable.OnSubscribe<T>] türü, [Action1<Subscriber<? super T>>] işlevsel arayüzünü genişleten bir işlevsel arayüzdür. Bu arayüzün [call] yöntemi, [Subscriber] türünde bir nesne (abone, izleyici) bekler. Bu belgenin devamında, [Observable.OnSubscribe<T>] türünü bazen eylem olarak adlandıracağız. Bir adı olacak özel eylemler oluşturacağız. Bunlar, aşağıdaki [IProcessAction] arayüzünün örnekleri olacaktır:

  

package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcessAction<T> extends Observable.OnSubscribe<T> {

    // eylemin bir adı var
    public String getName();
}
  • 5. satır: [IProcessAction<T>] arayüzü, [Observable.OnSubscribe<T>] arayüzünün tüm özelliklerine sahiptir;
  • 8. satır: Ayrıca, arayüzü uygulayan örneğin adını döndüren [getName] adlı bir yöntemi de vardır;

Aşağıdaki [ProcessAction01] adlı eylemi kullanacağız:


package dvp.rxjava.observables.utils;

import java.util.Random;

import rx.Subscriber;
import rx.functions.Func1;

public class ProcessAction01<T> implements IProcessAction<T> {

    // veri
    private String name;
    private int nbValues;
    private Func1<Integer, T> func1;

    // yapıcılar
    public ProcessAction01(String name, int nbValues, Func1<Integer, T> func1) {
        this.name = name;
        this.nbValues = nbValues;
        this.func1 = func1;
    }

    @Override
    public void call(Subscriber<? super T> subscriber) {
        ProcessUtils.showInfos.accept(String.format("Observable (%s) call start", getName()));
        for (int i = 0; i < nbValues; i++) {
            // bekleme
            try {
                Thread.sleep(new Random().nextInt(500));
            } catch (InterruptedException e) {
                // hata
                ProcessUtils.showInfos.accept(String.format("Observable (%s) onError", getName()));
                subscriber.onError(e);
            }
            // bir öğenin gönderilmesi
            T value = func1.call(i);
            ProcessUtils.showInfos.accept(String.format("Observable (%s,%s) onNext (%s)", getName(), i, value));
            subscriber.onNext(value);
        }
        // bitti
        ProcessUtils.showInfos.accept(String.format("Observable (%s) onCompleted", getName()));
        subscriber.onCompleted();
    }

    @Override
    public String getName() {
        return name;
    }

}
  • 8. satır: [ProcessAction01<T>] sınıfı, [IProcessAction<T>] arayüzünü ve dolayısıyla [Observable.OnSubscribe<T>] arayüzünü uygular;
  • 11. satır: eylemin adı;
  • 12. satır: gönderilecek değerlerin sayısı;
  • 13. satır: Bir tamsayıdan yola çıkarak gözlemlenebilir tarafından yayınlanacak bir T türünü oluşturan [Func1<Integer, T>] türünde bir örnek (35. ve 37. satırlar);
  • satır 16-20: oluşturucuya eylemin adı, gönderilecek değer sayısı ve gönderme işlevi aktarılır;
  • satır 23-42: süreç kodu;
  • 23. satır: [call] yöntemi, parametre olarak sürece bağlı gözlemlenebilirin abonesini alır;
  • 28. satır: süreç, rastgele bir bekleme süresinin ardından öğelerini yayınlar;
  • satır 32: bir hata yayımı;
  • satır 37: normal bir yayın;
  • satır 41: gönderim sonu bildiriminin gönderilmesi;
  • satır 25-38: eylem, rastgele bir bekleme süresinden sonra (satır 30) nbValues tabanlı gerçek değerler gönderir;
  • satır 35: gönderilecek değer, yapıcıya parametre olarak geçirilen [func1] işlevi tarafından sağlanır (satır 16);

[Process] sınıfını (bkz. paragraf 7.3.1), adlandırılmış bir eylemle de oluşturulabilmesi için yeniden düzenliyoruz. Sınıfa aşağıdaki oluşturucu eklenir:


public Process(IProcessAction<T> na, Scheduler schedulerObserved, Scheduler schedulerObserver) {
        // işlem adı=eylem adı
        name = na.getName();
        // eylem --> gözlemlenebilir
        observable = Observable.create(na);
        // gözlemlenen sürecin yürütme iş parçacığı
        if (schedulerObserved != null) {
            observable = observable.subscribeOn(schedulerObserved);
        }
        // gözlemcinin gözlem iş parçacığı
        if (schedulerObserver != null) {
            observable = observable.observeOn(schedulerObserver);
        }
    }
  • 1. satırda, oluşturucu 3 parametre kabul eder:
    1. gözlemleneni oluşturmak için kullanılacak adlandırılmış eylem (satır 5);
    2. gözlemlenen sürecin zamanlayıcısı (null olabilir);
    3. gözlemcinin zamanlayıcısı (örneğin null);
  • 5. satır: gözlemlenebilir, parametre olarak geçirilen eylemden oluşturulur;

Aşağıdaki kod [Exemple13], farklı gözlemlenebilirleri gözlemler:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple13 {
    public static void main(String[] args) throws InterruptedException {
        // süreç 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // süreç 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // işlem 3
        Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
                Schedulers.computation());
        // işlem 4
        Process<Boolean> process4 = new Process<>(new ProcessAction01<Boolean>("process4", 4, i -> i % 2 == 0), null, null);
        // abonelikler
        ProcessUtils.subscribe(1, process1);
        ProcessUtils.subscribe(1, process2);
        ProcessUtils.subscribe(1, process3);
        ProcessUtils.subscribe(1, process4);
    }
}
  • satır 13-15: process1 süreci, bir hesaplama iş parçacığında 1 gerçek sayı üretir ve bu sayı başka bir hesaplama iş parçacığında gözlemlenir;
  • satır 17-18: process2 süreci, bir hesaplama iş parçacığında 2 karakter dizisi üretir ve gözlemcinin iş parçacığı hakkında herhangi bir bilgi verilmez. Sonuçlar, gözlemin varsayılan olarak sürecin yürütüldüğü iş parçacığı üzerinde yapıldığını gösterir;
  • satır 20-21: process3 süreci, bir hesaplama iş parçacığı üzerinde gözlemlenecek olan 3 tamsayıyı, kısıtlanmamış bir iş parçacığı üzerinde üretir. Sonuçlar, sürecin yürütülmesinin varsayılan olarak ana iş parçacığı üzerinde gerçekleştiğini göstermektedir;
  • 23. satır: process4 süreci, belirlenmemiş bir iş parçacığında 4 boole değeri üretir ve bunlar da belirlenmemiş bir iş parçacığında gözlemlenir. Sonuçlar, sürecin yürütülmesinin ve gözlemlenmesinin varsayılan olarak ana iş parçacığında gerçekleştiğini göstermektedir;

Bu kodun yürütülmesinin sonucu şöyledir:

main : début observation ------Thread[main] ---- Time[18:642]
main : attente fin observation ------Thread[main] ---- Time[18:660]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[18:660]
Observable (process1,0) onNext (68.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[19:093]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[19:094]
Subscriber[observateur[0],process1] : onNext (68.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[19:396]
Subscriber[observateur[0],process1].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[19:397]
main : fin observation ------Thread[main] ---- Time[19:397]
main : début observation ------Thread[main] ---- Time[19:398]
main : attente fin observation ------Thread[main] ---- Time[19:399]
Observable (process2) call start ------Thread[RxComputationThreadPool-5] ---- Time[19:399]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-5] ---- Time[19:630]
Subscriber[observateur[0],process2] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-5] ---- Time[19:631]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-5] ---- Time[20:094]
Subscriber[observateur[0],process2] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-5] ---- Time[20:095]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[20:096]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[20:096]
main : fin observation ------Thread[main] ---- Time[20:097]
main : début observation ------Thread[main] ---- Time[20:097]
Observable (process3) call start ------Thread[main] ---- Time[20:098]
Observable (process3,0) onNext (0) ------Thread[main] ---- Time[20:188]
Subscriber[observateur[0],process3] : onNext (0) ------Thread[RxComputationThreadPool-6] ---- Time[20:213]
Observable (process3,1) onNext (2) ------Thread[main] ---- Time[20:336]
Subscriber[observateur[0],process3] : onNext (2) ------Thread[RxComputationThreadPool-6] ---- Time[20:338]
Observable (process3,2) onNext (4) ------Thread[main] ---- Time[20:676]
Observable (process3) onCompleted ------Thread[main] ---- Time[20:677]
main : attente fin observation ------Thread[main] ---- Time[20:677]
Subscriber[observateur[0],process3] : onNext (4) ------Thread[RxComputationThreadPool-6] ---- Time[20:678]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[20:679]
main : fin observation ------Thread[main] ---- Time[20:679]
main : début observation ------Thread[main] ---- Time[20:680]
Observable (process4) call start ------Thread[main] ---- Time[20:680]
Observable (process4,0) onNext (true) ------Thread[main] ---- Time[21:065]
Subscriber[observateur[0],process4] : onNext (true) ------Thread[main] ---- Time[21:067]
Observable (process4,1) onNext (false) ------Thread[main] ---- Time[21:187]
Subscriber[observateur[0],process4] : onNext (false) ------Thread[main] ---- Time[21:188]
Observable (process4,2) onNext (true) ------Thread[main] ---- Time[21:624]
Subscriber[observateur[0],process4] : onNext (true) ------Thread[main] ---- Time[21:625]
Observable (process4,3) onNext (false) ------Thread[main] ---- Time[21:765]
Subscriber[observateur[0],process4] : onNext (false) ------Thread[main] ---- Time[21:766]
Observable (process4) onCompleted ------Thread[main] ---- Time[21:767]
Subscriber[observateur[0],process4].onCompleted ------Thread[main] ---- Time[21:767]
main : attente fin observation ------Thread[main] ---- Time[21:767]
main : fin observation ------Thread[main] ---- Time[21:768]
  • process1 süreci, [RxComputationThreadPool-4] hesaplama iş parçacığı üzerinde 1 gerçek sayı üretir (4. satır) ve bu, [RxComputationThreadPool-3] hesaplama iş parçacığı üzerinde gözlemlenir (6. satır);
  • process2 süreci, [RxComputationThreadPool-5] hesaplama iş parçacığında 2 karakter dizesi (satır 12, 14) üretir; bu dizeler aynı iş parçacığında (satır 13, 15) gözlemlenir;
  • process3 süreci, ana iş parçacığı üzerinde 3 tamsayı üretir (satır 21, 23, 25) ve bunlar [RxComputationThreadPool-6] hesaplama iş parçacığı üzerinde gözlemlenir (satır 22, 24, 28);
  • process4 süreci, ana iş parçacığında 4 boole değeri üretir (satır 34, 36, 38, 40) ve bu değerler aynı ana iş parçacığında gözlemlenir (satır 33, 35, 37, 39);

Okuyucunun yukarıdakileri takip etmesi önerilir:

  • gözlemlenen işlemin ve iş parçacığının yaşam döngüsünü;
  • gözlemcisinin yaşam döngüsü ve iş parçacığı;

Rx kütüphanelerinin cazibesinin büyük bir kısmı, geliştiricinin kendi başına yönetmek zorunda olmadığı bu çoklu iş parçacığı özelliğine dayanmaktadır.

7.5. Birden fazla gözlemlenebilirin birleştirilmesi

7.5.1. Örnek-14: [Observable.merge] ile iki gözlemleneni birleştirmek

Şimdi, birden fazla gözlemleneni tek bir sonuç gözlemleneni içinde birleştirmeye olanak tanıyan [Observable] sınıfının statik yöntemlerini tanıtacağız.

Bu türden ilk örnek şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessAction01;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple14 {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // birleştirme
        Process<?> process12 = new Process<>("process12",
                Observable.merge(process1.getObservable(), process2.getObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process12);
    }
}
  • 15-17. satırlar: [process1] adlı bir işlem, bir hesaplama iş parçacığı üzerinden 3 gerçek sayı üretecektir. Bu işlem aynı zamanda bir hesaplama iş parçacığı üzerinde gözlemlenecektir;
  • satır 19-20: [process2] adlı bir işlem, bir hesaplama iş parçacığı üzerinden 2 karakter dizisi üretecektir. Gözlem iş parçacığı zorunlu değildir. Daha önce gördüğümüz gibi, bu durumda gözlem iş parçacığı hesaplama iş parçacığıdır;
  • satır 23: iki işlem birleştirilir, yani elemanları her iki işlemden aynı anda gelen bir gözlemlenebilir oluşturulur. Bunun için [Observable.merge] statik yöntemi kullanılır:
 

Yukarıdaki şemada gösterilenden farklı olarak, birleştirme sırasında akış 1'in elemanları, akış 2'nin elemanlarının arasına girebilir. Yürütme sonuçları da bunu göstermektedir:

main : début observation ------Thread[main] ---- Time[56:053]
main : attente fin observation ------Thread[main] ---- Time[56:073]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[56:073]
Observable (process2) call start ------Thread[RxComputationThreadPool-5] ---- Time[56:074]
Observable (process1,0) onNext (64.8) ------Thread[RxComputationThreadPool-4] ---- Time[56:263]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-5] ---- Time[56:403]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-5] ---- Time[56:515]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[56:516]
Subscriber[observateur[0],process12] : onNext (64.8) ------Thread[RxComputationThreadPool-3] ---- Time[56:552]
Subscriber[observateur[0],process12] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-3] ---- Time[56:553]
Subscriber[observateur[0],process12] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-3] ---- Time[56:553]
Observable (process1,1) onNext (56.4) ------Thread[RxComputationThreadPool-4] ---- Time[56:716]
Subscriber[observateur[0],process12] : onNext (56.4) ------Thread[RxComputationThreadPool-3] ---- Time[56:718]
Observable (process1,2) onNext (22.8) ------Thread[RxComputationThreadPool-4] ---- Time[57:082]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[57:083]
Subscriber[observateur[0],process12] : onNext (22.8) ------Thread[RxComputationThreadPool-3] ---- Time[57:084]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[57:085]
main : fin observation ------Thread[main] ---- Time[57:085]
  • 3. satır: [process1] işlemi, [RxComputationThreadPool-4] hesaplama iş parçacığı üzerinde yürütülür;
  • 4. satır: [process2] süreci, [RxComputationThreadPool-5] hesaplama iş parçacığı üzerinde çalışmaktadır;
  • 9. satır: [process12] işlemi, [RxComputationThreadPool-3] hesaplama iş parçacığı üzerinde gözlemleniyor. Bu seçime yol açan kuralı bilmiyorum;
  • 9-11. satırlar: Gözlemcinin, [process1] (5. satır) ve [process2] (6., 7. satırlar) süreçlerinin öğelerini gözlemlediği görülüyor, oysa bu iki süreçten hiçbiri tamamlanmamış (karışıklık var);
  • [process12] süreci, process1 ve process2 süreçlerinin her ikisi de sona erdiğinde (17. satır) sona erer;

7.5.2. Örnek-15: [Observable.concat] ile iki gözlemleneni birleştirmek

Şimdi şu kodu inceleyelim:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessAction01;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple15 {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, Schedulers.computation());
        // birleştirme
        Process<?> process12 = new Process<>("process12",
                Observable.concat(process1.getObservable(), process2.getObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process12);
    }
}
  • 15-17. satırlar: [process1] adlı bir işlem, bir hesaplama iş parçacığı üzerinden 3 gerçek sayı üretecektir. Bu işlem aynı zamanda bir hesaplama iş parçacığı üzerinden gözlemlenecektir;
  • 19-20. satırlar: [process2] adlı bir işlem, zorunlu olmayan bir iş parçacığına (burada varsayılan ana iş parçacığı) 2 karakter dizisi gönderecektir. Bu işlem, bir hesaplama iş parçacığında gözlemlenecektir;
  • 23. satır: iki işlem birleştirilir, yani elemanları her iki işlemden gelen bir gözlemlenebilir oluşturulur. Yayınlanan değerler birbirine karıştırılmaz. [process12] süreci, önce [process1] sürecinin tüm değerlerini, ardından da [process2] sürecinin değerlerini yayınlayacaktır. Bunun için [Observable.concat] statik yöntemi kullanılır:
 

Çalıştırma sonuçları şöyledir:

main : début observation ------Thread[main] ---- Time[30:162]
main : attente fin observation ------Thread[main] ---- Time[30:189]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[30:190]
Observable (process1,0) onNext (79.2) ------Thread[RxComputationThreadPool-4] ---- Time[30:681]
Observable (process1,1) onNext (98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[30:792]
Subscriber[observateur[0],process12] : onNext (79.2) ------Thread[RxComputationThreadPool-3] ---- Time[30:975]
Subscriber[observateur[0],process12] : onNext (98.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[30:976]
Observable (process1,2) onNext (84.0) ------Thread[RxComputationThreadPool-4] ---- Time[31:084]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[31:085]
Subscriber[observateur[0],process12] : onNext (84.0) ------Thread[RxComputationThreadPool-3] ---- Time[31:086]
Observable (process2) call start ------Thread[RxComputationThreadPool-3] ---- Time[31:087]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-3] ---- Time[31:556]
Subscriber[observateur[0],process12] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-5] ---- Time[31:557]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-3] ---- Time[31:608]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[31:609]
Subscriber[observateur[0],process12] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-5] ---- Time[31:609]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[31:610]
main : fin observation ------Thread[main] ---- Time[31:611]
  • 3-10. satırlar: [process1] süreci yürütülür ve [process12] süreci, [process1] tarafından üretilen değerleri gönderir;
  • 9. satır: [process1] süreci sona ermiştir;
  • 11-17. satırlar: [process2] süreci çalışıyor ve [process12] süreci, [process2] tarafından gönderilen değerleri iletiyor;

process2 işlemi ile ilgili tuhaf bir durum var: bu işlem için bir yürütme iş parçacığı belirlenmemişti. Bu durumda, varsayılan olarak ana iş parçacığının kullanılacağı beklenebilirdi. Ancak durum böyle değil. Yürütme iş parçacığı, [RxComputationThreadPool-3] hesaplama iş parçacığı olmuştur (11. satır). Dolayısıyla, yürütme veya gözlem iş parçacığı belirtilmediğinde, hangi iş parçacığının seçileceği konusunda bir varsayımda bulunulamaz.

7.5.3. Örnek-16: [Observable.zip] ile iki gözlemleneni birleştirmek

Şimdi şu kodu inceleyelim:


package dvp.rxjava.observables.exemples;

import java.util.Arrays;
import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.functions.FuncN;
import rx.schedulers.Schedulers;

public class Exemple16 {
    public static void main(String[] args) throws InterruptedException {
        // süreç 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // süreç 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, null);
        // 2 sürecin birleştirme işlevi
        FuncN<String> funcn = new FuncN<String>() {
            @Override
            public String call(Object... args) {
                if (args.length == 2) {
                    return String.format("double=%s, string=%s", args[0], args[1]);
                } else {
                    throw new RuntimeException("la fonction attend 2 paramètres exactement");
                }
            }
        };
        // 2 işlemin zip dosyası
        Process<String> process12 = new Process<>("process12",
                Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));
        // abonelikler
        ProcessUtils.subscribe(1, process12);
    }
}
  • 16-18. satırlar: [process1] adlı bir işlem, bir hesaplama iş parçacığı üzerinden 3 gerçek sayı üretecektir. Aynı zamanda bir hesaplama iş parçacığı üzerinden gözlemlenecektir;
  • 20-21. satırlar: [process2] adlı bir işlem, kısıtlanmamış bir iş parçacığı üzerinden 2 karakter dizisi gönderecektir. Gözlem iş parçacığı da kısıtlanmamıştır;
  • satır 23-32: [FuncN<String>] türünün anonim bir sınıfla örneklendirilmesi. FuncN işlevsel bir arayüzdür:
 

[FuncN.call] yöntemi bir nesne dizisi bekler ve bir R türü döndürür. [funcn] işlevi, process1 ve process2 işlemlerini bu sırayla birleştirmek için kullanılacaktır. [FuncN.call] yönteminde:

  • args[0], bir Double olacaktır;
  • args[1], String olacaktır;

Burada, [funcn.call]'in sonucu 27. satırdaki karakter dizisi olacaktır. Bu sonucun oluşturulması için call yönteminin argümanlarının türlerini bilmek gerekmez.

İki işlem şu şekilde birleştirilir:


// 2 işlemin zip dosyası
Process<String> process12 = new Process<>("process12",
Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));

[Observable.zip] yöntemi şu şekilde çalışır:

 

Görüldüğü gibi:

  • zip'in ilk argümanı bir Iterable<Observable>'tir. Örneğimizde, iki gözlemlenebilir değerimizden oluşan List<Observable> türünde bir etkin parametreye sahibiz;
  • zip'in ikinci argümanı bir FuncN türündedir. Örneğimizde, etkin parametre [funcn]'tir;

Çalıştırma sonucunda şu sonuçlar elde edilir:

main : début observation ------Thread[main] ---- Time[55:636]
Observable (process2) call start ------Thread[main] ---- Time[55:666]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[55:666]
Observable (process1,0) onNext (69.6) ------Thread[RxComputationThreadPool-4] ---- Time[55:902]
Observable (process2,0) onNext (valeur-0) ------Thread[main] ---- Time[56:076]
Observable (process1,1) onNext (82.8) ------Thread[RxComputationThreadPool-4] ---- Time[56:271]
Subscriber[observateur[0],process12] : onNext ("double=69.6, string=valeur-0") ------Thread[main] ---- Time[56:352]
Observable (process1,2) onNext (14.399999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[56:641]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[56:642]
Observable (process2,1) onNext (valeur-1) ------Thread[main] ---- Time[56:778]
Subscriber[observateur[0],process12] : onNext ("double=82.8, string=valeur-1") ------Thread[main] ---- Time[56:779]
Observable (process2) onCompleted ------Thread[main] ---- Time[56:779]
Subscriber[observateur[0],process12].onCompleted ------Thread[main] ---- Time[56:780]
main : attente fin observation ------Thread[main] ---- Time[56:781]
main : fin observation ------Thread[main] ---- Time[56:781]
  • 7. ve 11. satırlar: process12 süreci iki öğe üretir;
  • 8. satır: process1 işlemi tarafından gönderilen ve process2 işleminde eşi bulunmayan ek öğe, sonuç işlemi olan process12 tarafından gönderilmez;

Görüldüğü üzere, ne yürütme iş parçacığı ne de gözlem iş parçacığı atanmış olan process2 süreci, her ikisi için de ana iş parçacığını kullanmıştır.

7.5.4. Örnek-17: [Observable.combineLatest] ile iki gözlemleneni birleştirmek

Şimdi şu kodu inceleyelim:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple17 {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // 2. işlem
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null,
                Schedulers.computation());
        // 2 işlemin birleştirilmesi
        Process<Double> process12 = new Process<>("process12",
                Observable.combineLatest(process1.getObservable(), process2.getObservable(), (d1, d2) -> d1 + d2));
        // abonelikler
        ProcessUtils.subscribe(1, process12);
    }
}
  • 14-16. satırlar: [process1] adlı bir süreç, bir hesaplama iş parçacığı üzerinden 3 gerçek sayı üretecektir. Aynı süreç, bir hesaplama iş parçacığı üzerinden de gözlemlenecektir;
  • satır 18-20: [process2] adlı bir işlem, kısıtlanmamış bir iş parçacığı üzerinden 2 gerçek sayı üretecektir. Bu sayılar bir hesaplama iş parçacığı üzerinde gözlemlenecektir;
  • 23. satır: Her iki gözlemlenebilir değişken, aşağıdaki [Observable.combineLatest] statik yöntemi ile birleştirilir:
 

[combineLatest] gözlemlenebilir değeri şu şekilde çalışır: iki gözlemlenebilir değerden biri bir E1 öğesi yayınladığında, bu öğe [combineFunction] tarafından diğer gözlemlenebilir değerin yayınladığı son öğeyle birleştirilir.

Bu kodun çalıştırılması sonucunda şu çıktı elde edilir:

main : début observation ------Thread[main] ---- Time[01:768]
Observable (process2) call start ------Thread[main] ---- Time[01:791]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[01:791]
Observable (process1,0) onNext (54.0) ------Thread[RxComputationThreadPool-4] ---- Time[01:991]
Observable (process2,0) onNext (56.0) ------Thread[main] ---- Time[02:245]
Observable (process1,1) onNext (51.6) ------Thread[RxComputationThreadPool-4] ---- Time[02:358]
Subscriber[observateur[0],process12] : onNext (110.0) ------Thread[RxComputationThreadPool-5] ---- Time[02:521]
Subscriber[observateur[0],process12] : onNext (107.6) ------Thread[RxComputationThreadPool-5] ---- Time[02:522]
Observable (process2,1) onNext (261.8) ------Thread[main] ---- Time[02:595]
Observable (process2) onCompleted ------Thread[main] ---- Time[02:596]
main : attente fin observation ------Thread[main] ---- Time[02:596]
Subscriber[observateur[0],process12] : onNext (313.40000000000003) ------Thread[RxComputationThreadPool-5] ---- Time[02:597]
Observable (process1,2) onNext (80.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[02:790]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[02:791]
Subscriber[observateur[0],process12] : onNext (342.2) ------Thread[RxComputationThreadPool-3] ---- Time[02:792]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[02:792]
main : fin observation ------Thread[main] ---- Time[02:793]
  • 5. satır: process2 (56) tarafından gönderilen veri, process1 (54, 4. satır) tarafından gönderilen son öğeyle birleştirilir ve 7. satırdaki sonucu verir;
  • 6. satır: process1 (51,6) tarafından gönderilen veri, process2 (56, 5. satır) tarafından gönderilen son öğeyle birleştirilir ve 8. satırdaki sonucu verir;
  • 9. satır: process2 (261,8) tarafından gönderilen veri, process1 (51,6, 6. satır) tarafından gönderilen son öğeyle birleştirilir ve 12. satırdaki sonucu verir;
  • satır 13: process1 (80,39) tarafından gönderilen veri, process2 (261,8, satır 9) tarafından gönderilen son elemanla birleştirilir ve satır 15'teki sonucu verir;

Burada, [zip] gözlemlenebilirinin bir varyantıyla karşı karşıyayız; bu sefer birleştirilen öğeler, akışlardaki aynı konumdaki öğeler olmak zorunda değildir. Burada, yürütme iş parçacığı atanmamış olan process2 işleminin ana iş parçacığı üzerinde yürütüldüğü görülmektedir (satır 2).

7.5.5. Örnek-18: [Observable.amb] ile iki gözlemleneni birleştirmek

Şimdi şu kodu inceleyelim:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple18 {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // süreç 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null, null);
        // 2 işlemin birleştirilmesi
        Process<Double> process12 = new Process<>("process12",
                Observable.amb(process1.getObservable(), process2.getObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process12);
    }
}
  • 14-16. satırlar: [process1] adlı bir süreç, bir hesaplama iş parçacığı üzerinden 3 gerçek sayı üretecektir. Bu süreç aynı zamanda bir hesaplama iş parçacığı üzerinden gözlemlenecektir;
  • satır 18-20: [process2] adlı bir işlem, kısıtlanmamış bir iş parçacığı üzerinden 2 gerçek sayı üretecektir. Bu sayılar kısıtlanmamış bir iş parçacığı üzerinde gözlemlenecektir;
  • 22. satır: Her iki gözlemlenebilir değişken, aşağıdaki [Observable.amb] statik yöntemi ile birleştirilir:
 

Yukarıdaki şemada gösterildiği gibi, [Observable.amb(Observable o1, Observable o2)] gözlemlenebiliri, ilk olarak yayın yapan gözlemlenebilirin elemanlarını yayınlar. Bu durum, sunulan örneğin sonuçlarıyla da doğrulanmaktadır:

main : début observation ------Thread[main] ---- Time[21:594]
Observable (process2) call start ------Thread[main] ---- Time[21:612]
Observable (process1) call start ------Thread[RxComputationThreadPool-3] ---- Time[21:612]
Observable (process2,0) onNext (155.39999999999998) ------Thread[main] ---- Time[21:817]
Observable (process1) onError ------Thread[RxComputationThreadPool-3] ---- Time[21:820]
Observable (process1,0) onNext (90.0) ------Thread[RxComputationThreadPool-3] ---- Time[21:820]
Observable (process1,1) onNext (104.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[21:877]
Subscriber[observateur[0],process12] : onNext (155.39999999999998) ------Thread[main] ---- Time[22:105]
Observable (process1,2) onNext (44.4) ------Thread[RxComputationThreadPool-3] ---- Time[22:122]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[22:123]
Observable (process2,1) onNext (201.6) ------Thread[main] ---- Time[22:581]
Subscriber[observateur[0],process12] : onNext (201.6) ------Thread[main] ---- Time[22:583]
Observable (process2) onCompleted ------Thread[main] ---- Time[22:583]
Subscriber[observateur[0],process12].onCompleted ------Thread[main] ---- Time[22:584]
main : attente fin observation ------Thread[main] ---- Time[22:585]
main : fin observation ------Thread[main] ---- Time[22:586]
  • 4. satırda, ilk olarak process2 süreci yayınlar;
  • 8. ve 12. satırlar: process12 süreci, process2 süreci tarafından gönderilen tüm öğeleri (4. ve 11. satırlar) gönderir;

7.6. Bir gözlemlenebilirin işleme zinciri

7.6.1. Örnek-19: [Observable.map] ile bir gözlemleneni dönüştürme

Önceki örneklerde, iki gözlemlenebilirin bir üçüncü gözlemlenebilirde birleştirilmesine yönelik çeşitli kombinasyonları inceledik. Şimdi, bir gözlemlenebilir üzerinde dönüştürme, filtreleme ve toplama işlemlerine olanak tanıyan [Observable] sınıfının statik yöntemlerini sunacağız. Burada, 5. paragrafta incelediğimiz [Stream] sınıfındaki yöntemlere benzer yöntemlerle karşılaşacağız.

İlk örneğimiz şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple19 {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // süreç 2
        Process<String> process2 = new Process<>("process2",
                process1.getObservable().map(d -> String.format("valeur-%s", d)));
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}
  • 14-16. satırlar: process1 adlı bir süreç, bir hesaplama iş parçacığı üzerinden 3 gerçek sayı üretecektir. Bu süreç aynı zamanda bir hesaplama iş parçacığı üzerinden gözlemlenecektir;
  • 17-18. satırlar: process1 tarafından üretilen sayılar, process2 sürecinde karakter dizelerine dönüştürülecektir;
  • 20. satır: process2 izlenir;

18. satırdaki [Observable.map] yöntemi, 5.5. paragrafta incelenen [Stream.map] yöntemine benzerdir:

 

Örneğin sonuçları şunlardır:

main : début observation ------Thread[main] ---- Time[55:328]
main : attente fin observation ------Thread[main] ---- Time[55:346]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[55:347]
Observable (process1,0) onNext (21.599999999999998) ------Thread[RxComputationThreadPool-4] ---- Time[55:354]
Observable (process1,1) onNext (97.2) ------Thread[RxComputationThreadPool-4] ---- Time[55:512]
Subscriber[observateur[0],process2] : onNext ("valeur-21.599999999999998") ------Thread[RxComputationThreadPool-3] ---- Time[55:615]
Subscriber[observateur[0],process2] : onNext ("valeur-97.2") ------Thread[RxComputationThreadPool-3] ---- Time[55:616]
Observable (process1,2) onNext (98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[55:803]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[55:804]
Subscriber[observateur[0],process2] : onNext ("valeur-98.39999999999999") ------Thread[RxComputationThreadPool-3] ---- Time[55:804]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[55:805]
main : fin observation ------Thread[main] ---- Time[55:805]
  • 4., 5. ve 8. satırlar: process1'in çıktıları. Bunlar gerçek sayılardır;
  • 6., 7. ve 10. satırlar: gözlemlenen process2 emisyonları. Bunlar karakter dizileridir;

7.6.2. Örnek-20: [Observable.filter] ile bir gözlemleneni filtreleme

Örnek şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple20 {
    public static void main(String[] args) throws InterruptedException {
        // süreç 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}
  • 11-12. satırlar: process1 adlı bir işlem, bir hesaplama iş parçacığı üzerinde 0'dan 2'ye kadar tamsayılar üretecektir. Bu işlem aynı zamanda bir hesaplama iş parçacığı üzerinde gözlemlenecektir;
  • 14. satır: process1 tarafından üretilen sayılar filtrelenecek ve process2'te yalnızca çift sayılar tutulacaktır;
  • 20. satır: process2 izlenir;

18. satırdaki [Observable.filter] yöntemi, 5.4. paragrafta incelenen [Stream.filter] yöntemine benzerdir:

 

Örneğin sonuçları şunlardır:

main : début observation ------Thread[main] ---- Time[30:319]
main : attente fin observation ------Thread[main] ---- Time[30:335]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[30:336]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[30:388]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[30:625]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[30:703]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[30:704]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[30:705]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[30:706]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[30:707]
main : fin observation ------Thread[main] ---- Time[30:707]
  • 4., 5. ve 7. satırlar: process1'in yayınları;
  • 6. ve 9. satırlar: gözlemlenen process2 emisyonları. process1'in çift olan elemanlarıdır;

7.6.3. Örnek-21: [Observable.flatMap] ile bir gözlemleneni dönüştürme

Örnek şu şekilde olacaktır:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple21 {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }));
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}
  • 12-13. satırlar: process1 adlı bir işlem, bir hesaplama iş parçacığı üzerinde 0'dan 2'ye kadar tamsayılar üretecektir. Aynı zamanda bir hesaplama iş parçacığı üzerinde gözlemlenecektir;
  • 15-18. satırlar: process1 tarafından üretilen her n sayısı, (10*n, 10*n+1, 10*n+2) sayıları üreten bir gözlemlenebilir nesneye dönüştürülür. Eğer 15. satırda [map] yöntemi kullanılsaydı, process2, Integer türü yerine Observable<Integer> türünü üretirdi. Kullanılan [flatMap] yöntemi, (flatten) bu Observable<Integer> türündeki öğe dizisini, her bir Observable<Integer>'in her bir öğesinden oluşan bir Integer türündeki öğe dizisine düzleştirir;
  • 20. satır: process2 görülmektedir;

15. satırdaki [Observable.flatMap] yöntemi, 5.6.12. paragrafta incelenen [Stream.flatMap] yöntemine benzerdir:

 

Örneğin sonuçları şöyledir:

main : début observation ------Thread[main] ---- Time[31:466]
main : attente fin observation ------Thread[main] ---- Time[31:486]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[31:486]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[31:777]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[32:082]
Subscriber[observateur[0],process2] : onNext (1) ------Thread[RxComputationThreadPool-3] ---- Time[32:085]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[32:087]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[32:192]
Subscriber[observateur[0],process2] : onNext (10) ------Thread[RxComputationThreadPool-3] ---- Time[32:194]
Subscriber[observateur[0],process2] : onNext (11) ------Thread[RxComputationThreadPool-3] ---- Time[32:196]
Subscriber[observateur[0],process2] : onNext (12) ------Thread[RxComputationThreadPool-3] ---- Time[32:197]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[32:686]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[32:687]
Subscriber[observateur[0],process2] : onNext (20) ------Thread[RxComputationThreadPool-3] ---- Time[32:688]
Subscriber[observateur[0],process2] : onNext (21) ------Thread[RxComputationThreadPool-3] ---- Time[32:690]
Subscriber[observateur[0],process2] : onNext (22) ------Thread[RxComputationThreadPool-3] ---- Time[32:692]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[32:693]
main : fin observation ------Thread[main] ---- Time[32:693]
  • 5-7. satırlar: process1'in 4. satırının gönderilmesinin ardından gelen process2'in üç gönderimi;
  • 9-11. satırlar: process1'in 8. satırının gönderilmesinin ardından process2'in üç gönderimi;
  • satır 14-16: process1'in 12. satırının yayınlanmasının ardından process2'in üç yayını;

Aşağıdaki kod, process1 ve [Exemple21b]'ten Observable<Integer[]> türünün nasıl oluşturulacağını göstermektedir:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21b {
    public static void main(String[] args) throws InterruptedException {
        // süreç 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer[]> process2 = new Process<>("process2", process1.getObservable().map(i -> {
            int value = i * 10;
            return new Integer[] { value, value + 1, value + 2 };
        }));
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}
  • 14. satır: [Observable.map] yöntemi kullanılır;
  • 16. satır: Integer[] türünü döndürür;

Sonuçlar şöyledir:

main : début observation ------Thread[main] ---- Time[58:089]
main : attente fin observation ------Thread[main] ---- Time[58:107]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[58:108]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[58:503]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[58:762]
Subscriber[observateur[0],process2] : onNext ([0,1,2]) ------Thread[RxComputationThreadPool-3] ---- Time[58:792]
Subscriber[observateur[0],process2] : onNext ([10,11,12]) ------Thread[RxComputationThreadPool-3] ---- Time[58:795]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[58:851]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[58:852]
Subscriber[observateur[0],process2] : onNext ([20,21,22]) ------Thread[RxComputationThreadPool-3] ---- Time[58:853]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[58:854]
main : fin observation ------Thread[main] ---- Time[58:854]
  • 6., 7. ve 10. satırlar: map'in sonuçları görülmektedir;

Tüm bu gözlemlenebilir değişken dönüşümleri, her dönüşüm yeni bir gözlemlenebilir değişken ürettiği için zincirlenebilir. Aşağıdaki örnek [Exemple21c] bunu göstermektedir:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple21c {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }).filter(i -> i % 2 == 0));
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}
  • 15-18. satırlar: flatMap'in ardından bir filter gelir;

Çalıştırma sonuçları şöyledir:

main : début observation ------Thread[main] ---- Time[37:993]
main : attente fin observation ------Thread[main] ---- Time[38:016]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[38:017]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[38:124]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[38:366]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[38:380]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[38:381]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[38:436]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[38:439]
Subscriber[observateur[0],process2] : onNext (10) ------Thread[RxComputationThreadPool-3] ---- Time[38:441]
Subscriber[observateur[0],process2] : onNext (12) ------Thread[RxComputationThreadPool-3] ---- Time[38:443]
Subscriber[observateur[0],process2] : onNext (20) ------Thread[RxComputationThreadPool-3] ---- Time[38:445]
Subscriber[observateur[0],process2] : onNext (22) ------Thread[RxComputationThreadPool-3] ---- Time[38:446]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[38:447]
main : fin observation ------Thread[main] ---- Time[38:447]
  • satır 8-13: process2, flatMap'ten yalnızca çift numaralı öğeleri çıkarmıştır;

[flatMap]'e benzer bir yöntem, aşağıdaki [Exemple21d] örneğiyle gösterilen [flatMapIterable] yöntemidir:


package dvp.rxjava.observables.exemples;

import java.util.Arrays;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21d {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMapIterable(i -> {
            int value = i * 10;
            return Arrays.asList(value, value + 1, value + 2);
        }).filter(i -> i % 2 == 0));
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}

16. satırda, [flatMap] yöntemi yerine [flatMapIterable] yöntemi kullanılır. Bu durumda, dönüştürme işlevi Observable<T> türü yerine Iterable<T> türünü (18. satır) üretmelidir.

Öncekiyle aynı sonuçlar elde edilir.

[flatMap] yönteminin tanımına geri dönelim:

 

Yukarıda görüldüğü gibi, iki yeşil [1-2] öğesinin arasına mavi bir [3] öğesi eklenmiştir. Bu, Observable<T> öğelerini düzleştirme işlemi sırasında, [flatMap] yönteminin bu farklı iç gözlemlenebilir öğelerin yayın sırasını koruduğu anlamına gelir. Bu durum, aşağıdaki [Exemple21e] örneğiyle gösterilmektedir:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21e {
    public static void main(String[] args) throws InterruptedException {
        // süreç 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // işlem 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().flatMap(i -> process2.getObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process3);
    }
}
  • 11-12. satırlar: process1 işlemi, [0,1] tamsayılarını üretir;
  • satır 14-15: process2 süreci, [10,11,12] tamsayılarını üretir;
  • satır 17-18: process1 tarafından üretilen her bir öğeye, process2 sürecinin gözlemlenebilir değeri eşleştirilir. Bu şu anlama gelir:
    • process1'in [0] elemanına, [10,11,12] değerlerini üreten bir gözlemlenebilir atanacaktır;
    • aynı durum 1 numaralı eleman için de geçerlidir;

Sonuç olarak, 6 adet [10, 11, 12, 10, 11, 12] sayısı üretilecektir. Bunların hangi sırayla üretildiğini görmek istiyoruz.

Çalıştırma sonuçları şöyledir:

main : début observation ------Thread[main] ---- Time[22:540]
main : attente fin observation ------Thread[main] ---- Time[22:566]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[22:566]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[22:949]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[22:951]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[23:159]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[23:160]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[23:160]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[23:286]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[23:513]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[23:597]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[23:599]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[23:645]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-5] ---- Time[23:647]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[23:789]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[23:790]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-5] ---- Time[23:791]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[23:976]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[23:978]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[24:184]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[24:184]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[24:186]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[24:187]
main : fin observation ------Thread[main] ---- Time[24:187]

Görüldüğü gibi, process3 sürecinin yayın sırası şöyledir: [10, 10, 11, 12, 11, 12] (satır 11, 12, 14, 17, 19, 22). Dolayısıyla, process2 işlemi tarafından gönderilen öğeler arasında bir karışıklık olduğu açıktır. Bunu önlemek için, [flatMap] yöntemi yerine [concatMap] yöntemini kullanabilirsiniz. Aşağıdaki [Exemple21ef] kodu bunu göstermektedir:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21ef {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // işlem 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().concatMap(i -> process2.getObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process3);
    }
}

18. satırda, [flatMap], [concatMap] ile değiştirilmiştir. Yürütme sonuçları şöyledir:

main : début observation ------Thread[main] ---- Time[45:507]
main : attente fin observation ------Thread[main] ---- Time[45:530]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[45:530]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[45:775]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[45:778]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[45:846]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[45:890]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[45:947]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[45:948]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[46:096]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[46:097]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[46:144]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-5] ---- Time[46:147]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-5] ---- Time[46:148]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[46:149]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[46:364]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-7] ---- Time[46:366]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[46:529]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[46:531]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[46:558]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[46:559]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[46:560]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[46:562]
main : fin observation ------Thread[main] ---- Time[46:562]

Görüldüğü üzere, process3 işleminin gönderim sırası şöyledir: [10, 11, 12, 10, 11, 12] (12-14, 17, 19, 22. satırlar). process2 işlemi tarafından gönderilen öğeler karıştırılmamıştır.

[map] yönteminin bir başka varyantı da [switchMap] yöntemidir:

 

Yukarıda, [1] gözlemlenebilirinden, 2 elemanlı 3 başka [2] gözlemlenebiliri ortaya çıkar ve bunlar daha sonra [flatMap] ve [3]'te olduğu gibi düzleştirilir. Sonucun 6 değil 5 eleman içerdiği görülebilir. Bunun nedeni, ikinci gözlemlenebilirin 2 numaralı elemanı [6]'i yayınlamadan önce, üçüncü gözlemlenebilirin kendi ilk elemanı [5]'i yayınlamasıdır; bu da ikinci gözlemlenebilirin atılmasına neden olur. Dolayısıyla, [6] öğesi, [3] gözlemlenebilir sonuçta yer almamaktadır.

[switchMap]'i açıklamak için aşağıdaki [Exemple21eg] örneğini kullanacağız:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21eg {
    public static void main(String[] args) throws InterruptedException {
        // işlem 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // işlem 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().switchMap(i -> process2.getObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process3);
    }
}

Örneğin çalıştırılması şu sonuçları verir:

main : début observation ------Thread[main] ---- Time[02:388]
main : attente fin observation ------Thread[main] ---- Time[02:419]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[02:419]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[02:641]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[02:643]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[02:802]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[02:888]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[02:957]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[02:958]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[03:005]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[03:007]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[03:007]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[03:106]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[03:106]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[03:108]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[03:236]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[03:238]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[03:716]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[03:717]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[03:718]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[03:718]
main : fin observation ------Thread[main] ---- Time[03:719]
  • process1, 2 gözlemlenebilir process2'e yol açan 2 öğe yayar;
  • 14. satır: Gözlemci, 6. satırdaki 1. gözlemlenebilir process2 tarafından gönderilen 0 numaralı elemanı alır;
  • 15. satır: Gözlemci, 13. satırdaki 2. gözlemlenebilir process2 tarafından yayılan 0 numaralı öğeyi alır. Hikayede, gözlemcinin neden 1. gözlemlenebilir nesne process2 tarafından 7. ve 8. satırlarda gönderilen 1 ve 2 numaralı öğeleri daha önce almadığı belirtilmiyor. Her ne olursa olsun, 1. gözlemlenebilir nesne process2 terk edilir;
  • sonuç olarak, gözlemci gönderilen 6 öğe yerine sadece 4 öğeyi (satır 14, 15, 17, 20) görmektedir;

7.6.4. Örnekler-22: [Observable] sınıfının diğer yöntemleri

[Observable] sınıfı, [Stream] sınıfındaki birçok yöntemi benzer bir şekilde kullanır. İşte bunlardan birkaçı. Burada sadece kodu ve sonuçlarını veriyoruz.

[Exemple22a - take=limit]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22a {
    public static void main(String[] args) throws InterruptedException {
        // işlem
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).take(3));
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}

sonuçlar

1
2
3
4
5
6
7
main : début observation ------Thread[main] ---- Time[25:071]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[25:399]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[25:402]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[25:404]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[25:404]
main : attente fin observation ------Thread[main] ---- Time[25:406]
main : fin observation ------Thread[main] ---- Time[25:406]

[Exemple22b - takeLast]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22b {
    public static void main(String[] args) throws InterruptedException {
        // süreç
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).takeLast(2));
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}

sonuçlar

1
2
3
4
5
6
main : début observation ------Thread[main] ---- Time[19:440]
Subscriber[observateur[0],process] : onNext (9) ------Thread[main] ---- Time[19:726]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[19:728]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[19:728]
main : attente fin observation ------Thread[main] ---- Time[19:729]
main : fin observation ------Thread[main] ---- Time[19:730]

[Exemple22c - skip]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22c {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).skip(5).take(2));
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}

sonuçlar

1
2
3
4
5
6
main : début observation ------Thread[main] ---- Time[16:685]
Subscriber[observateur[0],process] : onNext (6) ------Thread[main] ---- Time[17:002]
Subscriber[observateur[0],process] : onNext (7) ------Thread[main] ---- Time[17:004]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[17:005]
main : attente fin observation ------Thread[main] ---- Time[17:006]
main : fin observation ------Thread[main] ---- Time[17:006]

[Exemple22d - reduce]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22d {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).reduce(0, (i, a) -> i + a));
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}
  • satır 10: gözlemlenebilirin elemanlarının toplamını hesaplar. Sonuç, bu toplamı veren bir gözlemlenebildir;

sonuçlar

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[52:412]
Subscriber[observateur[0],process] : onNext (55) ------Thread[main] ---- Time[52:640]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[52:640]
main : attente fin observation ------Thread[main] ---- Time[52:642]
main : fin observation ------Thread[main] ---- Time[52:642]

[Exemple22e - all]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22e {
    public static void main(String[] args) throws InterruptedException {
        // süreç
        Process<Boolean> process = new Process<>("process", Observable.range(1, 10).all(i -> i > 10));
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}
  • 10. satır: Observable<Boolean> döndürür; bu, true öğesini yayar; eğer [all] yönteminin yüklemi tüm öğeler için doğruysa, aksi takdirde false;

sonuçlar

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[59:866]
Subscriber[observateur[0],process] : onNext (false) ------Thread[main] ---- Time[00:069]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[00:070]
main : attente fin observation ------Thread[main] ---- Time[00:071]
main : fin observation ------Thread[main] ---- Time[00:071]

[Exemple22f - count]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22f {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).count());
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}
  • 10. satır: [Observable.count], gözlemlenen öğelerin toplamını içeren 1 öğeli bir gözlemlenebilir oluşturur;

sonuçlar

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[16:409]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[16:634]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[16:634]
main : attente fin observation ------Thread[main] ---- Time[16:635]
main : fin observation ------Thread[main] ---- Time[16:635]

[Exemple22g - distinct]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22g {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        Process<Integer> process = new Process<>("process", Observable.just(1, 2, 1, 3).distinct());
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}

sonuçlar

1
2
3
4
5
6
7
main : début observation ------Thread[main] ---- Time[05:373]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[05:594]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[05:595]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[05:596]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[05:597]
main : attente fin observation ------Thread[main] ---- Time[05:597]
main : fin observation ------Thread[main] ---- Time[05:597]

[Exemple22h - groupBy, asObservable]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.observables.GroupedObservable;

public class Exemple22h {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        Observable<GroupedObservable<Boolean, Integer>> obs = Observable.range(1, 10).groupBy(i -> i % 2 == 0);
        Process<Integer> process = new Process<>("process", obs.concatMap(g -> g.asObservable()));
        // abonelikler
        ProcessUtils.subscribe(1, process);
    }
}
  • 11. satır: [groupBy] yöntemi, gönderilen 10 öğeyi çift sayılar ve tek sayılar olmak üzere 2 gruba ayırır. Sonuç, Observable<GroupedObservable<Boolean, Integer>> türündedir; yani, elemanları GroupedObservable<Boolean, Integer> türünde olan bir gözlemlenebilir yapıdır; burada Boolean, grubun anahtarının türüdür (burada false, true) anahtarının türü olup, aynı zamanda [groupBy] yöntemine parametre olarak geçirilen lambda ifadesinin sonucunun türüdür; Integer ise grubun elemanlarının türüdür;
  • 12. satır: GroupedObservable türü, bu türden bir gözlemlenebilir oluşturmaya olanak tanıyan [asObservable] yöntemine sahiptir. Böylece, biri çift sayılar, diğeri tek sayılar için olmak üzere 2 adet Observable<Integer> türü elde edeceğiz. Bu iki gözlemlenebilir nesneden, [concatMap] yöntemi tek bir gözlemlenebilir nesne oluşturacaktır;

sonuçlar

main : début observation ------Thread[main] ---- Time[23:809]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[24:034]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[24:036]
Subscriber[observateur[0],process] : onNext (5) ------Thread[main] ---- Time[24:037]
Subscriber[observateur[0],process] : onNext (7) ------Thread[main] ---- Time[24:038]
Subscriber[observateur[0],process] : onNext (9) ------Thread[main] ---- Time[24:039]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[24:041]
Subscriber[observateur[0],process] : onNext (4) ------Thread[main] ---- Time[24:043]
Subscriber[observateur[0],process] : onNext (6) ------Thread[main] ---- Time[24:044]
Subscriber[observateur[0],process] : onNext (8) ------Thread[main] ---- Time[24:045]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[24:046]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[24:047]
main : attente fin observation ------Thread[main] ---- Time[24:047]
main : fin observation ------Thread[main] ---- Time[24:048]

[Exemple22i - timestamp]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;
import rx.schedulers.Timestamped;

public class Exemple22i {
    public static void main(String[] args) throws InterruptedException {
        // süreç 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // işlem 2
        Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
        // abonelikler
        ProcessUtils.subscribe(1, process2);
    }
}
  • 15. satırda, [timestamp] yöntemi, işlenen gözlemlenebilirin her bir elemanına bir saat atar;

sonuçlar

main : début observation ------Thread[main] ---- Time[59:362]
main : attente fin observation ------Thread[main] ---- Time[59:377]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[59:378]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[59:553]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[59:692]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975259555,"value":0}) ------Thread[RxComputationThreadPool-3] ---- Time[59:789]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975259789,"value":1}) ------Thread[RxComputationThreadPool-3] ---- Time[59:791]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[00:025]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[00:027]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975260026,"value":2}) ------Thread[RxComputationThreadPool-3] ---- Time[00:031]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[00:033]
main : fin observation ------Thread[main] ---- Time[00:034]

Bu örnekte, timestamp bilgisinin neyi temsil ettiğini söylemek zordur:

  • 4-5. satırlar: process1'in 1. elemanının, 0. elemandan 139 ms sonra gönderildiği görülmektedir;
  • 6. ve 7. satırlar: process2'in 1. elemanının, 0. elemandan 234 ms sonra gözlemlendiği görülüyor;
  • 5. ve 8. satırlar: process1'in 2. öğesinin, 1. öğeden 33 ms sonra gönderildiği görülmektedir;
  • 7. ve 10. satırlar: process2'in 2. elemanının, 1. elemandan 37 ms sonra gözlemlendiği görülmektedir;

Bu gecikmeler, gözlem ve gözlemlenebilirlerin yürütme iş parçacıklarının aynı olmaması nedeniyle ortaya çıkmaktadır. 12-13. satırları aşağıdaki satırlarla değiştirirsek (Örnek22j):


// süreç 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
  • 2-3. satırlar: gözlem iş parçacığı zorlanmaz. Bu durumda gözlemlenebilirin, yürütüldüğü yerde gözlemlendiği bilinmektedir;

Bu, aşağıdaki sonuçları verir:

main : début observation ------Thread[main] ---- Time[43:834]
main : attente fin observation ------Thread[main] ---- Time[43:845]
Observable (process1) call start ------Thread[RxComputationThreadPool-1] ---- Time[43:846]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-1] ---- Time[44:291]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976384293,"value":0}) ------Thread[RxComputationThreadPool-1] ---- Time[44:552]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-1] ---- Time[44:878]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976384879,"value":1}) ------Thread[RxComputationThreadPool-1] ---- Time[44:884]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-1] ---- Time[45:274]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976385275,"value":2}) ------Thread[RxComputationThreadPool-1] ---- Time[45:280]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[45:281]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[45:283]
main : fin observation ------Thread[main] ---- Time[45:284]
  • 4. ve 6. satırlar: process1 süreci, 0 numaralı öğesinden 587 ms sonra 1 numaralı öğesini gönderir;
  • 5. ve 7. satırlar: Gözlemci bu iki öğeyi 586 ms'lik bir aralıkla gözlemler;
  • 6. ve 8. satırlar: process1 süreci, 1 numaralı öğesinden 396 ms sonra 2 numaralı öğesini gönderir;
  • 7. ve 9. satırlar: Gözlemci bu iki öğeyi 396 ms'lik bir aralıkla gözlemler;

Burada, timestamp'in değerleri tutarlıdır: bunlar, öğenin gönderilme tarihini doğru bir şekilde yansıtmaktadır.

7.7. Zamanlayıcılar

7.7.1. Örnek-23: [Schedulers.computation] zamanlayıcısı

Şimdi yürütme zamanlayıcılarını inceleyeceğiz. Gözlem, yürütme iş parçacığı üzerinde yapılacaktır.

Zamanlayıcılar konusu biraz karmaşıktır. Farklı zamanlayıcılar, StackOverflow [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases] sitesindeki bu soruda tanıtılmaktadır:

 

Bu farklı zamanlayıcıların kullanımını örneklerle açıklamaya çalışacağız. İlk örnek, [Schedulers.computation] zamanlayıcısını göstermektedir:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple23 {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.computation(), null);
        }
        // abonelikler
        ProcessUtils.subscribe(1, processes);
    }
}
  • 14-19. satırlar: bir hesaplama iş parçacığı üzerinde çalışan 10 işlemden oluşan bir dizi oluşturulur;
  • 17. satır: her işlem rastgele bir gerçek sayı üretir;
  • 21. satır: tüm bu süreçlere abone olunur;

Sonuçlar şöyledir:

main : début observation ------Thread[main] ---- Time[01:034]
Observable (process0) call start ------Thread[RxComputationThreadPool-1] ---- Time[01:042]
Observable (process2) call start ------Thread[RxComputationThreadPool-3] ---- Time[01:042]
Observable (process1) call start ------Thread[RxComputationThreadPool-2] ---- Time[01:042]
Observable (process5) call start ------Thread[RxComputationThreadPool-6] ---- Time[01:043]
Observable (process7) call start ------Thread[RxComputationThreadPool-8] ---- Time[01:043]
Observable (process4) call start ------Thread[RxComputationThreadPool-5] ---- Time[01:042]
Observable (process3) call start ------Thread[RxComputationThreadPool-4] ---- Time[01:042]
main : attente fin observation ------Thread[main] ---- Time[01:043]
Observable (process6) call start ------Thread[RxComputationThreadPool-7] ---- Time[01:043]
Observable (process3,0) onNext (70.8) ------Thread[RxComputationThreadPool-4] ---- Time[01:115]
Observable (process1,0) onNext (13.2) ------Thread[RxComputationThreadPool-2] ---- Time[01:153]
Observable (process0,0) onNext (63.599999999999994) ------Thread[RxComputationThreadPool-1] ---- Time[01:215]
Subscriber[observateur[0],process0] : onNext (63.599999999999994) ------Thread[RxComputationThreadPool-1] ---- Time[01:326]
Subscriber[observateur[0],process3] : onNext (70.8) ------Thread[RxComputationThreadPool-4] ---- Time[01:326]
Subscriber[observateur[0],process1] : onNext (13.2) ------Thread[RxComputationThreadPool-2] ---- Time[01:326]
Observable (process3) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[01:326]
Observable (process0) onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[01:326]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[01:327]
Subscriber[observateur[0],process0].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[01:327]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[01:327]
Subscriber[observateur[0],process1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[01:327]
Observable (process8) call start ------Thread[RxComputationThreadPool-1] ---- Time[01:329]
Observable (process9) call start ------Thread[RxComputationThreadPool-2] ---- Time[01:329]
...
main : fin observation ------Thread[main] ---- Time[01:610]
  • 2-10. satırlar: İlk 8 işlem, 8 farklı iş parçacığı üzerinde başlatılır (kullanılan makine 8 çekirdeğe sahiptir). Hepsinin yaklaşık olarak aynı anda başladığı görülebilir;
  • satır 17-19: 3 işlem sona erer ve böylece 3 iş parçacığı serbest kalır;
  • 23-24. satırlar: Son iki işlem, serbest kalan iş parçacıklarından 2'sini kullanarak başlatılabilir;

Dolayısıyla, [Schedulers.computation] zamanlayıcısının, n iş parçacığı havuzu sağladığını unutmayalım; burada n, makinenin çekirdek sayısıdır. İş parçacıkları bu çekirdekler üzerinde paralel olarak yürütülür.

7.7.2. Örnek-24: [Schedulers.io] zamanlayıcı

Önceki kodu [Schedulers.io] zamanlayıcısıyla çalıştırıyoruz:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple24 {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.io(), null);
        }
        // abonelikler
        ProcessUtils.subscribe(1, processes);
    }
}
  • 18. satır: İşlemler, [Schedulers.io] zamanlayıcısının iş parçacıklarıyla yürütülür;

Bu, aşağıdaki sonuçları verir:

main : début observation ------Thread[main] ---- Time[03:451]
Observable (process0) call start ------Thread[RxCachedThreadScheduler-1] ---- Time[03:459]
Observable (process1) call start ------Thread[RxCachedThreadScheduler-2] ---- Time[03:459]
Observable (process2) call start ------Thread[RxCachedThreadScheduler-3] ---- Time[03:460]
Observable (process3) call start ------Thread[RxCachedThreadScheduler-4] ---- Time[03:460]
Observable (process4) call start ------Thread[RxCachedThreadScheduler-5] ---- Time[03:464]
Observable (process5) call start ------Thread[RxCachedThreadScheduler-6] ---- Time[03:464]
Observable (process6) call start ------Thread[RxCachedThreadScheduler-7] ---- Time[03:465]
Observable (process8) call start ------Thread[RxCachedThreadScheduler-9] ---- Time[03:465]
Observable (process9) call start ------Thread[RxCachedThreadScheduler-10] ---- Time[03:465]
main : attente fin observation ------Thread[main] ---- Time[03:465]
Observable (process7) call start ------Thread[RxCachedThreadScheduler-8] ---- Time[03:465]
Observable (process7,0) onNext (54.0) ------Thread[RxCachedThreadScheduler-8] ---- Time[03:473]
Observable (process8,0) onNext (116.39999999999999) ------Thread[RxCachedThreadScheduler-9] ---- Time[03:500]
Observable (process6,0) onNext (105.6) ------Thread[RxCachedThreadScheduler-7] ---- Time[03:506]
Observable (process0,0) onNext (96.0) ------Thread[RxCachedThreadScheduler-1] ---- Time[03:509]
Observable (process5,0) onNext (25.2) ------Thread[RxCachedThreadScheduler-6] ---- Time[03:583]
Observable (process3,0) onNext (97.2) ------Thread[RxCachedThreadScheduler-4] ---- Time[03:684]
Subscriber[observateur[0],process7] : onNext (54.0) ------Thread[RxCachedThreadScheduler-8] ---- Time[03:685]
Subscriber[observateur[0],process6] : onNext (105.6) ------Thread[RxCachedThreadScheduler-7] ---- Time[03:685]
Subscriber[observateur[0],process0] : onNext (96.0) ------Thread[RxCachedThreadScheduler-1] ---- Time[03:685]
Subscriber[observateur[0],process8] : onNext (116.39999999999999) ------Thread[RxCachedThreadScheduler-9] ---- Time[03:685]
Observable (process0) onCompleted ------Thread[RxCachedThreadScheduler-1] ---- Time[03:686]
Observable (process6) onCompleted ------Thread[RxCachedThreadScheduler-7] ---- Time[03:686]
Observable (process7) onCompleted ------Thread[RxCachedThreadScheduler-8] ---- Time[03:685]
...
main : fin observation ------Thread[main] ---- Time[03:933]
  • 2-10. satırlar: 10 işlem, her biri farklı bir iş parçacığı üzerinde başlatılır. Önceki durumun aksine, tüm işlemler başarıyla başlatılabilmiştir. Bu başlatma işlemlerinin 6 ms sürdüğü görülürken, önceki durumda bu süre 1 ms idi;
  • 13-18. satırlar: Gözlemlenebilirler, önceki örnekte olduğu gibi neredeyse paralel olarak değil, birbiri ardına sinyal gönderiyor;

[Schedulers.io] ve [Schedulers.computation] zamanlayıcıları arasındaki fark nedir? Cevap, URL ve [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]'te bulunabilir:

 

7.7.3. Örnek-25: [Schedulers.newThread] zamanlayıcısı

Önceki kodu [Schedulers.newThread] zamanlayıcısıyla çalıştırıyoruz:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple25 {
    public static void main(String[] args) throws InterruptedException {
        // işlemler
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.newThread(), null);
        }
        // abonelikler
        ProcessUtils.subscribe(1, processes);
    }
}

Elde edilen sonuçlar, [Schedulers.io] zamanlayıcıyla elde edilen sonuçlarla aynıdır:

main : début observation ------Thread[main] ---- Time[17:058]
Observable (process0) call start ------Thread[RxNewThreadScheduler-1] ---- Time[17:065]
Observable (process1) call start ------Thread[RxNewThreadScheduler-2] ---- Time[17:065]
Observable (process2) call start ------Thread[RxNewThreadScheduler-3] ---- Time[17:066]
Observable (process3) call start ------Thread[RxNewThreadScheduler-4] ---- Time[17:066]
Observable (process4) call start ------Thread[RxNewThreadScheduler-5] ---- Time[17:068]
Observable (process5) call start ------Thread[RxNewThreadScheduler-6] ---- Time[17:069]
Observable (process6) call start ------Thread[RxNewThreadScheduler-7] ---- Time[17:069]
Observable (process8) call start ------Thread[RxNewThreadScheduler-9] ---- Time[17:069]
Observable (process7) call start ------Thread[RxNewThreadScheduler-8] ---- Time[17:069]
Observable (process9) call start ------Thread[RxNewThreadScheduler-10] ---- Time[17:069]
main : attente fin observation ------Thread[main] ---- Time[17:069]
Observable (process6,0) onNext (25.2) ------Thread[RxNewThreadScheduler-7] ---- Time[17:120]
Observable (process3,0) onNext (39.6) ------Thread[RxNewThreadScheduler-4] ---- Time[17:193]
Observable (process5,0) onNext (21.599999999999998) ------Thread[RxNewThreadScheduler-6] ---- Time[17:212]
Observable (process0,0) onNext (19.2) ------Thread[RxNewThreadScheduler-1] ---- Time[17:273]
Observable (process8,0) onNext (81.6) ------Thread[RxNewThreadScheduler-9] ---- Time[17:308]
Subscriber[observateur[0],process3] : onNext (39.6) ------Thread[RxNewThreadScheduler-4] ---- Time[17:331]
Subscriber[observateur[0],process0] : onNext (19.2) ------Thread[RxNewThreadScheduler-1] ---- Time[17:331]
Subscriber[observateur[0],process6] : onNext (25.2) ------Thread[RxNewThreadScheduler-7] ---- Time[17:331]
Subscriber[observateur[0],process8] : onNext (81.6) ------Thread[RxNewThreadScheduler-9] ---- Time[17:331]
Subscriber[observateur[0],process5] : onNext (21.599999999999998) ------Thread[RxNewThreadScheduler-6] ---- Time[17:331]
Observable (process8) onCompleted ------Thread[RxNewThreadScheduler-9] ---- Time[17:333]
Observable (process5) onCompleted ------Thread[RxNewThreadScheduler-6] ---- Time[17:333]
Observable (process6) onCompleted ------Thread[RxNewThreadScheduler-7] ---- Time[17:332]
Observable (process0) onCompleted ------Thread[RxNewThreadScheduler-1] ---- Time[17:332]
Observable (process3) onCompleted ------Thread[RxNewThreadScheduler-4] ---- Time[17:332]
...
main : fin observation ------Thread[main] ---- Time[17:571]

URL ve [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io] bölümlerinde, [Schedulers.io] zamanlayıcısının bir iş parçacığı havuzu sağladığı, ancak [Schedulers.newThread] zamanlayıcısının bunu sağlamadığı açıklanmaktadır. Bir iş parçacığı havuzu, otomatik olarak n sayıda iş parçacığı oluşturur. Bunları, ihtiyaç duyan işlemlere tahsis eder. İşlemler tamamlandığında, iş parçacıkları silinmez, havuza geri döner ve başka bir işlem tarafından yeniden kullanılabilir. Bu, iş parçacıklarını sürekli olarak oluşturup silmekten daha verimlidir. Dolayısıyla, [Schedulers.io] zamanlayıcısını kullanmanın daha uygun olduğu düşünülebilir.

7.7.4. Örnek-26: [Schedulers.immediate, Schedulers.trampoline] zamanlayıcıları

Bu iki zamanlayıcı için verilen açıklamaya geri dönelim:

 

Açıklama anlaşılması oldukça basit olsa da, bunu örnekle göstermeye çalıştığımızda aslında tam olarak anlamadığımızı fark ediyoruz. [Learning Reactive Programming With Java 8] kitabı sayesinde, bu kitapta bulunan bir örneği temel alarak onu basitleştiren bir örnek oluşturabildim. Örnek şu şekildedir:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Scheduler;
import rx.Scheduler.Worker;
import rx.functions.Action0;
import rx.schedulers.Schedulers;

public class Exemple26 {
    public static void main(String[] args) throws InterruptedException {

        // bir zamanlayıcı
        Scheduler scheduler = Schedulers.immediate();
        // bu zamanlayıcının bir işleyicisi
        Worker worker = scheduler.createWorker();
        // işçide yürütülecek Action0 türü
        Action0 action02 = new Action0() {
            @Override
            public void call() {
                // Action02 günlüğü
                ProcessUtils.showInfos.accept("action02");
            }
        };

        // işçide yürütülecek bir Action0 türü
        Action0 action01 = new Action0() {
            @Override
            public void call() {
                // aynı işleyicide yeni bir eylem programlanıyor
                worker.schedule(action02);
                // action01 günlüğü
                ProcessUtils.showInfos.accept("action01");
            }
        };
        // action01, işleyicide programlandı
        worker.schedule(action01);
    }

    // görüntülemeler
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));

}
  • 17. satır: bir zamanlayıcı. Bu, ya burada olduğu gibi [Schedulers.immediate] ya da daha sonra [Schedulers.trampoline] olacaktır;
  • 19. satır: Zamanlayıcının işleyicileri üzerinde Action0 türündeki eylemleri (21. ve 20. satırlar) çalıştırabiliriz. [Scheduler.createWorker] yöntemi, bir işçi oluşturmaya olanak tanır. [Worker.schedule(Action0)] yöntemi, bir işçi tarafından Action0 türünde bir eylemin yürütülmesini sağlar;
  • 21-27. satırlar: [action02] adlı ilk eylem, 19. satırdaki işçi tarafından (40. satırda) yürütülecektir;
  • satır 30-38: [action01] adlı ikinci bir eylem. Bu eylemin özelliği, action02 eylemini kendisiyle aynı işçi üzerinde çalıştırmasıdır (satır 34). [Schedulers.immediate] ile [Schedulers.trampoline] arasındaki fark işte burada yatmaktadır:
    • eğer zamanlayıcı [Schedulers.immediate] ise, 34. satırda action02 eylemi hemen çalıştırılır (zamanlayıcının adı da buradan gelir) ve o anda yürütülmekte olan action01 eylemi kesilir. Bunun üzerine 25. satırdaki mesaj görüntülenecektir. action02 eylemi tamamlandığında, action01 eylemi yeniden başlayacak ve 36. satırdaki mesaj görüntülenecektir;
    • eğer zamanlayıcı [Schedulers.trampoline] ise, 34. satırda action02 eylemi beklemeye alınır. Bu eylem, yalnızca devam eden action01 görevi tamamlandığında yürütülecektir. Böylece 36. satırdaki mesaj görünecektir. action01 eylemi tamamlandığında, action02 eylemi yürütülecek ve 25. satırdaki mesaj görünecektir;

Yukarıdaki kodun çalıştırılması sonucunda şu sonuçlar elde edilir:

action02 ------Thread[main] ---- Time[38:480]
action01 ------Thread[main] ---- Time[38:485]

17. satırda [Schedulers.trampoline] zamanlayıcısını kullanırsak, tam tersi sonuçlar elde edilir:

action01 ------Thread[main] ---- Time[42:972]
action02 ------Thread[main] ---- Time[42:976]

Bununla birlikte, gözlemlenebilirlerle bir bağlantı kurmak zordur. Bu iki iş parçacığından birinde bir gözlemlenebilirin çalıştırılmasının yararını gösterebilecek ikna edici bir örnek bulamadım. Yine de işte bir örnek, ancak bunu hiç de doğal bulmuyorum:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.Scheduler.Worker;
import rx.functions.Action1;
import rx.schedulers.Schedulers;

public class Exemple27 {
    public static void main(String[] args) throws InterruptedException {

        // İşçi
        Worker worker = Schedulers.immediate().createWorker();
        // İşçi işçi = Schedulers.trampoline().createWorker();
        // işçi üzerinde gözlemlenebilir 1
        worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {

            @Override
            public void call(Integer i) {
                ProcessUtils.showInfos.accept(String.valueOf(i));
                // aynı işçi üzerinde gözlemlenebilir 2
                worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
                    @Override
                    public void call(Integer i) {
                        ProcessUtils.showInfos.accept(String.valueOf(i));
                    }
                }));
            }
        }));
    }
}
  • 13-14. satırlar: [Schedulers.immediate] ve [Schedulers.trampoline] adlı iki zamanlayıcıdan birinden bir işçi oluşturulur;
  • 16. satır: obs1 adlı ilk gözlemlenebilir, [1,2] numaralarını üretmek üzere bu işleyiciye programlanıyor
  • 22. satır: Bu obs1 gözlemlenebilirinin bir elemanı her gözlemlendiğinde, aynı işçi üzerinde ikinci bir gözlemlenebilir olan obs2'in gözlemi başlatılır ve [100,101] sayıları üretilir;

[Schedulers.immediate] zamanlayıcısıyla şu sonuçlar elde edilir:

1
2
3
4
5
6
1 ------Thread[main] ---- Time[44:604]
100 ------Thread[main] ---- Time[44:610]
101 ------Thread[main] ---- Time[44:610]
2 ------Thread[main] ---- Time[44:612]
100 ------Thread[main] ---- Time[44:612]
101 ------Thread[main] ---- Time[44:612]

Oysa [Schedulers.trampoline] zamanlayıcısıyla şu sonuçlar elde edilir:

1
2
3
4
5
6
1 ------Thread[main] ---- Time[14:107]
2 ------Thread[main] ---- Time[14:114]
100 ------Thread[main] ---- Time[14:115]
101 ------Thread[main] ---- Time[14:115]
100 ------Thread[main] ---- Time[14:115]
101 ------Thread[main] ---- Time[14:116]

7.8. Conclusion

Hala yapılacak çok iş var. RxJava kütüphanesini daha derinlemesine incelemek için, okuyucunun bu belgenin başında verilen kaynaklara başvurarak eğitimine devam etmesi önerilir. Yine de, Swing ve Android ortamlarında RxJava'i kullanmak için gerekli temellere sahibiz. Şimdi bunu göstereceğiz.