Skip to content

7. کتابخانه RxJava

کتابخانه RxJava بر اساس مفهوم زیر است: یک جریان از عناصر از نوع T که با Observable<T> مشاهده می‌شود توسط یک یا چند مشترک (مشترک‌ها، ناظران، مصرف‌کنندگان) از نوع Subscriber<T> دنبال می‌شود. کتابخانه RxJava این امکان را فراهم می‌کند که جریان Observable<T> در یک نخ T1 و ناظر Subscriber<T> آن در یک نخ T2 اجرا شود، بدون اینکه توسعه‌دهندهنیازی به نگرانی در مورد مدیریت چرخه عمر این نخ‌ها یا رسیدگی به مسائل ذاتاً دشوار، مانند اشتراک‌گذاری داده بین نخ‌ها و همگام‌سازی آن‌ها برای اجرای یک وظیفه کلی، ندارد. بنابراین، این امر برنامه‌نویسی ناهمزمان را تسهیل می‌کند.

یک جریان Observable<T> عناصر از نوع T را تولید می‌کند که می‌توان آن‌ها را هر زمان که تولید می‌شوند مشاهده کرد. اگر ناظر و قابل مشاهده (اصطلاحی کلی برای اشاره به نوع Observable<T>) در یک نخ باشند، آنگاه قابل مشاهده تنها پس از مصرف عنصر i توسط ناظر، می‌تواند عنصر (i+1) را تولید کند. موارد اندکی وجود دارد که این معماری کاربرد دارد. اگر ناظر و قابل مشاهده در یک نخ نباشند، آنگاه قابل مشاهده و ناظرش به طور مستقل عمل می‌کنند: قابل مشاهده با سرعت خود عناصر را منتشر می‌کند و ناظر با سرعت خود آن‌ها را مصرف می‌کند. ارزش کتابخانه در همین‌جاست. تا به اینجا، ما همیشه به یک ناظر واحد اشاره کرده‌ایم. در واقع، یک قابل مشاهده می‌تواند هر تعداد ناظری داشته باشد.

کتابخانه RxJava به‌ویژه برای معماری توصیف‌شده در بند ۲ مقدمه مناسب است که در اینجا خلاصه شده است:

Image

  • در [1]، یک لایه خدماتی، خدماتی را ارائه می‌دهد که برخی از آن‌ها زمان زیادی برای دریافت نیاز دارند (مانند درخواست‌های شبکه)؛
  • این لایه خدماتی توسط یک رابط کاربری گرافیکی [1] (Swing، Android، JavaFx) فراخوانی می‌شود. اگر لایه سرویس در همان تِردِ متد [swing] که از آن استفاده می‌کند اجرا شود، رابط کاربری گرافیکی هنگام انتظار برای نتیجه سرویس، فریز می‌شود (غیرپاسخگو می‌گردد)؛
  • در [2]، یک لایه تطبیق نازک که با استفاده از RxJava پیاده‌سازی شده است، امکان ارائه پیاده‌سازی ناهمزمان همان سرویس به لایه رابط کاربری گرافیکی را فراهم می‌کند: این سرویس می‌تواند در یک نخ (thread) متفاوت از نخ متد لایه رابط کاربری گرافیکی که آن را فراخوانی می‌کند، اجرا شود. در این حالت، رابط گرافیکی [3] همچنان پاسخگو باقی می‌ماند: کاربر می‌تواند تعامل با آن را ادامه دهد، برای مثال با راه‌اندازی یک درخواست شبکه‌ای جدید به موازات درخواست اول؛ و نکته بسیار مهم این است که به کاربر این امکان داده می‌شود که فرآیندهایی را که بیش از حد طول می‌کشند لغو کند – کاری که اگر رابط گرافیکی منجمد شده بود، غیرممکن بود؛
  • فراخوانی [4] همگام است، در حالی که فراخوانی [5-6] غیرهمگام است؛

در این معماری، لایه [2] خدماتی را ارائه می‌دهد که انواع Observable<T> را بازمی‌گردانند، که متدهای لایه گرافیکی [3] می‌توانند به آن‌ها مشترک شوند. سپس یک سرویس در لایه [2] نتایج خود را یکی‌یکی ارائه می‌دهد و لایه [3] می‌تواند به هر یک از آن‌ها واکنش نشان دهد، برای مثال با به‌روزرسانی یک یا چند مؤلفه رابط کاربری گرافیکی.

کلاس </span>**Observable&lt;T&gt;** ده‌ها متد دارد. این یکی از چالش‌های این کتابخانه است: بسیار جامع است و درک تمام امکانات آن دشوار است. ما در اینجا برخی از آنها را معرفی خواهیم کرد. تسلط بر سایر متدها به مرور زمان حاصل خواهد شد.

7.1. ایجاد ناظرها و مشترک شدن در آن‌ها

7.1.1. مثال-۰۱: متد [Observable.from]

  

کد زیر را در نظر بگیرید:


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) {
    // متغیرهای عددی قابل مشاهده
    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");
      }
    });
  }
}
  • خط ۱۲: یک نوع Observable<Integer> از یک لیست اعداد صحیح ایجاد می‌شود.

کلاس Observable<T> یک جریان از عناصر با نوع T است که می‌توان آن‌ها را – ترجیحاً به‌صورت ناهمزمان، اما لزوماً نه – در حین تولید مشاهده کرد. تعریف آن به شرح زیر است:

 

همان‌طور که قبلاً ذکر شد، کلاس Observable<T> دارای چندین دوجین متد است. برخی از آن‌ها مشابه متدهای کلاس Stream<T> هستند که در بند ۵ مورد بحث قرار گرفته‌اند. مستندات RxJava شامل «نمودارهای مرواریدی» ([2]) است که نحوه عملکرد این متدها را نشان می‌دهد:

  • خط ۳ انتشارها را از مشاهده‌پذیر در طول زمان نشان می‌دهد؛
  • روش [4] بر روی عناصر منتشر شده توسط مشاهده‌پذیر اعمال می‌شود. این روش عموماً یک مشاهده‌پذیر جدید تولید می‌کند؛
  • خط ۵ مشاهده‌پذیر جدید به‌دست‌آمده را نشان می‌دهد؛

روش [Observable.from] دارای امضای زیر است:

 

متد استاتیک [Observable.from] به شما امکان می‌دهد یک Observable<T> را از یک مجموعه عناصر از نوع T ایجاد کنید. این یک روش بسیار ساده برای شروع کار با observables است. خط:


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

بنابراین سه عنصر را منتشر خواهد کرد. این کار را فوراً انجام نمی‌دهد. هر بار که یک مشترک ثبت‌نام کند، آن‌ها را به‌طور کامل منتشر خواهد کرد. این به عنوان یک مشاهده‌پذیر سرد (cold observable) شناخته می‌شود. مشاهده‌پذیر برای هر مشترک جدید، عناصر خود را مجدداً منتشر می‌کند.

می‌توانیم عبارت قبلی را به عنوان یک عمل پیکربندی برای مشاهده‌پذیر در نظر بگیریم. این عمل یک‌بار پیکربندی می‌شود و در صورت ثبت n مشترک، n بار اجرا می‌شود.

چگونه مشترک می‌شوید؟

یک راه برای انجام این کار استفاده از متد [Observable.subscribe] است که تعریف آن به شرح زیر است:

 
  • پارامتر اول [Action1<T> onNext] (به بخش 6.2 مراجعه کنید) این متد، متدی است که باید زمانی که قابل مشاهده عنصر جدید T را منتشر می‌کند، اجرا شود؛
  • پارامتر دوم، [Action1<Throwable> onError]، از این متد، متدی است که باید هنگام پرتاب یک استثنا توسط مشاهده‌پذیر اجرا شود؛
  • پارامتر سوم [Action0 onComplete] (به بخش 6.1 مراجعه کنید) از متد، متدی است که باید هنگام پرتاب یک استثنا توسط مشاهده‌پذیر اجرا شود؛
  • این متد یک نوع [Subscription] را برمی‌گرداند؛

نوع [Subscription] نمایانگر یک اشتراک به مشاهده‌پذیر است. تعریف آن به شرح زیر است:

 

مزیت این رابط [1] در متد [2] آن نهفته است که امکان لغو اشتراک را فراهم می‌کند.

در مثال ما، کد برای اشتراک‌پذیری قابل مشاهده به شرح زیر است:


    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");
      }
});
  • خط ۱: نتیجه با نوع [Subscription] نادیده گرفته می‌شود؛
  • خطوط ۱–۱۵: سه پارامتر نمونه‌هایی از کلاس‌های ناشناس هستند. ما همچنین از لامبداها استفاده خواهیم کرد. مزیت کلاس‌های ناشناس این است که به وضوح انواع داده‌ای مورد انتظار توسط تنها متد این کلاس‌ها را نشان می‌دهند؛
  • خطوط ۲–۵: پیاده‌سازی پارامتر اول از نوع [Action1<Integer>];
  • خطوط ۶–۱۰: پیاده‌سازی پارامتر دوم از نوع [Action1<Throwable>];
  • خطوط ۱۱–۱۵: پیاده‌سازی پارامتر سوم از نوع [Action0];

کد کامل به شرح زیر است:


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) {
    // متغیرهای عددی قابل مشاهده
    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");
      }
    });
  }
}

مشاهده‌پذیر در خط ۱۲ بلافاصله پس از فراخوانی متد [subscribe] در خط ۱۴، شروع به ارسال سه عنصر خود می‌کند. از آن نقطه به بعد:

  • هرگاه یک عنصر صادر شود، خطوط ۱۵–۱۸ اجرا می‌شوند.
  • پس از اینکه هر سه عنصر ارسال شدند، خطوط ۲۴ تا ۲۹ اجرا می‌شوند؛
  • خطوط ۱۹–۲۴ هرگز اجرا نخواهند شد زیرا مشاهده‌پذیر در اینجا خطا صادر نمی‌کند؛

به‌طور پیش‌فرض، ناظرپذیر و ناظر در همان نخ اجرا می‌شوند. چند ناظرپذیر از پیش تعریف‌شده وجود دارند که در نخی غیر از نخ اصلی (در این مورد، نخ متد main) اجرا می‌شوند، اما برای اکثر آن‌ها اینطور نیست. بنابراین در اینجا، همه چیز در نخ متد [main] رخ می‌دهد:

  • ناظر، عنصر 1 را منتشر می‌کند؛
  • خطوط ۱۵ تا ۱۸ اجرا شده و این عنصر را نمایش می‌دهند؛
  • ناظر عنصر ۲ را منتشر می‌کند؛
  • خطوط ۱۵–۱۸ اجرا شده و این عنصر را نمایش می‌دهند؛
  • مشاهده‌پذیر عنصر ۳ را منتشر می‌کند؛
  • خطوط ۱۵–۱۸ اجرا شده و این عنصر را نمایش می‌دهند؛
  • مشاهده‌پذیر اعلان [completed] را منتشر می‌کند؛
  • خطوط 24–29 اجرا می‌شوند؛

این موضوع با نتایج به‌دست‌آمده نشان داده می‌شود:

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

کلاس [Exemple02] کلاس [Exemple01] را بازتولید می‌کند، این بار با استفاده از توابع لامبدا به عنوان پارامتر برای متد [Observable.subscribe]:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;

public class Exemple02 {
  public static void main(String[] args) {
    // متغیرهای عددی قابل مشاهده
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // اشتراک
    obs1.subscribe(
      (integer) -> System.out.printf("next : %s%n", integer),
      (th) -> System.out.println(th),
      () -> System.out.println("completed"));
  }
}

7.1.2. مثال ۰۳: کلاس ناظر

  

متد [Observable.subscribe] که به شما امکان می‌دهد به یک مشاهده‌پذیر مشترک شوید، نسخه‌های مختلفی دارد، از جمله موارد زیر:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;
import rx.Observer;

public class Exemple03 {
    public static void main(String[] args) {
        // متغیرهای عددی قابل مشاهده
        Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
        // اشتراک
        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);
            }
        });
    };
}

در خط ۱۳، به جای ارسال سه پارامتر به متد [subscribe]، نوع زیر از [Observer] به آن ارسال می‌شود:

 

نوع [Observer] یک رابط با سه متد است:

  • [onNext(T t)]، که هرگاه مشاهده‌پذیر یک عنصر t را منتشر می‌کند، فراخوانی می‌شود؛
  • [onError(Throwable th)]، که زمانی فراخوانی می‌شود که مشاهده‌پذیر یک استثنا th را پرتاب می‌کند؛
  • [onCompleted]، که زمانی فراخوانی می‌شود که مشاهده‌پذیر نشان دهد ارسال را متوقف کرده است؛

کد به روشی مشابه آنچه قبلاً توضیح داده شد کار می‌کند. نتایج زیر به دست می‌آیند:

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

7.1.3. مثال-۰۴: متد [Observable.create]

  

متد استاتیک Observable.create به صورت زیر تعریف شده است:

 
  • متد [create] یک نوع از نوع Observable<T> را برمی‌گرداند؛
  • پارامتر متد [create] تابعی از نوع [Observable.OnSubscribe<T>] است که به شرح زیر تعریف شده است:
 

نوع [Observable.OnSubscribe<T>] یک رابط تابعی است که خود رابط تابعی [Action1<Subscriber<? super T>>] را گسترش می‌دهد. متد [call] این رابط انتظار یک نوع [Subscriber] (مشترک، ناظر) را دارد که به شرح زیر تعریف شده است:

 

در [1] می‌بینیم که کلاس [Subscriber<T>] رابط [Observer<T>] را که در بخش 7.1.2 ارائه شده است، پیاده‌سازی می‌کند.

در نهایت، متد [<T> Observable.create]:

  • به‌عنوان پارامتر یک نمونه از نوع [Observable.OnSubscribe<T>] را می‌پذیرد که تنها یک متد با امضای زیر دارد: void call(Subscriber<T> s). نوع [Subscriber<T>] از نوع [Observer<T>] ارث می‌برد و بنابراین متدهای onNext، onError و onCompleted را دارد؛
  • یک `Observable<T>` را بازمی‌گرداند؛

متد [<T> Observable.create] یک observable پیکربندی‌شده را بازمی‌گرداند. هنوز هیچ عنصری ارسال نشده است. هنگامی که یک مشترک [Subscriber<T> s] به این قابل‌مشاهده مشترک می‌شود، متد [void call(s)] از تابعی که به‌عنوان پارامتر به متد [<T> Observable.create] ارسال شده است، فراخوانی می‌شود. نقش آن ارسال عناصر t از نوع T و فراخوانی متد ناظر [s.onNext(t)] در هر ارسال است. پس از اتمام این مرحله، روش ناظر [s.onCompleted(t)] باید فراخوانی شود و روش [call] باید خاتمه یابد. اگر متد [call] با یک استثنا th مواجه شود، متد ناظر [s.onError(th)] باید فراخوانی شود و متد [call] باید خاتمه یابد؛

برای تشریح این رفتار پیچیده، از کد زیر استفاده خواهیم کرد: [Exemple04]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.util.Random;

public class Exemple04 {
    public static void main(String[] args) {
        //پیکربندی قابل مشاهده اعداد حقیقی
        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
                    subscriber.onNext(new Random((i + 1)).nextDouble());
                }
                //پایان پخش
                subscriber.onCompleted();
            }
        });
        // اشتراک و در نتیجه پخش
        obs1.subscribe((d) -> System.out.printf("onNext %s%n", d), (th) -> System.out.printf("onError %s%n", th),
                () -> System.out.println("onCompleted"));
    }
}
  • خط ۱۱: یک مشاهده‌پذیر (observable) از نوع Double ایجاد می‌شود؛
  • خطوط ۱۱–۲۱: پارامتر متد [create] با یک کلاس ناشناس که شامل تنها متد [call] از خطوط ۱۲–۲۰ است، نمونه سازی می‌شود. مشاهده‌پذیر ایجادشده در خط ۱۱ آماده ارسال است، اما تنها زمانی ارسال خواهد کرد که یک ناظر برسد؛
  • خطوط ۱۳–۲۱: متد [call] یک مرجع به یک ناظر دریافت می‌کند؛
  • خطوط ۱۴–۱۷: سه عنصر به ناظر ارسال می‌شوند؛
  • خط ۱۹: اطلاع‌رسانی پایان انتقال به ناظر؛
  • خطوط ۲۳–۲۴: مشترک‌شدن برای قابل‌مشاهده از خط ۱۱. سه پارامتر [onNext, onError, onCompleted] در متد [subscribe] با استفاده از سه عبارت لامبدا پیاده‌سازی شده‌اند. این اشتراک، مشترک [Subscriber<Double>] را ایجاد می‌کند که در خط ۱۳ به متد [call] ارسال خواهد شد. سپس انتشار عناصر آغاز می‌شود؛
  • همه چیز در همان نخ رخ می‌دهد: مشاهده‌پذیر و مشاهده‌گر؛

نتایج زیر به دست می‌آیند:

1
2
3
4
onNext 0.7308781907032909
onNext 0.7311469360199058
onNext 0.731057369148862
onCompleted

متد [Observable.create] امکان ایجاد یک قابل مشاهده از هر رویدادی را فراهم می‌کند. این متد همان چیزی است که ما در بخش ۲ از مقدمه برای تبدیل یک رابط همگام به رابطی ناهمزمان استفاده کردیم.

7.1.4. مثال-۰۵: بازسازی [Exemple-04]

  

مثال زیر نسخه جدیدی از متد استاتیک [Observable.subscribe] را نشان می‌دهد:


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) {
        //پیکربندی یک متغیر قابل مشاهده با مقادیر حقیقی
        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++) {
                    //در انتظار
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // خطا
                        subscriber.onError(e);
                    }
                    // اقدام
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // تکمیل شد
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // یک مشترک
        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));
            }
        };

        //اشتراک
        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()));
    }
}
  • خط ۵۶: نسخه جدید متد استاتیک [Observable.subscribe] نوع [Subscriber] را که در پاراگراف قبلی معرفی کردیم، به عنوان پارامتر می‌پذیرد؛
  • خطوط ۳۷–۵۲: مشترک (مشترک، ناظر). این رابط Observer را با سه متد خود پیاده‌سازی می‌کند: onNext، onError و onCompleted؛
  • خطوط ۶۱–۶۴: از این پس، ما بر روی رشته‌هایی تمرکز خواهیم کرد که در آن‌ها مشاهده‌شدنی و مشاهده‌گر آن اجرا می‌شوند؛
  • خط ۶۲: نام تِرد؛
  • خط ۶۳: زمان فعلی بیان‌شده به ثانیه‌ها و میلی‌ثانیه‌ها. این به ما امکان می‌دهد تا در طول زمان، انتشار عناصر توسط مشاهده‌شونده و پردازش آن‌ها توسط مشاهده‌گر را پیگیری کنیم؛
  • این کد همان کارایی کد قبلی را دارد. ما صرفاً کد دوم را بازسازی کرده‌ایم؛

نتایج به‌دست‌آمده به شرح زیر است:

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]
  • خط ۱ نتایج: قبل از خط ۵۶ کد، هنوز هیچ اتفاقی نیفتاده است. متغیر مشاهدنی صرفاً پیکربندی شده است؛
  • خط ۲ نتایج: خط ۵۶ کد باعث فراخوانی متد [call] در خط ۱۵ می‌شود. در خط ۳، عدد حقیقی ۸۰.۳۹ به ناظر ارسال می‌شود؛
  • خط ۴: ناظر عدد ارسال‌شده را دریافت می‌کند؛
  • خطوط ۵–۸: فرآیند قبلی دو بار تکرار می‌شود؛
  • خط ۹: مشاهده‌پذیر اعلان پایان انتقال را ارسال می‌کند؛
  • خط ۱۰: ناظر آن را دریافت می‌کند؛
  • خط ۱۱: توسط خط ۵۷ کد نمایش داده می‌شود؛

بنابراین می‌توانیم ببینیم که تنها خط ۵۶ — اشتراک — باعث نمایش خطوط ۲ تا ۱۰ نتایج شد. وقتی با کتابخانه RxJava شروع می‌کنیم، این سؤال پیش می‌آید که همه چیز چگونه در کنار هم قرار می‌گیرد، به‌ویژه پیوندهای بین ناظر و مشاهده‌شونده. در اینجا می‌توانیم ببینیم که خط ۵۶، اشتراک‌گذاری با قابلیت مشاهده‌شدنی،

  • انتشار تمام عناصر قابل مشاهده را تحریک کرد؛
  • که مشاهده‌گر و مشاهده‌شدنی در همان نخ اجرا می‌شوند؛
  • و در نتیجه، ما توالی زیر را مشاهده می‌کنیم: ارسال عنصر i، مشاهده عنصر i، ارسال عنصر (i+1)، مشاهده عنصر (i+1)، …

به یاد می‌آوریم که فرستنده قبل از ارسال عناصر خود در انتظار بود:


                    //در حال انتظار
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // خطا
                        subscriber.onError(e);
}

که i در خط ۳ نمایانگر شماره انتقال (0 ≤ i < 3) است. اگر به زمان‌های انتقال عناصر قابل مشاهده نگاه کنیم:

  • رده‌های ۲ و ۳: عنصر ۰ تقریباً ۵۰۰ میلی‌ثانیه پس از شروع اشتراک ارسال شد؛
  • خطوط ۳ و ۵: عنصر ۱ تقریباً ۴۰۰ میلی‌ثانیه پس از عنصر ۰ ارسال شد؛
  • خطوط ۵ و ۷: عنصر ۲ تقریباً ۳۰۰ میلی‌ثانیه پس از عنصر ۱ ارسال شد؛

7.2. رشته اجرایی، رشته مشاهده

7.2.1. مثال-06: قابل مشاهده و مشاهده‌گر در تریدی غیر از [main]

  

ما مثال قبلی را به صورت زیر بازسازی می‌کنیم: [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) {

        // نگهبان مانع
        CountDownLatch latch = new CountDownLatch(1);

        // پیکربندی یک متغیر مشاهدنی با مقادیر حقیقی
        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++) {
                    //در حال انتظار
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // خطا
                        subscriber.onError(e);
                    }
                    // اقدام
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // تکمیل شد
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // یک مشترک
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // کاهش مانع
                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));
            }
        };

        // پی‌گیری پیکربندی قابل مشاهده ادامه دارد
        obs1 = obs1.subscribeOn(Schedulers.computation());
        // اشتراک
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // منتظر در مقابل مانع
        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()));
    }
}
  • خط ۱۶: ما یک مانع (سِمافور) با استفاده از یک شی از نوع [CountDownLatch] ایجاد می‌کنیم. این شی برای همگام‌سازی نخ‌ها با یکدیگر استفاده می‌شود. در اینجا، مقدار اولیه آن ۱ است که ما آن را به عنوان مقدار مانع (یا مقدار سِمافور) می‌نامیم. یک نخ با استفاده از عملیات زیر منتظر مانع می‌ماند:

latch.await();

اگر مقدار مانع >0 باشد، نخ مسدود می‌شود. یک نخ می‌تواند مقدار داخلی مانع را افزایش یا کاهش دهد. در خط 48، مقدار مانع به اندازه 1 کاهش می‌یابد.

  • خط ۶۳: observable طوری پیکربندی شده است که روی یک تار (thread) که توسط زمان‌بندی‌کننده [Schedulers.computation()] فراهم می‌شود، اجرا شود. این زمان‌بندی‌کننده می‌تواند به تعداد هسته‌های ماشین اجرایی، تار فراهم کند. بخش مربوط به مثال کاربردی، استفاده از برنامه‌ریزهای دیگر را نشان داد (به بخش ۲.۸ مراجعه کنید)؛

اصل پشت این کد به شرح زیر است:

  • متد [main] در نخ اصلی اجرا می‌شود؛
  • خط ۶۶: انتشار عناصر از ناظر را آغاز می‌کند. این عناصر بر روی رشته‌ای غیر از رشته اصلی منتشر خواهند شد؛
  • خط ۷۰: نخ اصلی مسدود شده است زیرا مانع مقدار ۱ دارد (به خط ۱۶ مراجعه کنید). این نخ تنها زمانی می‌تواند ادامه دهد که این مقدار به ۰ تغییر کند. این اتفاق در خط ۴۸ رخ می‌دهد. این ناظر است که مانع را پایین می‌آورد، زمانی که اعلان دریافت می‌کند که شیء قابل مشاهده ارسال خود را به پایان رسانده است؛

اجرا نتایج زیر را به دست می‌دهد:

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]
  • خط ۱: اشتراک‌گذاری در شرف انجام است؛
  • خط ۲: این باعث اجرای متد [call] در نخ [RxComputationThreadPool-1] می‌شود. اکنون ما یک اجرای موازی با دو نخ داریم؛
  • خط ۳: به دلیلی نامعلوم، نخ [RxComputationThreadPool-1] کنترل را واگذار کرده است. سپس نخ [main] کنترل را به دست می‌گیرد و توسط گاردریل (خط ۷۰ کد) مسدود می‌شود. از این نقطه به بعد، تنها نخ [RxComputationThreadPool-1] می‌تواند عمل کند؛
  • خطوط ۴–۱۱: ما رفتار مشاهده‌شده قبلی بین مشاهده‌شونده و ناظر آن را می‌بینیم، اما اکنون همه چیز در داخل نخ [RxComputationThreadPool-1] رخ می‌دهد؛
  • خطوط ۱۲–۱۳: ناظر مانع را پایین آورده است (خط ۴۸ کد) و نخ [RxComputationThreadPool-1] خاتمه یافته است. نخ [main] کنترل را به دست می‌گیرد و دو پیام را نمایش می‌دهد؛

7.2.2. مثال-۰۷: مشاهده‌پذیر و مشاهده‌گر در دو نخ مختلف

  

ما مثال قبلی را به شرح زیر تغییر می‌دهیم:


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) {

        // خدمه مانع
        CountDownLatch latch = new CountDownLatch(1);

        // پیکربندی یک مشاهده‌پذیر از اعداد حقیقی
        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++) {
                    // انتظار
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // خطا
                        subscriber.onError(e);
                    }
                    // اقدام
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // تکمیل شد
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // یک مشترک
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // کاهش مانع
                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));
            }
        };

        // پی‌گیری پیکربندی قابل مشاهده ادامه دارد
        obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
        // اشتراک
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // منتظر بالا آمدن مانع
        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()));
    }
}

کد با کد مثال قبلی یکسان است، به جز خط ۶۳:


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

که مشاهده‌پذیر (subscribeOn) و مشاهده‌گر (observeOn) را برای اجرا روی یکی از رشته‌های ارائه‌شده توسط زمان‌بندی‌کننده [Schedulers.computation()] پیکربندی می‌کند.

نتایج به‌دست‌آمده به شرح زیر است:

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]

نکات زیر قابل توجه است:

  • قابل مشاهده در نخ [RxComputationThreadPool-4] اجرا می‌شود (خطوط ۳–۴، ۶، ۸–۹)؛
  • ناظر در نخ [RxComputationThreadPool-3] (خطوط 5، 7، 10–11) اجرا می‌شود؛
  • آنها به‌طور مستقل اجرا می‌شوند. بنابراین، در خطوط ۸–۹، قابل‌مشاهده دو اعلان (onNext, onCompleted) را قبل از اینکه ناظر اعلان [onNext] را بازیابی کند (خط ۱۰) منتشر می‌کند؛

کتابخانه RxJava انتقال داده‌ها (انتشارها) را از نخ قابل مشاهده به نخ ناظر مدیریت می‌کند. توسعه‌دهنده نیازی به نگرانی در این مورد ندارد.

ما دیدیم چگونه observableها را ایجاد کنیم (Observable.from, Observable.create). اکنون به observableهای از پیش تعریف‌شده در کتابخانه RxJava می‌پردازیم.

7.3. ناظرهای از پیش تعریف‌شده

7.3.1. مثال-۰۸: متد [Observable.range]

 

از این پس، ما از کلاس‌های اختصاصی برای فرآیندهای مشاهده‌شده و ناظران آن‌ها استفاده خواهیم کرد. ایده این است که بتوانیم نام‌ها، نخ‌های اجرایی و زمان‌های اجرای آن‌ها را ثبت کنیم تا بتوانیم آن‌ها را در طول زمان ردیابی کنیم.

کلاس [Process] صرفاً یک Observable خواهد بود که می‌توان آن را نام‌گذاری کرد. این کلاس رابط زیر را پیاده‌سازی خواهد کرد: [IProcess]:


package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcess<T> {

    // نام قابل مشاهده
    public String getName();

    //قابل مشاهده
    public Observable<T> getObservable();

}

این رابط می‌تواند توسط کلاس زیر پیاده‌سازی شود، [Process<T>]:


package dvp.rxjava.observables.utils;

import rx.Observable;
import rx.Scheduler;

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

    // نام قابل مشاهده
    protected String name;
    // فرآیند مشاهده‌شده
    protected Observable<T> observable;

    // سازنده‌ها
    public Process(String name, Observable<T> observable) {
        // ابتکاری‌های محلی
        this.name = name;
        this.observable = observable;
    }

    // گیرنده و تنظیم‌کننده
    public String getName() {
        return name;
    }

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

}
  • خط ۹: نام فرایند؛
  • خط ۱۱: متغیر مشاهداتی مشاهده‌شده؛
  • خطوط 14–18: سازنده؛

ناظر توسط کلاس زیر تعریف خواهد شد: [Observateur]:


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> {

...
}
  • خط ۱۱: کلاس Observateur<T> از کلاس Subscriber<T> که در بخش ۷.۱.۳ به‌طور مختصر معرفی شد، ارث می‌برد. ما از آن به‌عنوان آرگومان برای متد [Observable.subscribe] استفاده خواهیم کرد:

// اجرای قابل مشاهده (مشاهده)
obs1.subscribe(observateur);

متد [Observable.subscribe] که در خط ۲ بالا استفاده شده است، تعریف زیر را دارد:

 

نقش [Subscriber] عمدتاً مدیریت عناصری است که توسط مشاهده‌پذیری که در آن مشترک شده است صادر می‌شوند، با استفاده از روش‌های رابط [Observer]: onNext، onError، onCompleted. کلاس [Subscriber] دارای متدهای زیر است:

 

در کد کلاس [Observateur]، از متدهای [1] و isUnsubscribed برای تعیین اینکه آیا اشتراک مشترک لغو شده است یا خیر، استفاده خواهیم کرد. کلاس کامل [Observateur<T>] به شرح زیر است:


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> {

    // یک نگهبان (سِمافور)
    private CountDownLatch latch;
    // یک روش نمایش
    private Consumer<String> showInfos;
    // نام ناظر
    private String observerName;
    // نام فرآیند مشاهده‌شده
    private String processName;

    // سازنده‌ها
    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>
    @Override
    public void onCompleted() {
        // پایان پخش‌ها
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s,%s].onCompleted", observerName, processName));
        }
        // پایان بلوک نخ اصلی
        latch.countDown();
    }

    @Override
    public void onError(Throwable e) {
        // خطای پخش
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s, %s].onError (%s)", observerName, processName, e));
        }
    }

    @Override
    public void onNext(T value) {
        // انتشار اضافی
        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));
            }
        }
    }
}
  • علاوه بر ویژگی‌های یک Subscriber، ناظر Observateur اطلاعات زیر را حمل خواهد کرد:
    • خط ۱۴: یک گارد یا نیم‌فور که برای مسدود کردن نخ اصلی تا زمانی که ناظر تمام عناصری را که توسط مشاهده‌پذیر منتشر شده‌اند دریافت نکرده است، استفاده می‌شود. این کار در خط ۳۶ کد انجام می‌شود، زمانی که ناظر اعلان پایان انتشار را از مشاهده‌پذیر دریافت می‌کند؛
    • خط ۱۶: یک نمونه از Consumer<String>، که برای نمایش یک پیام در کنسول استفاده خواهد شد؛
    • خط ۱۸: نام ناظر، برای تمایز بین آن‌ها در صورت وجود چندین ناظر؛
    • خط ۲۰: نام فرآیند مشاهده‌شده؛
  • خطوط ۳۶، ۴۶، ۵۴: متدهای [onCompleted, onError, onNext] از رابط [Observer<T>] که توسط کلاس انتزاعی [Subscriber<T>] پیاده‌سازی شده‌اند. این کلاس آن‌ها را پیاده‌سازی نمی‌کند. بنابراین این کار باید در کلاس‌های فرزند آن انجام شود. قبل از انجام هرگونه عملیاتی در این متدها، بررسی می‌کنیم که آیا ناظر از مشاهده‌شونده‌ای که زیر نظر دارد، لغو اشتراک شده است یا خیر؛
  • خط ۵۹: متد ناظر [onNext] رشته jSON را از عنصر دریافتی می‌نویسد. این به ما امکان می‌دهد تا انواع مختلف عناصر را نمایش دهیم؛

با این اوصاف، بیایید روش جدیدی از کلاس Observable را بررسی کنیم، یعنی روش [range]:

 

مشاهده‌پذیر Observable.range(n,m) اعداد صحیحی بین n تا n+m-1 را به صورت (m) تکه منتشر می‌کند. ما آن را با استفاده از کد زیر [Exemple08] بررسی خواهیم کرد:


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 {

        // تعداد ناظران
        final int nbObservateurs = 2;

        // سِمافور
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        //پیکربندی قابل مشاهده
        Observable<Integer> obs1 = Observable.range(15, 3).subscribeOn(Schedulers.computation());
        // اجرای قابل مشاهده (مشاهده)
        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"));
        }
        // انتظار
        showInfos.accept("main : attente fin observation");
        latch.await();
        // پایان
        showInfos.accept("main : fin observation");
    }

    //نمایش‌ها
    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()));
}
  • خط ۱۶: ما از دو ناظر استفاده خواهیم کرد؛
  • خط ۱۹: مانع (سِمافور) روی دو مقداردهی اولیه می‌شود زیرا هر ناظر را روی یک نخ مجزا قرار می‌دهیم. بنابراین نخ اصلی باید منتظر پایان هر دو نخ ناظر بماند؛
  • خط ۲۲: ما قابل‌مشاهده را طوری پیکربندی می‌کنیم که روی یک تِردِ برنامه‌ریز [Schedulers.computation()] اجرا شود. ناظر روی همان تِردِ قابل‌مشاهده خواهد بود؛
  • خطوط ۲۵–۲۷: ما دو ناظر را برای قابل مشاهده ثبت‌نام می‌کنیم. این کار باعث می‌شود که قابل مشاهده برای هر ناظر به طور کامل اجرا شود: اعداد صحیح ۱۵، ۱۶ و ۱۷ ارسال خواهند شد؛
  • خط ۳۰: نخ اصلی منتظر پایان کار ناظران می‌ماند؛

نتایج به‌دست‌آمده به شرح زیر است:

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]
  • خط ۲: نخ اصلی مسدود شده و منتظر پایان کار دو ناظر است؛
  • خطوط ۳–۴: می‌بینیم که ناظر ۰ روی نخ [RxComputationThreadPool-1] و ناظر ۱ روی نخ [RxComputationThreadPool-2] قرار دارد؛
  • خطوط ۳–۱۰: می‌بینیم که هر دو ناظر دقیقاً عناصر یکسانی دریافت می‌کنند؛

ما از کلاس Observateur که به این صورت تعریف شده است، برای نشان دادن رفتار انواع دیگر قابل‌مشاهده‌ها استفاده خواهیم کرد.

7.3.2. مثال-۰۹: متدهای Observable.[interval, take, doNext]

  
 

این مثال استفاده از قابل مشاهده Observable.interval (فاصله زمانی طولانی، واحد TimeUnit) را نشان می‌دهد که اعداد صحیح بزرگ را در فواصل زمانی منظم منتشر می‌کند. به نکته [1] توجه کنید: به طور پیش‌فرض، مشاهده‌پذیر [Observable.interval] روی یکی از رشته‌های زمان‌بندی‌کننده [Schedulers.computation] اجرا می‌شود.

کد به شرح زیر خواهد بود:


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 {

        // تعداد ناظران
        final int nbObservateurs = 2;

        // سِمافور
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // پیکربندی قابل مشاهده
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // اجرای قابل مشاهده (مشاهده)
        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"));
        }
        // انتظار
        showInfos.accept("main : attente fin observation");
        latch.await();
        // پایان
        showInfos.accept("main : fin observation");
    }

    //نمایش‌ها
    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()));
}
  • خط ۲۲: این قابل مشاهده هر ۵۰۰ میلی‌ثانیه اعداد صحیح بلند را منتشر می‌کند. دنباله با عدد ۰ شروع می‌شود؛
  • خط ۲۲: این قابل مشاهده تعداد نامحدودی مقدار صادر می‌کند. متد [Observable.take(n)] یک قابل مشاهده جدید ایجاد می‌کند که تنها n عنصر اول صادر شده را حفظ می‌کند؛
 

بیایید بار دیگر به کد این observable نگاهی بیندازیم:


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

خط ۲: متد [Observable.doOnNext] هر بار که observable یک عنصر جدید صادر می‌کند، اجرا می‌شود. این کار اغلب برای ثبت اطلاعات استفاده می‌شود. در اینجا، ما می‌خواهیم تاریخ انتشار عناصر را ثبت کنیم تا بررسی کنیم که آیا فاصله زمانی ۵۰۰ میلی‌ثانیه‌ای به درستی حفظ می‌شود یا خیر. متد [Observable.doOnNext]، ناظر (observable) مورد اعمال خود را تغییر نمی‌دهد. تعریف آن به شرح زیر است:

 

اجرا نتایج زیر را تولید می‌کند:

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]
  • رده‌های ۳، ۷ و ۱۱: می‌توانیم ببینیم که فاصله انتقال تقریباً ۵۰۰ میلی‌ثانیه است؛
  • دو ناظر، البته، روی دو نخ مختلف قرار دارند، حتی اگر قابل مشاهده طوری پیکربندی نشده باشد که با یک زمان‌بندی‌کنندهٔ خاص اجرا شود. این رفتار پیش‌فرض قابل مشاهدهٔ [Observable.interval] است که در اینجا می‌بینیم؛

7.3.3. Examples-10/12: متدهای Observable.[error, empty, never]

 

از این پس، در ارائهٔ مثال‌های روش‌های کلاس [Observable] مختصرتر خواهیم بود. کد قبلی به شرح زیر بود:


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 {

        // تعداد ناظران
        final int nbObservateurs = 2;

        // سِمافور
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // پیکربندی قابل مشاهده
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // اجرای قابل مشاهده (مشاهده)
        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"));
        }
        // انتظار
        showInfos.accept("main : attente fin observation");
        latch.await();
        // پایان
        showInfos.accept("main : fin observation");
    }

    // نمایش‌ها
    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()));
}

این کد قبلاً برای مثال قبلی استفاده شده بود. تنها خطوط ۲۱–۲۲ متفاوت بودند. بنابراین، بیشتر این کد را به کلاس زیر، [ProcessUtils]، منتقل خواهیم کرد:


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 {

        // سِمافور
        CountDownLatch latch = new CountDownLatch(nbObservateurs * processes.length);

        // اجرای قابل مشاهده (مشاهده)
        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()));
            }
        }
        // انتظار
        showInfos.accept("main : attente fin observation");
        latch.await();
        // پایان
        showInfos.accept("main : fin observation");
    }

    //نمایش‌ها
    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()));
}
  • خط ۱۳: متد دو پارامتر می‌گیرد:
    • nbObservateurs: تعداد ناظران برای فرآیندها که به‌عنوان پارامتر دوم ارسال می‌شود؛
    • processes: فرآیندها (قابل‌مشاهده‌های نام‌گذاری‌شده) که باید مشاهده شوند. به لطف نشانه‌گذاری [IProcess<?>]، فرآیندها قادر خواهند بود عناصر از انواع مختلف را منتشر کنند؛
  • خط ۱۶: سِمافور باید وقتی همه ناظران تمام مشاهدات خود را تکمیل کردند، سبز شود. بنابراین مقدار اولیه سِمافور برابر است با حاصل ضرب تعداد ناظران در تعداد مشاهدات؛
  • خطوط ۲۰–۲۵: هر ناظر به تمام فرآیندهایی که نیاز به مشاهده دارند، مشترک می‌شود؛
  • خط ۲۳: مقدار قابل مشاهده از فرآیند بازیابی می‌شود (به بخش ۷.۳.۱ مراجعه کنید)؛
  • خط ۲۳: یک ناظر به آن مشترک می‌شود. چهار مورد اطلاعات به ناظر ارسال می‌شود:
    • نام آن؛
    • سِمافورِ که باید هنگام دریافت اعلان مبنی بر پایان ارسالِ قابل‌مشاهده‌ای که نظارت می‌کند، مقدار آن را کاهش دهد؛
    • روش مورد استفاده زمانی که می‌خواهد اطلاعات را در کنسول ثبت کند؛
    • نام فرایندی که آن را مشاهده خواهد کرد؛

با تعریف این کلاس‌ها، مثال ۱۰ به شرح زیر خواهد بود:


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 {
        //پیکربندی قابل مشاهده
        Observable<?> obs = Observable.error(new RuntimeException("Erreur !!!")).subscribeOn(Schedulers.computation());
        // اجرای قابل مشاهده (مشاهده)
        ProcessUtils.subscribe(2,new Process<>("process1", obs));
    }
}

در خط ۱۱، متد استاتیک [Observable.error] به صورت زیر تعریف شده است:

 

بنابراین خط ۸ یک مشاهده‌پذیر را پیکربندی می‌کند که به سادگی یک استثنا را به متد [onError] مشترکین خود پرتاب می‌کند. اجرای آن نتایج زیر را می‌دهد:


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]

در خطوط ۳ و ۴، متد [onError] هر دو مشترک، استثنایی را که توسط مشاهده‌پذیر پرتاب شده بود، دریافت کردند.

این اجرا یک ویژگی خاص دارد: متدهای [onCompleted] هر دو ناظر فراخوانی نشدند. در نتیجه، مانع پایین نیامد و نخ اصلی در متد استاتیک [ProcessUtils.subscribe] در خط ۳ زیر همچنان مسدود باقی می‌ماند:


// در حال انتظار
showInfos.accept("main : attente fin observation");
latch.await();
//پایان
showInfos.accept("main : fin observation");

در اینجا می‌بینیم که در صورت بروز خطا در observable، متد مشترکین [onCompleted] فراخوانی نمی‌شود. بنابراین متد [Observateur.onError] را به شرح زیر اصلاح می‌کنیم:


    @Override
    public void onError(Throwable e) {
        // خطای انتقال
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber[%s, %s].onError (%s)", observerName, processName, e));
        }
        // پایان بلوک نخ اصلی
        latch.countDown();
}

ما خطوط ۷–۸ را برای حذف مانع در صورت وقوع خطای قابل مشاهده اضافه می‌کنیم. با این کد جدید، اجرای برنامه نتایج زیر را می‌دهد:


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]

اکنون خط ۵ را به دست می‌آوریم که قبلاً نداشتیم.

مثال ۱۱ به شرح زیر خواهد بود:


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 {
        //پیکربندی قابل مشاهده
        Observable<?> obs1 = Observable.empty();
        // اجرا (مشاهده) قابل مشاهده
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

در خط ۱۰، متد استاتیک [Observable.empty] یک observable ایجاد می‌کند که هیچ عنصری را ارسال نمی‌کند. این متد تنها اعلان پایان ارسال (end-of-emission) را ارسال می‌کند؛

 

اجرای کد در مثال بالا نتایج زیر را تولید می‌کند:

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]
  • در خطوط ۲ و ۳ می‌بینیم که هر دو ناظر، اعلان پایان انتشار را بدون دریافت هیچ عنصری در beforehand دریافت می‌کنند.

شاید این سؤال پیش بیاید که این روش چه کاربردی دارد. این روش را می‌توان به روشی مشابه یک مجموعه که در ابتدا خالی است و سپس عناصر به آن اضافه می‌شوند، استفاده کرد:

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

در خط ۳، ناظر قابل مشاهده اولیه obs (خط ۱) با سایر ناظرهای قابل مشاهده ادغام می‌شود.

مثال ۱۲ روش استاتیک [Observable.never] را نشان می‌دهد:


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 {
        //پیکربندی قابل مشاهده
        Observable<?> obs1 = Observable.never();
        // قابل مشاهده اجرا (مشاهده)
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

متد استاتیک [Observable.never] یک observable ایجاد می‌کند که هرگز emit نمی‌کند:

 

اجرای مثال نتایج زیر را تولید می‌کند:

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

خط ۲: نخ اصلی به طور نامحدود منتظر می‌ماند. این به این دلیل است که هیچ observable‌ای اعلان [onCompleted] را ارسال نمی‌کند، که این اجازه را می‌دهد تا سِمافور (مانع) روی سبز (کاهش مانع) تنظیم شود.

7.4. Multi-threading

7.4.1. مثال ۱۳: نخ اقدام، نخ ناظر

در بخش 7.1.3، ما یک observable را با استفاده از متد استاتیک [Observable.create] ایجاد کردیم:

 
  • متد [create] یک نوع Observable<T> را برمی‌گرداند؛
  • پارامتر متد [create] تابعی از نوع [Observable.OnSubscribe<T>] است که به صورت زیر تعریف شده است:
 

نوع [Observable.OnSubscribe<T>] یک رابط تابعی است که خود رابط تابعی [Action1<Subscriber<? super T>>] را گسترش می‌دهد. متد [call] این رابط انتظار یک نوع [Subscriber] (مشترک، ناظر) را دارد. در ادامه این سند، گاهی اوقات به نوع [Observable.OnSubscribe<T>] به عنوان یک اکشن اشاره خواهیم کرد. ما قصد داریم اکشن‌های سفارشی ایجاد کنیم که هر یک نامی خواهند داشت. این‌ها نمونه‌هایی از رابط زیر [IProcessAction] خواهند بود:

  

package dvp.rxjava.observables.utils;

import rx.Observable;

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

    // اقدام دارای نام است
    public String getName();
}
  • خط ۵: رابط [IProcessAction<T>] تمام ویژگی‌های رابط [Observable.OnSubscribe<T>] را دارد؛
  • خط ۸: این همچنین دارای متدی به نام [getName] است که نام نمونهٔ پیاده‌ساز رابط را برمی‌گرداند؛

ما از اقدام زیر با نام [ProcessAction01] استفاده خواهیم کرد:


package dvp.rxjava.observables.utils;

import java.util.Random;

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

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

    // داده‌ها
    private String name;
    private int nbValues;
    private Func1<Integer, T> func1;

    // سازنده‌ها
    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++) {
            // انتظار
            try {
                Thread.sleep(new Random().nextInt(500));
            } catch (InterruptedException e) {
                // خطا
                ProcessUtils.showInfos.accept(String.format("Observable (%s) onError", getName()));
                subscriber.onError(e);
            }
            //ارسال یک عنصر
            T value = func1.call(i);
            ProcessUtils.showInfos.accept(String.format("Observable (%s,%s) onNext (%s)", getName(), i, value));
            subscriber.onNext(value);
        }
        // تکمیل‌شده
        ProcessUtils.showInfos.accept(String.format("Observable (%s) onCompleted", getName()));
        subscriber.onCompleted();
    }

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

}
  • خط ۸: کلاس [ProcessAction01<T>] رابط [IProcessAction<T>] و در نتیجه رابط [Observable.OnSubscribe<T>] را پیاده‌سازی می‌کند؛
  • خط ۱۱: نام اقدام؛
  • خط ۱۲: تعداد مقادیری که باید صادر شوند؛
  • خط ۱۳: یک نمونه از نوع [Func1<Integer, T>] که با دریافت یک عدد صحیح، یک نوع T ایجاد می‌کند که توسط ناظر (lines 35 and 37) صادر خواهد شد؛
  • خطوط ۱۶–۲۰: به سازنده، نام اقدام، تعداد مقادیری که باید صادر شوند و تابع صدور پاس می‌شود؛
  • خطوط ۲۳–۴۲: کد فرآیند؛
  • خط ۲۳: متد [call] مشترکِ ناظرِ مرتبط با فرآیند را به عنوان پارامتر می‌پذیرد؛
  • خط ۲۸: فرآیند عناصر خود را پس از یک وقفه با مدت زمان تصادفی منتشر می‌کند؛
  • خط ۳۲: ارسال یک خطا؛
  • خط ۳۷: یک انتشار عادی؛
  • خط ۴۱: ارسال اعلان پایان انتشار؛
  • خطوط ۲۵–۳۸: این عمل پس از یک زمان انتظار تصادفی، مقدار حقیقی nbValues را صادر می‌کند (خط ۳۰);
  • خط ۳۵: مقداری که باید صادر شود توسط تابع [func1] که به‌عنوان پارامتر به سازنده (خط ۱۶) پاس شده است، فراهم می‌شود؛

ما کلاس [Process] را (رجوع کنید به بخش 7.3.1) بازسازی می‌کنیم تا بتوان آن را با یک اکشن نام‌گذاری‌شده نیز نمونه‌سازی کرد. ما کانتراکتور زیر را به آن اضافه می‌کنیم:


public Process(IProcessAction<T> na, Scheduler schedulerObserved, Scheduler schedulerObserver) {
        // نام فرآیند=نام اقدام
        name = na.getName();
        // اقدام --> قابل مشاهده
        observable = Observable.create(na);
        // رشته اجرایی فرآیند مشاهده‌شده
        if (schedulerObserved != null) {
            observable = observable.subscribeOn(schedulerObserved);
        }
        // رشتهٔ نظارت ناظر
        if (schedulerObserver != null) {
            observable = observable.observeOn(schedulerObserver);
        }
    }
  • خط ۱، سازنده ۳ پارامتر می‌گیرد:
    1. اقدام نام‌گذاری‌شده‌ای که برای ساختنیاب (observable) استفاده خواهد شد (خط ۵)؛
    2. برنامه‌ریز برای فرآیند مشاهده‌شده (می‌تواند null باشد)؛
    3. برنامه‌ریز برای ناظر (می‌تواند null باشد)؛
  • خط ۵: مشاهده‌پذیر از اکشنی که به‌عنوان پارامتر ارسال شده، ایجاد می‌شود؛

کد زیر، [Exemple13]، چندین قابل مشاهده را مشاهده می‌کند:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // فرآیند ۳
        Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
                Schedulers.computation());
        // فرآیند ۴
        Process<Boolean> process4 = new Process<>(new ProcessAction01<Boolean>("process4", 4, i -> i % 2 == 0), null, null);
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process1);
        ProcessUtils.subscribe(1, process2);
        ProcessUtils.subscribe(1, process3);
        ProcessUtils.subscribe(1, process4);
    }
}
  • خطوط ۱۳–۱۵: فرآیند process1 یک عدد حقیقی را روی یک نخ محاسباتی تولید می‌کند که در یک نخ محاسباتی دیگر مشاهده خواهد شد؛
  • خطوط ۱۷–۱۸: فرآیند process2 دو رشتهٔ کاراکتری را روی یک نخ محاسباتی تولید می‌کند و هیچ نشانه‌ای در مورد نخ ناظر داده نشده است. نتایج نشان می‌دهند که، به‌طور پیش‌فرض، مشاهده روی همان نخی انجام می‌شود که فرآیند روی آن اجرا می‌شود؛
  • خطوط ۲۰–۲۱: فرآیند process3 سه عدد صحیح را روی یک رشته محاسباتی تولید می‌کند که روی یک رشته محاسباتی دیگر مشاهده خواهد شد. نتایج نشان می‌دهد که این فرآیند به طور پیش‌فرض روی رشته اصلی اجرا می‌شود؛
  • خط ۲۳: فرآیند process4 چهار مقدار بولی را روی یک رشته نامشخص تولید می‌کند که روی یک رشته نامشخص مشاهده خواهد شد. نتایج نشان می‌دهد که هم اجرای فرآیند و هم مشاهده آن به طور پیش‌فرض روی رشته اصلی انجام می‌شود؛

نتیجه اجرای این کد به شرح زیر است:

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 یک عدد حقیقی (خط ۴) را روی نخ محاسباتی [RxComputationThreadPool-4] تولید می‌کند، که روی نخ محاسباتی [RxComputationThreadPool-3] (خط ۶) مشاهده می‌شود؛
  • فرآیند process2 دو رشته کاراکتری (خطوط 12 و 14) را در نخ محاسباتی [RxComputationThreadPool-5] تولید می‌کند که در همان نخ (خطوط 13 و 15) مشاهده می‌شوند؛
  • فرآیند process3 سه عدد صحیح (خطوط 21، 23، 25) را در نخ اصلی تولید می‌کند که در نخ محاسباتی [RxComputationThreadPool-6] مشاهده می‌شوند (خطوط 22، 24، 28);
  • فرآیند process4 سه عدد صحیح (خطوط 21، 23، 25) را در نخ اصلی تولید می‌کند که در همان نخ اصلی (خطوط 22، 24، 28) مشاهده می‌شوند؛

از خوانندگان دعوت می‌شود موارد فوق را دنبال کنند:

  • چرخهٔ عمر فرآیند مشاهده‌شده و نخ آن؛
  • چرخهٔ عمر ناظر آن و نخ آن؛

بخش عمده‌ای از جذابیت کتابخانه‌های Rx در همین چندریسمانی نهفته است، که توسعه‌دهنده نیازی به مدیریت آن به صورت دستی ندارد.

7.5. ترکیب چندین مشاهده‌پذیر

7.5.1. مثال ۱۴: ادغام دو مشاهده‌پذیر با [Observable.merge]

اکنون به متدهای استاتیک کلاس [Observable] می‌پردازیم که امکان ترکیب چندین قابل‌مشاهده در یک قابل‌مشاهده نتیجه واحد را فراهم می‌کنند.

اولین مثال از این نوع به شرح زیر است:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        //ادغام
        Process<?> process12 = new Process<>("process12",
                Observable.merge(process1.getObservable(), process2.getObservable()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process12);
    }
}
  • خطوط ۱۵–۱۷: فرایندی به نام [process1] سه عدد حقیقی را روی یک نخ محاسباتی خروجی خواهد داد. همچنین روی یک نخ محاسباتی مشاهده خواهد شد؛
  • خطوط ۱۹–۲۰: فرایندی به نام [process2] دو رشتهٔ کاراکتری را روی یک نخ محاسباتی صادر خواهد کرد. نخ مشاهده مشخص نشده است. قبلاً دیده‌ایم که در این مورد، نخ مشاهده همان نخ محاسباتی است؛
  • خط ۲۳: دو فرآیند ادغام می‌شوند، یعنی یک متغیر قابل مشاهده ایجاد می‌شود که عناصر آن به‌طور همزمان از هر دو فرآیند می‌آیند. برای این کار از روش استاتیک [Observable.merge] استفاده می‌شود:
 

برخلاف آنچه نمودار بالا ممکن است نشان دهد، در طول ادغام، عناصر جریان ۱ ممکن است بین عناصر جریان ۲ درهم‌تنیده شوند. این موضوع توسط نتایج اجرای برنامه نشان داده شده است:

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]
  • خط ۳: فرآیند [process1] روی نخ محاسباتی [RxComputationThreadPool-4] اجرا می‌شود؛
  • خط ۴: فرآیند [process2] روی نخ محاسباتی [RxComputationThreadPool-5] در حال اجرا است؛
  • خط ۹: فرآیند [process12] در نخ محاسباتی [RxComputationThreadPool-3] مشاهده می‌شود. من قاعده‌ای را که به این انتخاب منجر شده است نمی‌دانم؛
  • خطوط ۹–۱۱: می‌بینیم که ناظر در حال مشاهده عناصر از هر دو فرآیند [process1] (خط ۵) و [process2] (خطوط ۶ و ۷) است، هرچند هیچ‌کدام از آن‌ها به پایان نرسیده‌اند (آمیختگی وجود دارد)؛
  • فرآیند [process12] (خط 17) زمانی خاتمه می‌یابد که هر دو فرآیند process1 و process2 خاتمه یافته باشند؛

7.5.2. مثال ۱۵: الحاق دو مشاهده‌پذیر با [Observable.concat]

اکنون کد زیر را بررسی خواهیم کرد:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, Schedulers.computation());
        //الحاق
        Process<?> process12 = new Process<>("process12",
                Observable.concat(process1.getObservable(), process2.getObservable()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process12);
    }
}
  • خطوط ۱۵–۱۷: فرایندی به نام [process1] سه عدد حقیقی را روی یک نخ محاسباتی خروجی خواهد داد. همچنین روی یک نخ محاسباتی مشاهده خواهد شد؛
  • خطوط ۱۹–۲۰: فرایندی به نام [process2] دو رشتهٔ کاراکتری را روی یک نخ نامشخص، در این مورد نخ اصلی پیش‌فرض، صادر خواهد کرد. این فرایند روی یک نخ محاسباتی مشاهده خواهد شد؛
  • خط ۲۳: دو فرآیند با هم الحاق می‌شوند، یعنی یک مشاهده‌پذیر ایجاد می‌شود که عناصر آن از هر دو فرآیند می‌آیند. مقادیر صادرشده با هم مخلوط نمی‌شوند. فرآیند [process12] ابتدا تمام مقادیر را از فرآیند [process1] صادر می‌کند، سپس مقادیر فرآیند [process2] را. برای این کار از متد استاتیک [Observable.concat] استفاده می‌شود:
 

نتایج اجرای آن به شرح زیر است:

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]
  • خطوط ۳–۱۰: فرآیند [process1] در حال اجرا است و فرآیند [process12] مقادیر تولید شده توسط [process1] را خروجی می‌دهد؛
  • خط ۹: فرآیند [process1] به پایان رسیده است؛
  • خطوط ۱۱–۱۷: فرآیند [process2] در حال اجرا است و فرآیند [process12] مقادیر خروجی [process2] را خروجی می‌دهد؛

در فرآیند process2 یک ناهنجاری وجود دارد: هیچ نخ اجرایی مشخص نشده است. بنابراین ممکن بود انتظار رود که نخ اصلی به‌طور پیش‌فرض استفاده شود. با این حال، این‌گونه نیست. رشته اجرایی، رشته محاسباتی [RxComputationThreadPool-3] (خط ۱۱) بود. بنابراین، هنگامی که هیچ رشته اجرایی یا مشاهده‌ای مشخص نشده باشد، نمی‌توان هیچ فرضی در مورد انتخاب شدن کدام رشته کرد.

7.5.3. مثال ۱۶: ترکیب دو متغیر مشاهداتی با [Observable.zip]

اکنون کد زیر را بررسی خواهیم کرد:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, null);
        // تابع ترکیب دو فرآیند
        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");
                }
            }
        };
        //فایل فشردهٔ دو فرایند
        Process<String> process12 = new Process<>("process12",
                Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process12);
    }
}
  • خطوط ۱۶–۱۸: فرایندی به نام [process1] سه عدد حقیقی را روی یک نخ محاسباتی خروجی خواهد داد. همچنین روی یک نخ محاسباتی مشاهده خواهد شد؛
  • خطوط ۲۰–۲۱: فرایندی به نام [process2] دو رشتهٔ کاراکتری را روی یک نخ محاسباتی آزاد (unbound) صادر می‌کند. نخ مشاهده نیز آزاد است؛
  • خطوط ۲۳–۳۲: نمونه‌سازی یک نوع [FuncN<String>] با یک کلاس ناشناس. FuncN یک رابط تابعی است:
 

متد [FuncN.call] یک آرایه از اشیاء را می‌پذیرد و یک نوع R را برمی‌گرداند. تابع [funcn] برای ترکیب فرآیندهای process1 و process2 به ترتیب استفاده خواهد شد. در روش [FuncN.call]:

  • args[0] یک Double خواهد بود؛
  • args[1] یک String خواهد بود؛

در اینجا، نتیجهٔ [funcn.call] رشتهٔ کاراکتری در خط ۲۷ خواهد بود. محاسبهٔ این نتیجه نیازمند دانستن نوع آرگومان‌های متد call نیست.

این دو فرایند به شرح زیر ترکیب می‌شوند:


//فایل زیپ شامل دو فرآیند
Process<String> process12 = new Process<>("process12",
Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));

روش [Observable.zip] به شرح زیر عمل می‌کند:

 

می‌توانیم ببینیم که:

  • آرگومان اول `zip` از نوع Iterable<Observable> است. در مثال ما، پارامتر واقعی از نوع List<Observable> است که از دو متغیر مشاهده‌پذیر ما تشکیل شده است؛
  • آرگومان دوم `zip` از نوع `FuncN` است. در مثال ما، پارامتر واقعی `[funcn]` است؛

اجرا نتایج زیر را به دست می‌دهد:

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]
  • خطوط ۷ و ۱۱: فرآیند process12 دو عنصر را صادر می‌کند؛
  • خط ۸: عنصر اضافی که توسط فرآیند process1 تولید شده و در فرآیند process2 شریکی ندارد، توسط فرآیند نتیجه process12 تولید نمی‌شود؛

می‌توانیم ببینیم که فرآیند process2، که به آن نه نخ اجرایی و نه نخ مشاهده‌ای اختصاص داده نشده بود، از نخ اصلی برای هر دو استفاده کرد.

7.5.4. مثال ۱۷: ترکیب دو مشاهد‌شدنی با [Observable.combineLatest]

اکنون کد زیر را بررسی خواهیم کرد:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null,
                Schedulers.computation());
        //ترکیب دو فرآیند
        Process<Double> process12 = new Process<>("process12",
                Observable.combineLatest(process1.getObservable(), process2.getObservable(), (d1, d2) -> d1 + d2));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process12);
    }
}
  • خطوط ۱۴–۱۶: فرایندی به نام [process1] سه عدد حقیقی را روی یک نخ محاسباتی صادر خواهد کرد. این فرایند همچنین روی یک نخ محاسباتی مشاهده خواهد شد؛
  • خطوط ۱۸–۲۰: فرایندی به نام [process2] دو عدد حقیقی را روی یک نخ نامحدود منتشر می‌کند. این‌ها روی یک نخ محاسباتی مشاهده خواهند شد؛
  • خط ۲۳: دو متغیر مشاهداتی با استفاده از روش استاتیک زیر [Observable.combineLatest] ترکیب می‌شوند:
 

مشاهده‌پذیر [combineLatest] به شرح زیر عمل می‌کند: هنگامی که یکی از دو مشاهده‌پذیر یک عنصر E1 را منتشر می‌کند، این عنصر توسط [combineFunction] با آخرین عنصر منتشر شده توسط مشاهده‌پذیر دیگر ترکیب می‌شود.

اجرای این کد نتیجه زیر را تولید می‌کند:

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]
  • خط ۵: انتقال از process2 (۵۶) با آخرین عنصر ارسال‌شده توسط process1 (۵۴، خط ۴) ترکیب شده و نتیجهٔ نمایش‌داده‌شده در خط ۷ را تولید می‌کند؛
  • خط ۶: خروجی از process1 (51.6) با آخرین عنصر خروجی شده توسط process2 (56، خط ۵) ترکیب شده و نتیجه نشان داده شده در خط ۸ را تولید می‌کند؛
  • خط ۹: انتقال از process2 (261.8) با آخرین عنصر ارسال‌شده توسط process1 (51.6، خط ۶) ترکیب شده و نتیجه نشان‌داده‌شده در خط ۱۲ را تولید می‌کند؛
  • خط ۱۳: ارسال از process1 (۸۰.۳۹) با آخرین عنصر ارسال‌شده توسط process2 (۲۶۱.۸، خط ۹) ترکیب شده و نتیجه نشان‌داده‌شده در خط ۱۵ را تولید می‌کند؛

این یک واریانت از مشاهده‌پذیر [zip] است، که در این مورد، عناصر ترکیبی لزوماً عناصر در همان موقعیت در جریان‌ها نیستند. در اینجا باید توجه داشت که فرآیند process2، که برای آن هیچ نخ اجرایی مشخص نشده بود، در اینجا روی نخ اصلی اجرا شد (خط ۲).

7.5.5. مثال ۱۸: ترکیب دو مشاهده‌پذیر با [Observable.amb]

اکنون کد زیر را بررسی خواهیم کرد:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null, null);
        //ترکیب دو فرآیند
        Process<Double> process12 = new Process<>("process12",
                Observable.amb(process1.getObservable(), process2.getObservable()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process12);
    }
}
  • خطوط ۱۴–۱۶: فرایندی به نام [process1] سه عدد حقیقی را روی یک نخ محاسباتی صادر خواهد کرد. این اعداد همچنین روی یک نخ محاسباتی مشاهده خواهند شد؛
  • خطوط ۱۸–۲۰: فرایندی به نام [process2] دو عدد حقیقی را در یک نخ محاسباتی بدون محدودیت صادر می‌کند. آن‌ها در یک نخ محاسباتی بدون محدودیت مشاهده خواهند شد؛
  • خط ۲۲: دو مشاهده‌پذیر با استفاده از روش استاتیک زیر [Observable.amb] ترکیب می‌شوند:
 

همان‌طور که در نمودار بالا نشان داده شده است، مشاهده‌پذیر [Observable.amb(Observable o1, Observable o2)] عناصر مشاهده‌پذیری را که ابتدا صادر می‌شوند، منتشر می‌کند. این موضوع با نتایج مثال ارائه‌شده تأیید می‌شود:

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]
  • خط ۴: این فرآیند process2 است که ابتدا صادر می‌کند؛
  • خطوط ۸ و ۱۲: فرآیند process12 تمام عناصری را که توسط فرآیند process2 (خطوط ۴ و ۱۱) منتشر شده‌اند، منتشر می‌کند؛

7.6. زنجیره پردازش برای یک مشاهده‌پذیر

7.6.1. مثال ۱۹: تبدیل یک مشاهده‌پذیر با [Observable.map]

در مثال‌های قبلی، ما ترکیب‌های مختلفی از دو مشاهده‌پذیر را برای تشکیل یک مشاهده‌پذیر سوم بررسی کردیم. اکنون روش‌های ایستا کلاس [Observable] را معرفی می‌کنیم که امکان عملیات تبدیل، فیلتر کردن و تجمیع را بر روی یک مشاهده‌پذیر فراهم می‌کنند. در اینجا روش‌هایی مشابه روش‌های کلاس [Stream] که در بخش ۵ مطالعه شد، خواهیم یافت.

نخستین مثال ما به شرح زیر خواهد بود:


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 {
        // فرآیند ۱
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<String> process2 = new Process<>("process2",
                process1.getObservable().map(d -> String.format("valeur-%s", d)));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}
  • خطوط 14–16: فرایندی به نام process1 سه عدد حقیقی را در یک نخ محاسباتی منتشر خواهد کرد. همچنین روی یک نخ محاسباتی مشاهده خواهد شد؛
  • خطوط 17–18: اعداد خروجی process1 در فرایندی به نام process2 به رشته‌های کاراکتری تبدیل خواهند شد؛
  • خط ۲۰: process2 مشاهده می‌شود؛

متد [Observable.map] در خط ۱۸ مشابه متد [Stream.map] مورد بحث در بخش 5.5 است:

 

نتایج مثال به شرح زیر است:

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]
  • خطوط ۴، ۵ و ۸: خروجی‌های process1. این‌ها اعداد حقیقی هستند؛
  • خطوط ۶، ۷ و ۱۰: خروجی‌های process2 که مشاهده شده‌اند. این‌ها رشته‌های کاراکتری هستند؛

7.6.2. مثال-20: فیلتر کردن یک مشاهده‌پذیر با [Observable.filter]

مثال به شرح زیر خواهد بود:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}
  • خطوط ۱۱–۱۲: فرایندی به نام process1 اعداد صحیح ۰ تا ۲ را در یک نخ محاسباتی صادر می‌کند. همچنین در یک نخ محاسباتی مشاهده خواهد شد؛
  • خط 14: اعداد تولید شده توسط process1 فیلتر می‌شوند تا فقط اعداد زوج در process2 باقی بمانند؛
  • خط ۲۰: process2 نظارت می‌شود؛

روش [Observable.filter] در خط ۱۸ مشابه روش [Stream.filter] مورد بحث در بخش 5.4 است:

 

نتایج مثال به شرح زیر است:

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]
  • خطوط ۴، ۵ و ۷: ارسال‌ها از process1؛
  • خطوط ۶ و ۹: انتشارهای مشاهده‌شده از process2. این‌ها عناصر زوج process1 هستند؛

7.6.3. مثال ۲۱: تبدیل یک قابل مشاهده با [Observable.flatMap]

مثال به شرح زیر خواهد بود:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}
  • خطوط ۱۲–۱۳: فرایندی به نام process1 اعداد صحیح ۰ تا ۲ را در یک نخ محاسباتی منتشر می‌کند. همچنین در یک نخ محاسباتی مشاهده خواهد شد؛
  • خطوط ۱۵–۱۸: هر عددی n که توسط process1 صادر می‌شود، به یک مشاهده‌پذیر تبدیل می‌شود که سه عدد (10*n, 10*n+1, 10*n+2) را صادر می‌کند. اگر در خط ۱۵، متد [map] استفاده می‌شد، process2 به جای نوع Integer، نوع Observable<Integer> را صادر می‌کرد. روش [flatMap] به‌کاررفته اجازه می‌دهد (flatten) این دنباله از عناصر از نوع Observable<Integer> را به دنباله‌ای از عناصر از نوع Integer تبدیل می‌کند که شامل هر عنصر از هر یک از Observable<Integer> است؛
  • خط ۲۰: process2 مشاهده می‌شود؛

روش [Observable.flatMap] در خط ۱۵ مشابه روش [Stream.flatMap] مورد بحث در بخش 5.6.12 است:

 

نتایج مثال به شرح زیر است:

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]
  • خطوط ۵–۷: سه انتقال process2 پس از انتقال خط ۴ از process1؛
  • خطوط ۹–۱۱: سه انتقال process2 پس از انتقال خط ۸ از process1؛
  • خطوط 14–16: سه انتقال process2 پس از انتقال خط 12 از process1;

کد زیر نشان می‌دهد چگونه یک نوع Observable<Integer[]> را از process1 و [Exemple21b] ایجاد کنیم:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Integer[]> process2 = new Process<>("process2", process1.getObservable().map(i -> {
            int value = i * 10;
            return new Integer[] { value, value + 1, value + 2 };
        }));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}
  • خط ۱۴: از متد [Observable.map] استفاده می‌شود؛
  • خط 16: که یک نوع Integer[] را بازمی‌گرداند؛

نتایج به شرح زیر است:

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]
  • خطوط ۶، ۷، ۱۰: نتایج map را مشاهده می‌کنیم؛

تمام این تبدیل‌های قابل مشاهده را می‌توان پشت سر هم زنجیر کرد، زیرا هر تبدیل یک مشاهده‌پذیر جدید تولید می‌کند. این موضوع با مثال زیر [Exemple21c] نشان داده شده است:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        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));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}
  • خطوط ۱۵–۱۸: flatMap با یک filter دنبال می‌شود؛

نتایج اجرا به شرح زیر است:

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]
  • خطوط ۸–۱۳: process2 تنها عناصر زوج را از flatMap خروجی داده است؛

یک روش مشابه [flatMap]، روش [flatMapIterable] است که با مثال زیر نشان داده شده است، [Exemple21d]:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        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));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}

در خط ۱۶، به جای استفاده از روش [flatMap]، از روش [flatMapIterable] استفاده می‌شود. در این مورد، تابع تبدیل باید به جای نوع Observable<T>، نوع Iterable<T> را تولید کند (خط ۱۸).

این همان نتایج قبلی را تولید می‌کند.

بیایید به تعریف متد [flatMap] بازگردیم:

 

همانطور که در بالا مشاهده می‌شود، یک عنصر آبی [3] بین دو عنصر سبز [1-2] قرار گرفته است. این بدان معناست که هنگام مسطح‌سازی عناصر Observable<T>، متد [flatMap] به ترتیبی که این مشاهدات داخلی مختلف صادر شده‌اند، احترام می‌گذارد. این موضوع با مثال زیر نشان داده شده است: [Exemple21e]:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // فرآیند ۳
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().flatMap(i -> process2.getObservable()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process3);
    }
}
  • خطوط ۱۱–۱۲: فرآیند process1 اعداد صحیح [0,1] را خروجی می‌دهد؛
  • خطوط ۱۴–۱۵: فرآیند process2 اعداد صحیح [10,11,12] را منتشر می‌کند؛
  • خطوط 17–18: هر عنصری که توسط process1 منتشر می‌شود با مشاهده‌پذیر فرآیند process2 مرتبط است. این بدان معناست که:
    • عنصر [0] از process1 با یک مشاهده‌پذیر که [10,11,12] را منتشر می‌کند، مرتبط خواهد بود؛
    • همین امر در مورد عنصر 1 نیز صدق می‌کند؛

در نهایت، ۶ عدد [10, 11, 12, 10, 11, 12] منتشر خواهند شد. می‌خواهیم ببینیم در چه ترتیبی.

نتایج اجرای برنامه به شرح زیر است:

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]

می‌توانیم ببینیم که ترتیبی که فرآیند process3 اعداد را تولید کرد، به این صورت بود: [10, 10, 11, 12, 11, 12] (خطوط ۱۱، ۱۲، ۱۴، ۱۷، ۱۹، ۲۲). بنابراین، عناصر خروجی‌شده توسط فرآیند process2 در واقع با هم مخلوط شده‌اند. این مشکل را می‌توان با استفاده از متد [concatMap] به جای متد [flatMap] حل کرد. این موضوع در کد زیر برای [Exemple21ef] نشان داده شده است:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // فرآیند ۳
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().concatMap(i -> process2.getObservable()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process3);
    }
}

در خط ۱۸، [flatMap] با [concatMap] جایگزین شده است. نتایج اجرای آن به شرح زیر است:

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]

می‌توانیم ببینیم که ترتیبی که فرآیند process3 با آن اجرا شد، به این صورت بود: [10, 11, 12, 10, 11, 12] (خطوط ۱۲–۱۴، ۱۷، ۱۹، ۲۲). عناصر خروجی‌شده توسط فرآیند process2 به‌صورت درهم‌تنیده نبودند.

یک واریانت دیگر از روش [map]، روش [switchMap] است:

 

در بالا، از مشاهده‌پذیر [1]، سه مشاهده‌پذیر دیگر [2]، که هر کدام شامل دو عنصر هستند، تولید می‌شوند؛ این‌ها سپس همانند [flatMap] و [3] مسطح می‌شوند. مشاهده می‌شود که نتیجه به جای ۶ عنصر، دارای ۵ عنصر است. این به آن دلیل است که پیش از آنکه مشاهده‌پذیر دوم عنصر دوم خود، [6]، را منتشر کند، مشاهده‌پذیر سوم عنصر اول خود، [5]، را منتشر می‌کند، که به این معناست که مشاهده‌پذیر دوم نادیده گرفته می‌شود. در نتیجه، عنصر [6] در مشاهده‌پذیر حاصل [3] ظاهر نمی‌شود.

برای روشن شدن [switchMap]، از مثال زیر استفاده می‌کنیم: [Exemple21eg]:


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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // فرآیند ۳
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().switchMap(i -> process2.getObservable()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process3);
    }
}

اجرای مثال نتایج زیر را تولید می‌کند:

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 دو عنصر را منتشر می‌کند که منجر به دو مشاهده‌پذیر process2 می‌شوند، که هر کدام از سه عنصر تشکیل شده‌اند؛
  • خط 14: ناظر عنصر شماره 0 را که توسط مشاهده‌پذیر اول process2 در خط 6 منتشر شده است، دریافت می‌کند؛
  • خط 15: ناظر عنصر شماره 0 را که از مشاهده‌گر دوم process2 در خط 13 منتشر شده است، دریافت می‌کند. مشخص نیست چرا ناظر پیش از این عناصر شماره ۱ و ۲ را که توسط مشاهده‌پذیر اول، process2، در خطوط ۷ و ۸ ارسال شده بودند، دریافت نکرده است. در هر صورت، مشاهده‌پذیر اول، process2، کنار گذاشته می‌شود؛
  • در نهایت، ناظر تنها ۴ عنصر (خطوط ۱۴، ۱۵، ۱۷ و ۲۰) را به جای ۶ عنصری که ارسال شده بود، مشاهده می‌کند؛

7.6.4. مثال‌ها-22: سایر متدهای کلاس [Observable]

کلاس [Observable] بسیاری از متدهای کلاس [Stream] را که به شیوه‌ای مشابه عمل می‌کنند، در خود جای داده است. در اینجا به چند نمونه از آن‌ها اشاره می‌شود. ما صرفاً کد و نتایج آن را ارائه خواهیم داد.

[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 {
        // فرآیند
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).take(3));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}

نتایج

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 {
        // فرآیند
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).takeLast(2));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}

نتایج

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 {
        // پردازش‌ها
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).skip(5).take(2));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}

نتایج

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 {
        // پردازش‌ها
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).reduce(0, (i, a) -> i + a));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}
  • خط ۱۰: مجموع عناصر مشاهده‌پذیر را محاسبه می‌کند. نتیجه یک مشاهده‌پذیر است که این مجموع را منتشر می‌کند؛

نتایج

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 {
        // پردازش‌ها
        Process<Boolean> process = new Process<>("process", Observable.range(1, 10).all(i -> i > 10));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}
  • خط ۱۰: یک Observable<Boolean> را بازمی‌گرداند که عنصر true را صادر می‌کند، اگر گزاره روش [all] برای همه عناصر درست باشد، در غیر این صورت false؛

نتایج

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 {
        // پردازش‌ها
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).count());
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}
  • خط ۱۰: [Observable.count] یک مشاهده‌پذیر تک‌عنصری ایجاد می‌کند که جمع عناصر مشاهده‌شده است؛

نتایج

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 {
        // پردازش‌ها
        Process<Integer> process = new Process<>("process", Observable.just(1, 2, 1, 3).distinct());
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}

نتایج

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 {
        // فرآیندها
        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()));
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process);
    }
}
  • خط ۱۱: متد [groupBy] ده عنصر خروجی را به دو گروه تقسیم می‌کند: اعداد زوج و اعداد فرد. نتیجه از نوع `Observable<GroupedObservable<Boolean, Integer>>` است، یعنی یک observable که عناصر آن از نوع `GroupedObservable<Boolean, Integer>` هستند، که در آن `Boolean` نوع کلید گروه است (در این مورد false و true) و همچنین نوع نتیجه لَمبدا (lambda) است که به عنوان پارامتر به متد [groupBy] پاس می‌شود، در حالی که Integer نوع عناصر گروه است؛
  • خط ۱۲: نوع GroupedObservable دارای متدی به نام [asObservable] است که امکان ایجاد یک مشاهده‌پذیر از این نوع را فراهم می‌کند. بنابراین ما دو نوع خواهیم داشت، Observable<Integer>، یکی برای اعداد زوج و دیگری برای اعداد فرد. از این دو متغیر مشاهدنی، متد [concatMap] یک مورد واحد ایجاد خواهد کرد؛

نتایج

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 {
        // فرآیند ۱
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // فرآیند ۲
        Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
        // اشتراک‌ها
        ProcessUtils.subscribe(1, process2);
    }
}
  • در خط ۱۵، متد [timestamp] یک زمان را به هر عنصر از مشاهده‌پذیر پردازش‌شده مرتبط می‌کند؛

نتایج

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]

در این مثال، دشوار است بگوییم اطلاعات timestamp نمایانگر چیست:

  • خطوط ۴–۵: می‌بینیم که عنصر ۱ از process1 در ۱۳۹ میلی‌ثانیه پس از عنصر ۰ منتقل شده است؛
  • خطوط ۶ و ۷: می‌توانیم ببینیم که عنصر ۱ از process2، ۲۳۴ میلی‌ثانیه پس از عنصر ۰ مشاهده شده است؛
  • خطوط ۵ و ۸: می‌بینیم که عنصر ۲ از process1 با تأخیر ۳۳ میلی‌ثانیه‌ای پس از عنصر ۱ ارسال شده است؛
  • خطوط ۷ و ۱۰: می‌بینیم که عنصر ۲ از process2، ۲۳ میلی‌ثانیه پس از عنصر ۱ مشاهده شده است؛

این تأخیرهای زمانی به این دلیل است که نخ‌های مسئول مشاهده و اجرای مشاهدات یکسان نیستند. اگر خطوط ۱۲–۱۳ را با خطوط زیر (Example22j) جایگزین کنیم:


// فرآیند ۱
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
  • خطوط ۲–۳: نخ مشاهده مشخص نشده است. می‌دانیم که در این مورد، مشاهده‌شدنی در جایی که اجرا می‌شود مشاهده می‌شود؛

این نتایج زیر را می‌دهد:

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]
  • خطوط ۴ و ۶: فرآیند process1 عنصر شماره ۱ خود را ۵۸۷ میلی‌ثانیه پس از عنصر شماره ۰ خود منتشر می‌کند؛
  • خطوط ۵ و ۷: ناظر این دو عنصر را با اختلاف زمانی ۵۸۶ میلی‌ثانیه مشاهده می‌کند؛
  • خطوط ۶ و ۸: فرآیند process1 عنصر شماره ۲ خود را ۳۹۶ میلی‌ثانیه پس از عنصر شماره ۱ خود منتشر می‌کند؛
  • خطوط ۷ و ۹: ناظر این دو عنصر را با اختلاف زمانی ۳۹۶ میلی‌ثانیه مشاهده می‌کند؛

در اینجا، مقادیر مربوط به timestamp سازگار هستند: آنها تاریخ انتقال عنصر را به دقت نشان می‌دهند.

7.7. برنامه‌ریزها

7.7.1. مثال ۲۳: زمان‌بندی‌کننده [Schedulers.computation]

اکنون برنامه‌ریزهای اجرا را بررسی خواهیم کرد. تحلیل بر روی نخ اجرای برنامه متمرکز خواهد بود.

موضوع برنامه‌ریزها تا حدی مبهم است. برنامه‌ریزهای مختلف در این سؤال در وب‌سایت StackOverflow [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases] ارائه شده‌اند:

 

ما تلاش خواهیم کرد تا با مثال‌ها کاربرد این برنامه‌ریزهای مختلف را نشان دهیم. اولین مثال برنامه‌ریز [Schedulers.computation] را نشان می‌دهد:


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 {
        // فرآیندها
        @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);
        }
        // اشتراک‌ها
        ProcessUtils.subscribe(1, processes);
    }
}
  • خطوط 14–19: یک آرایه از 10 فرآیند که روی یک نخ محاسباتی در حال اجرا هستند ایجاد می‌شود؛
  • خط ۱۷: هر فرآیند یک عدد حقیقی تصادفی تولید می‌کند؛
  • خط ۲۱: ما به همه این فرآیندها مشترک می‌شویم؛

نتایج به شرح زیر است:

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]
  • خطوط ۲–۱۰: ۸ فرآیند اول روی ۸ نخ مختلف شروع به کار می‌کنند (ماشین مورد استفاده دارای ۸ هسته است). مشاهده می‌شود که همگی تقریباً هم‌زمان شروع می‌شوند؛
  • خطوط 17–19: سه فرآیند خاتمه می‌یابند و در نتیجه سه نخ آزاد می‌شوند؛
  • خطوط ۲۳–۲۴: دو فرآیند آخر می‌توانند سپس با استفاده از دو تا از رشته‌های آزادشده، آغاز شوند؛

بنابراین می‌توان نتیجه گرفت که برنامه‌ریز [Schedulers.computation] یک مجموعه از n نخ را فراهم می‌کند، که در آن n تعداد هسته‌های موجود در ماشین است. این نخ‌ها به صورت موازی روی این هسته‌ها اجرا می‌شوند.

7.7.2. مثال ۲۴: برنامه‌ریز [Schedulers.io]

ما کد قبلی را با استفاده از زمان‌بندی‌کننده [Schedulers.io] اجرا می‌کنیم:


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 {
        // فرآیندها
        @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);
        }
        // اشتراک‌ها
        ProcessUtils.subscribe(1, processes);
    }
}
  • خط ۱۸: فرآیندها با استفاده از رشته‌های زمان‌بندی‌کننده [Schedulers.io] اجرا می‌شوند؛

این نتایج زیر را تولید می‌کند:

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]
  • خطوط ۲–۱۰: هر یک از ۱۰ فرآیند روی یک رشته (thread) مجزا شروع می‌شوند. برخلاف مورد قبلی، همه فرآیندها قادر به راه‌اندازی بودند. شایان ذکر است که این راه‌اندازی‌ها ۶ میلی‌ثانیه طول کشید، در حالی که قبلاً ۱ میلی‌ثانیه طول می‌کشید؛
  • خطوط ۱۳–۱۸: مشاهده‌شدنی‌ها به‌صورت پشت سر هم صادر می‌شوند، نه تقریباً به‌صورت موازی مانند مورد قبلی؛

تفاوت بین برنامه‌ریزهای [Schedulers.io] و [Schedulers.computation] چیست؟ پاسخ را می‌توان در URL و [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases] یافت:

 

7.7.3. مثال ۲۵: زمان‌بندی‌کننده [Schedulers.newThread]

ما کد قبلی را با استفاده از زمان‌بندی‌کننده [Schedulers.newThread] اجرا می‌کنیم:


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 {
        // پردازش‌ها
        @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);
        }
        // اشتراک‌ها
        ProcessUtils.subscribe(1, processes);
    }
}

نتایج به‌دست‌آمده با برنامه‌ریز [Schedulers.io] یکسان است:

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 و [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io] توضیح داده شده است که زمان‌بندی‌کننده [Schedulers.io] یک استخر نخ فراهم می‌کند، در حالی که زمان‌بندی‌کننده [Schedulers.newThread] این کار را انجام نمی‌دهد. یک استخر نخ به‌طور خودکار n نخ ایجاد می‌کند. این نخ‌ها را به فرآیندهایی که به آن‌ها نیاز دارند اختصاص می‌دهد. وقتی این فرآیندها به پایان می‌رسند، نخ‌هایشان حذف نمی‌شوند، بلکه به استخر بازمی‌گردند و می‌توانند توسط فرآیند دیگری دوباره استفاده شوند. این روش نسبت به ایجاد و حذف مداوم نخ‌ها کارآمدتر است. بنابراین، به نظر می‌رسد استفاده از برنامه‌ریز [Schedulers.io] ترجیح دارد.

7.7.4. مثال ۲۶: برنامه‌ریزهای [Schedulers.immediate, Schedulers.trampoline]

بیایید به توضیحات ارائه‌شده برای این دو برنامه‌ریز بازگردیم:

 

درک این توضیح نسبتاً ساده است، اما وقتی سعی می‌کنید آن را تشریح کنید، متوجه می‌شوید که در واقع آن را به درستی نفهمیده‌اید. کتاب [Learning Reactive Programming With Java 8] بود که به من امکان داد مثالی بر اساس یکی از مثال‌های موجود در آن کتاب، اما به صورت ساده‌تر، ایجاد کنم. این مثال به شرح زیر است:


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 {

        //یک برنامه‌ریز
        Scheduler scheduler = Schedulers.immediate();
        //یک کارگر برای این زمان‌بندی‌کننده
        Worker worker = scheduler.createWorker();
        // نوع Action0 برای اجرا روی کارگر
        Action0 action02 = new Action0() {
            @Override
            public void call() {
                // لاگ Action02
                ProcessUtils.showInfos.accept("action02");
            }
        };

        // نوع Action0 برای اجرا روی کارگر
        Action0 action01 = new Action0() {
            @Override
            public void call() {
                // یک اقدام جدید روی همان کارگر زمان‌بندی شده است
                worker.schedule(action02);
                // لاگ Action01
                ProcessUtils.showInfos.accept("action01");
            }
        };
        //action01 روی کارگر زمان‌بندی شده است
        worker.schedule(action01);
    }

    // دیدها
    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()));

}
  • خط ۱۷: یک زمان‌بندی‌کننده. این می‌تواند یا [Schedulers.immediate]، همان‌طور که در اینجا نشان داده شده است، یا [Schedulers.trampoline] در مرحله بعد باشد؛
  • خط ۱۹: عملیات از نوع Action0 (خطوط ۲۱ و ۲۰) می‌توانند روی کارگران برنامه‌ریز اجرا شوند. متد [Scheduler.createWorker] برای ایجاد یک کارگر استفاده می‌شود. متد [Worker.schedule(Action0)] برای وادار کردن یک کارگر به اجرای نوع Action0 استفاده می‌شود؛
  • خطوط ۲۱–۲۷: یک اقدام اول به نام [action02] که توسط کارگر در خط ۱۹ اجرا خواهد شد (خط ۴۰)؛
  • خطوط ۳۰–۳۸: یک اقدام دوم به نام [action01]. ویژگی متمایز آن این است که باعث اجرای اقدام action02 بر روی همان کارگری که خود آن قرار دارد، می‌شود (خط ۳۴). اینجا جایی است که تفاوت بین [Schedulers.immediate] و [Schedulers.trampoline] نهفته است:
    • اگر زمان‌بندی‌کننده [Schedulers.immediate] باشد، در خط ۳۴، اقدام action02 بلافاصله اجرا خواهد شد (از این رو نام زمان‌بندی‌کننده) و اقدام در حال اجرای action01 متوقف خواهد شد. سپس پیام خط 25 نمایش داده می‌شود. پس از پایان اقدام action02، اقدام action01 از سر گرفته شده و پیام خط 36 نمایش داده می‌شود؛
    • اگر برنامه‌ریز [Schedulers.trampoline] باشد، در خط ۳۴، اقدام action02 به حالت تعلیق درمی‌آید. این اقدام تنها پس از اتمام وظیفه جاری action01 اجرا خواهد شد. سپس پیام خط ۳۶ را مشاهده خواهیم کرد. پس از اتمام اقدام action01، اقدام action02 اجرا می‌شود و پیام خط ۲۵ را مشاهده خواهیم کرد؛

اجرای کد بالا نتایج زیر را تولید می‌کند:

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

اگر در خط 17 از زمان‌بندی‌کننده [Schedulers.trampoline] استفاده کنیم، نتایج متضادی به دست می‌آوریم:

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

با این حال، درک ارتباط با observables دشوار است. من مثال قانع‌کننده‌ای نیافته‌ام که بتواند مزیت اجرای یک observable را روی یکی از این دو thread نشان دهد. با این حال، یک مثال وجود دارد، هرچند اصلاً طبیعی به نظر نمی‌رسد:


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 {

        // Worker
        Worker worker = Schedulers.immediate().createWorker();
        // Worker worker = Schedulers.trampoline().createWorker();
        //قابل مشاهده ۱ روی کارگر
        worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {

            @Override
            public void call(Integer i) {
                ProcessUtils.showInfos.accept(String.valueOf(i));
                // مشاهده‌پذیر ۲ روی همان کارگر
                worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
                    @Override
                    public void call(Integer i) {
                        ProcessUtils.showInfos.accept(String.valueOf(i));
                    }
                }));
            }
        }));
    }
}
  • خطوط ۱۳–۱۴: یک کارگر با استفاده از یکی از دو زمان‌بندی‌کننده، [Schedulers.immediate] و [Schedulers.trampoline]، ایجاد می‌شود؛
  • خط 16: یک مشاهده‌پذیر اول، obs1، روی این کارگر زمان‌بندی می‌شود تا اعداد [1,2] را منتشر کند
  • خط ۲۲: هر بار که یک عنصر از این مشاهد‌شدنی obs1 مشاهده می‌شود، مشاهده یک مشاهد‌شدنی دوم obs2 روی همان کارگر برای تولید اعداد [100,101] تحریک می‌شود؛

با برنامه‌ریز [Schedulers.immediate]، نتایج زیر به‌دست می‌آید:

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]

در حالی که با برنامه‌ریز [Schedulers.trampoline]، نتایج زیر به دست می‌آیند:

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

هنوز کارهای زیادی باقی مانده است. برای بررسی عمیق‌تر کتابخانه RxJava، از خواننده دعوت می‌شود تا با استفاده از مراجع ارائه‌شده در ابتدای این سند، به آموزش خود ادامه دهد. با این حال، اکنون ما اصول اولیه برای استفاده از RxJava در محیط‌های Swing و Android را داریم. این چیزی است که اکنون به نمایش خواهیم گذاشت.