7. مكتبة RxJava
تعتمد مكتبة RxJava على المفهوم التالي: يتم مراقبة تدفق العناصر من النوع T Observable<T> بواسطة مشترك واحد أو أكثر (مشتركين، مراقبين، مستهلكين) Subscriber<T>. تسمح المكتبة RxJava بتشغيل تدفق Observable<T> في مؤشر ترابط T1 ومراقبه Subscriber<T> في مؤشر ترابط T2 دون أن يضطر المطوريضطر إلى القلق بشأن إدارة دورة حياة هذه الخيوط والمشكلات الصعبة بطبيعتها، مثل مشاركة البيانات بين الخيوط ومزامنتها لتنفيذ مهمة شاملة. وبالتالي، فهي تسهل البرمجة غير المتزامنة.
ينتج تدفق Observable<T> عناصر من النوع T، والتي يمكن ملاحظتها فور إنتاجها. إذا كان المراقب والمراقب (وهو مصطلح يُستخدم بشكل غير دقيق للإشارة إلى النوع Observable<T>) موجودين في نفس الخيط، فإن المراقب لا يمكنه إنتاج العنصر (i+1) إلا بعد أن يكون المراقب قد استهلك العنصر i. وهناك حالات قليلة تكون فيها هذه البنية ذات فائدة. إذا لم يكن المراقب والمراقب فيهما في نفس الخيط، فإن المراقب فيهما والمراقب يتصرفان بشكل مستقل: المراقب فيهما ينتج وفقًا لوتيرته الخاصة، والمراقب يستهلك وفقًا لوتيرته الخاصة. وهنا تكمن فائدة المكتبة. لقد تحدثنا حتى الآن عن مراقب واحد. في الواقع، يمكن أن يكون للمراقب فيهما أي عدد من المراقبين.
تتناسب مكتبة RxJava بشكل خاص مع البنية التي تمت مناقشتها في الفقرة 2 من «الاستكشاف» والتي نذكرها هنا:

- في [1]، توفر طبقة الخدمات خدمات يستغرق الحصول على بعضها وقتًا طويلاً (طلبات الشبكة على سبيل المثال)؛
- يتم استدعاء طبقة الخدمات هذه بواسطة واجهة مستخدم رسومية [1] (Swing، Android، JavaFx). إذا تم تنفيذ طبقة الخدمات في نفس الخيط (thread) الذي تعمل فيه الطريقة [swing] التي تستخدمها، فإن الواجهة الرسومية تتجمد (لا تستجيب) أثناء انتظار نتيجة الخدمة؛
- في [2]، تسمح طبقة تكييف رفيعة مُنفَّذة باستخدام RxJava بتقديم تنفيذ غير متزامن لنفس الخدمة إلى الطبقة الرسومية: حيث يمكن تنفيذ هذه الخدمة في مؤشر ترابط مختلف عن مؤشر ترابط الطريقة في الطبقة الرسومية التي تستدعيها. في هذه الحالة، تظل الواجهة الرسومية [3] سريعة الاستجابة: يمكن للمستخدم الاستمرار في التفاعل معها، على سبيل المثال إطلاق طلب شبكة جديد بالتوازي مع الطلب الأول، والأهم من ذلك، يمكن منحه إمكانية إلغاء العمليات التي تستغرق وقتًا طويلاً، وهو أمر مستحيل إذا كانت الواجهة الرسومية متجمدة؛
- يُعد الاستدعاء [4] متزامنًا، في حين أن الاستدعاء [5-6] غير متزامن؛
في هذه البنية، توفر الطبقة [2] خدمات تُرجع أنواعًا من Observable<T> يمكن لأساليب الطبقة الرسومية [3] الاشتراك فيها. ثم تقوم إحدى خدمات الطبقة [2] بتسليم نتائجها واحدة تلو الأخرى، ويمكن للطبقة [3] الاستجابة لكل منها، على سبيل المثال عن طريق تحديث مكون واحد أو أكثر من مكونات واجهة المستخدم الرسومية.
تحتوي فئة Observable<T> على عشرات الطرق. وهذه إحدى الصعوبات التي تنطوي عليها المكتبة: فهي غنية جدًّا ويصعب استيعاب جميع إمكانياتها. سنعرض بعضًا منها هنا. أما إتقان الطرق الأخرى فسيأتي لاحقًا مع مرور الوقت.
7.1. إنشاء العناصر القابلة للمراقبة والاشتراك فيها
7.1.1. المثال-01: الطريقة [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");
}
});
}
}
- السطر 12: يتم إنشاء نوع Observable<Integer> من قائمة من الأعداد الصحيحة.
الفئة Observable<T> هي تدفق لعناصر من النوع T يمكن ملاحظتها، ويفضل أن يكون ذلك بشكل غير متزامن ولكن ليس بالضرورة، كلما تم إنتاجها. وتعريفها كما يلي:
![]() |
كما سبق ذكره، تحتوي الفئة Observable<T> على عشرات الطرق. بعضها مشابه لتلك الموجودة في الفئة Stream<T> التي تمت دراستها في الفقرة 5. تتضمن وثائق RxJava «مخططات الرخام» [2] التي توضح كيفية عمل هذه الطرق:
- يوضح السطر 3 انبعاثات المتغير القابل للرصد بمرور الوقت؛
- يتم تطبيق الطريقة [4] على العناصر المنبعثة من المقياس. وتنتج هذه الطريقة عمومًا مقياسًا جديدًا؛
- يُظهر السطر 5 المقياس الجديد الذي تم الحصول عليه؛
الطريقة [Observable.from] لها التوقيع التالي:
![]() |
تسمح الطريقة الثابتة [Observable.from] بإنشاء Observable<T> من مجموعة عناصر من النوع T. وهي طريقة بسيطة جدًا للبدء في استخدام العناصر القابلة للمراقبة. السطر:
Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
ستقوم بإصدار ثلاثة عناصر. وهي لا تصدرها على الفور. بل ستصدرها بالكامل في كل مرة يعلن فيها مراقب عن نفسه. وهذا ما يُسمى بالمُراقب البارد. يقوم المُراقب بإعادة إصدار عناصره لكل مشترك جديد.
يمكن اعتبار التعليمات السابقة بمثابة إجراء لتكوين المراقب. يتم تكوينه مرة واحدة ويتم تنفيذه 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");
}
});
- السطر 1: يتم تجاهل النتيجة من النوع [Subscription]؛
- الأسطر 1-15: المعلمات الثلاثة هي مثيلات لفئات مجهولة. سنستخدم أيضًا دال لامدا. وتكمن فائدة الفئات المجهولة في أنها تتيح رؤية أنواع البيانات المتوقعة من الطريقة الوحيدة لهذه الفئات بوضوح؛
- الأسطر 2-5: تنفيذ المعلمة الأولى من النوع [Action1<Integer>]؛
- الأسطر 6-10: تنفيذ المعلمة الثانية من النوع [Action1<Throwable>]؛
- الأسطر 11-15: تنفيذ المعلمة الثالثة من النوع [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");
}
});
}
}
يبدأ العنصر القابل للمراقبة في السطر 12 في إصدار عناصره الثلاثة بمجرد استدعاء الأسلوب [subscribe] في السطر 14. ومنذ تلك اللحظة:
- مع كل عنصر يتم إرساله، يتم تنفيذ الأسطر 15-18.
- عند انتهاء العناصر الثلاثة، يتم تنفيذ الأسطر 24-29؛
- لن يتم تنفيذ الأسطر 19-24 أبدًا لأن المراقب لا يصدر استثناءً هنا؛
بشكل افتراضي، يتم تنفيذ المراقب والمراقب في نفس الخيط. توجد بعض المراقبين المُعرَّفين مسبقًا التي يتم تنفيذها في خيط مختلف عن الخيط الرئيسي (هنا خيط الأسلوب main)، ولكن هذا ليس هو الحال بالنسبة لمعظمها. لذلك، كل ما يحدث هنا يجري في مؤشر الترابط الخاص بالطريقة [main]:
- تقوم «المراقبة» بإصدار العنصر 1؛
- يتم تنفيذ الأسطر 15-18 وعرض هذا العنصر؛
- تُصدر المراقبة العنصر 2؛
- يتم تنفيذ الأسطر 15-18 وعرض هذا العنصر؛
- تُصدر المراقبة العنصر 3؛
- يتم تنفيذ الأسطر 15-18 وعرض هذا العنصر؛
- تُصدر المراقبة الإشعار [completed]؛
- يتم تنفيذ الأسطر 24-29؛
وهذا ما تظهره النتائج التي تم الحصول عليها:
تستند الفئة [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. مثال-03: فئة Observer
![]() |
تتوفر الطريقة [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);
}
});
};
}
في السطر 13، بدلاً من تمرير ثلاثة معلمات إلى الطريقة [subscribe]، يتم تمرير النوع [Observer] التالي إليها:
![]() |
النوع [Observer] هو واجهة تحتوي على ثلاث طرق:
- [onNext(T t)] التي يتم استدعاؤها في كل مرة يصدر فيها المراقب عنصرًا t؛
- [onError(Throwable th)] التي يتم استدعاؤها عندما تصدر المراقبة استثناءً th؛
- [onCompleted] التي يتم استدعاؤها عندما يشير المراقب إلى أنه انتهى من الإرسال؛
طريقة عمل الكود مشابهة لتلك التي تم شرحها سابقًا. ونحصل على النتائج التالية:
7.1.3. المثال-04: الطريقة [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] كائنًا قابلًا للمراقبة مُهيأًّا. لم يتم إصدار أي عناصر حتى الآن. عندما يشترك مشترك [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"));
}
}
- السطر 11: يتم إنشاء عنصر قابل للمراقبة يصدر أنواعًا من نوع Double؛
- الأسطر 11-21: يتم إنشاء مثيل لمعلمة الدالة [create] باستخدام فئة مجهولة تحتوي على الدالة الوحيدة [call] الواردة في الأسطر 12-20. المراقب الذي تم إنشاؤه في السطر 11 جاهز للإرسال، لكنه لن يرسل إلا عند وصول مراقب؛
- الأسطر 13-21: تتلقى الطريقة [call] مرجع المراقب؛
- الأسطر 14-17: إرسال 3 عناصر إلى المراقب؛
- السطر 19: إرسال إشعار بنهاية الإرسال إلى المراقب؛
- الأسطر 23-24: الاشتراك في المتغير القابل للمراقبة الوارد في السطر 11. يتم تنفيذ المعلمات الثلاثة [onNext, onError, onCompleted] الخاصة بالطريقة [subscribe] بواسطة ثلاث دال لامدا. سيؤدي هذا الاشتراك إلى إنشاء المشترك [Subscriber<Double>] الذي سيتم تمريره إلى الدالة [call] في السطر 13. عندئذ سيبدأ بث العناصر؛
- ويحدث كل ذلك في نفس الخيط: المراقب والمراقب؛
ونحصل على النتائج التالية:
تسمح الطريقة [Observable.create] بإنشاء عنصر قابل للمراقبة من أي ظاهرة. هذه هي الطريقة التي استخدمناها في الفقرة 2 من «الاكتشاف»، لتحويل واجهة متزامنة إلى واجهة غير متزامنة.
7.1.4. المثال-05: إعادة هيكلة [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()));
}
}
- السطر 56: النسخة الجديدة من الأسلوب الثابت [Observable.subscribe] تقبل كمعلمة النوع [Subscriber] الذي عرضناه في الفقرة السابقة؛
- الأسطر 37-52: المشترك (المشترك، المراقب). وهو ينفذ الواجهة Observer مع طرقها الثلاث onNext، onError، onCompleted؛
- الأسطر 61-64: من الآن فصاعدًا سنركز على الخيوط التي يتم فيها تنفيذ العنصر القابل للمراقبة ومراقبه؛
- السطر 62: اسم الخيط؛
- السطر 63: الوقت الحالي معبَّرًا عنه بالثواني والميلي ثانية. سيسمح لنا ذلك بمتابعة إرسال العناصر من قبل العنصر القابل للمراقبة ومعالجتها من قبل المراقب عبر الزمن؛
- هذا الكود له نفس وظيفة الكود السابق. لقد قمنا ببساطة بإعادة هيكلة هذا الأخير؛
النتائج التي تم الحصول عليها هي كما يلي:
- السطر 1 من النتائج: قبل السطر 56 من الكود، لم يحدث أي شيء بعد. تم فقط تكوين المتغير القابل للمراقبة؛
- السطر 2 من النتائج: يؤدي السطر 56 من الكود إلى استدعاء الطريقة [call] الموجودة في السطر 15. في السطر 3، يتم إرسال القيمة العددية 80.39 إلى المراقب؛
- السطر 4: يتلقى المراقب الرقم المرسل؛
- الأسطر 5-8: تتكرر العملية السابقة مرتين؛
- السطر 9: يرسل المراقب إشعارًا بانتهاء الإرسال؛
- السطر 10: يتلقى المراقب الإشعار؛
- السطر 11: يتم عرضه بواسطة السطر 57 من الكود؛
نرى إذن أن السطر 56 الخاص بالاشتراك وحده هو الذي تسبب في عرض الأسطر 2-10 من النتائج. عند البدء في استخدام المكتبة RxJava، نتساءل عن كيفية تسلسل الأحداث، ولا سيما الروابط التي تربط المراقب بالمراقب. ونلاحظ هنا أن السطر 56، أي الاشتراك في القابل للملاحظة،
- قد تسبب في إرسال جميع عناصر القابل للملاحظة؛
- وأن القابل للملاحظة والمراقب يتم تنفيذهما في نفس الخيط؛
- وبسبب ذلك، نلاحظ التسلسل التالي: إرسال العنصر i، مراقبة العنصر i، إرسال العنصر (i+1)، مراقبة العنصر (i+1)، ...
نتذكر أن المُرسِل كان ينتظر قبل إرسال عناصره:
// في انتظار
try {
Thread.sleep(500 - i * 100);
} catch (InterruptedException e) {
// خطأ
subscriber.onError(e);
}
حيث يمثل i في السطر 3 رقم الإرسال (0<=i<3). إذا لاحظنا أوقات إرسال عناصر القابل للملاحظة:
- السطران 2 و3: تم إرسال العنصر 0 بعد حوالي 500 مللي ثانية من بدء الاشتراك؛
- السطران 3 و5: تم إرسال العنصر 1 بعد حوالي 400 مللي ثانية من إرسال العنصر 0؛
- السطران 5 و7: تم إرسال العنصر 2 بعد حوالي 300 مللي ثانية من إرسال العنصر 1؛
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()));
}
}
- السطر 16: نقوم بإنشاء حاجز (إشارة) باستخدام كائن من النوع [CountDownLatch]. يستخدم هذا الكائن لمزامنة الخيوط فيما بينها. يتم تهيئته هنا بالقيمة 1 التي سنسميها قيمة الحاجز (أو الإشارة). ينتظر الخيط الحاجز من خلال عملية:
latch.await();
يتم حظر الخيط إذا كانت قيمة الحاجز >0. يمكن للخيط زيادة/تقليل القيمة الداخلية للحاجز. في السطر 48، يتم تقليل قيمة الحاجز بمقدار 1.
- السطر 63: يتم تكوين المراقب بحيث يتم تنفيذه على خيط يوفره المجدول [Schedulers.computation()]. يمكن لهذا المجدول توفير عدد من الخيوط يساوي عدد النوى الموجودة في جهاز التنفيذ. أظهرت الفقرة الخاصة بالتطبيق النموذجي استخدام مجدولات أخرى (انظر الفقرة 2.8)؛
ويتمثل مبدأ العمل في الكود فيما يلي:
- يتم تنفيذ الطريقة [main] في الخيط الرئيسي (main)؛
- السطر 66: يبدأ إرسال عناصر المراقب. سيتم إرسال هذه العناصر عبر خيط مختلف عن الخيط الرئيسي؛
- السطر 70: يتم حظر مؤشر الترابط الرئيسي لأن قيمة الحاجز تساوي 1 (انظر السطر 16). ولن يتمكن من المتابعة إلا عندما تتغير هذه القيمة إلى 0. ويحدث ذلك في السطر 48. ويقوم المراقب بخفض الحاجز عندما يتلقى إشعارًا بأن العنصر القابل للمراقبة قد أنهى عمليات الإرسال؛
يُنتج التنفيذ النتائج التالية:
- السطر 1: سيتم الاشتراك؛
- السطر 2: يؤدي هذا إلى تشغيل الأسلوب [call] على الخيط [RxComputationThreadPool-1]. لدينا الآن تنفيذ متوازي باستخدام خيطين؛
- السطر 3: لسبب غير واضح، قام الخيط [RxComputationThreadPool-1] بالتنازل عن السيطرة. ثم يتولى الخيط [main] زمام الأمور ويتم حظره بواسطة حاجز الحماية (السطر 70 من الكود). ومنذ تلك اللحظة، لا يمكن إلا للخيط [RxComputationThreadPool-1] أن يعمل؛
- الأسطر 4-11: نلاحظ السلوك الذي لوحظ سابقًا بين المتغير القابل للمراقبة ومراقبه، لكن كل شيء يحدث الآن في الخيط [RxComputationThreadPool-1]؛
- الأسطر 12-13: قام المراقب بإنزال الحاجز (السطر 48 من الكود) وانتهى عمل الخيط [RxComputationThreadPool-1]. يتولى الخيط [main] زمام الأمور ويعرض رسالتين؛
7.2.2. المثال-07: العنصر القابل للمراقبة والمراقب في خيطين مختلفين
![]() |
نقوم بتعديل المثال السابق على النحو التالي:
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()));
}
}
الكود مطابق لكود المثال السابق باستثناء السطر 63:
obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
الذي يهيئ المتغير القابل للمراقبة (subscribeOn) والمراقب (observeOn) ليتم تنفيذهما على أحد الخيوط التي يوفرها المجدول [Schedulers.computation()].
والنتائج التي تم الحصول عليها هي كما يلي:
يمكن ملاحظة النقاط التالية:
- يتم تنفيذ المراقب في مؤشر الترابط [RxComputationThreadPool-4] (الأسطر 3-4، 6، 8-9)؛
- يتم تنفيذ المراقب في مؤشر الترابط [RxComputationThreadPool-3] (الأسطر 5 و7 و10-11)؛
- أن كلاهما يعمل بشكل مستقل. وهكذا، في السطرين 8-9، يرسل المراقب 2 إشعارات (onNext، onCompleted) قبل أن يستلم المراقب الإشعار [onNext] (السطر 10)؛
تتولى المكتبة RxJava عملية نقل البيانات (الإرسالات) من مؤشر ترابط العنصر القابل للمراقبة إلى مؤشر ترابط المراقب. ولا داعي للمطور أن يقلق بشأن ذلك.
لقد رأينا كيفية إنشاء العناصر القابلة للمراقبة (Observable.from، Observable.create). سنستعرض الآن العناصر القابلة للمراقبة المُعرَّفة مسبقًا في المكتبة RxJava.
7.3. المراقبون المُعرَّفون مسبقًا
7.3.1. المثال-08: الطريقة [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;
}
}
- السطر 9: اسم العملية؛
- السطر 11: المتغير المرصود؛
- الأسطر 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> {
...
}
- السطر 11، الفئة Observateur<T> تمتد من الفئة Subscriber<T> التي قدمناها بإيجاز في الفقرة 7.1.3. وسنستخدمها كمعلمة للطريقة [Observable.subscribe]:
// التنفيذ القابل للمراقبة (المراقبة)
obs1.subscribe(observateur);
الطريقة [Observable.subscribe] المستخدمة في السطر 2 أعلاه لها التعريف التالي:
![]() |
تتمثل وظيفة [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 معه المعلومات التالية:
- السطر 14: حاجز أو إشارة ستُستخدم لحجب الخيط الرئيسي حتى يتلقى المراقب جميع العناصر التي أرسلها العنصر القابل للمراقبة. سيتم ذلك في السطر 36 من الكود عندما يتلقى المراقب من العنصر القابل للمراقبة إشعارًا بانتهاء الإرسال؛
- السطر 16: مثيل Consumer<String> الذي سيُستخدم لعرض رسالة على وحدة التحكم؛
- السطر 18: اسم المراقب لتمييزه عن غيره في حالة وجود عدة مراقبين؛
- السطر 20: اسم العملية المراقبة؛
- الأسطر 36 و46 و54: الطرق [onCompleted, onError, onNext] الخاصة بواجهة [Observer<T>] التي تنفذها الفئة المجردة [Subscriber<T>]. هذه الفئة لا تنفذ هذه الطرق. لذلك يجب تنفيذها في فئاتها الفرعية. قبل القيام بأي شيء في هذه الطرق، نتحقق مما إذا كان المراقب قد تم إلغاء اشتراكه في العنصر القابل للمراقبة الذي يراقبه؛
- السطر 59: تكتب الطريقة [onNext] الخاصة بالمراقب السلسلة jSON للعنصر المستلم. سيسمح لنا ذلك بعرض أنواع مختلفة من العناصر؛
وبناءً على ذلك، دعونا ندرس طريقة جديدة من الفئة Observable، وهي الطريقة [range]:
![]() |
تُصدر المتغيرات القابلة للملاحظة Observable.range(n,m) (m) أعدادًا صحيحة تتراوح من n إلى n+m-1. سندرسها باستخدام الكود [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()));
}
- السطر 16: سنستخدم مراقبين اثنين؛
- السطر 19: يتم تهيئة حارس الحاجز (السيمافور) على القيمة 2 لأننا سنضع كل مراقب على خيط مختلف. وبالتالي، سيتعين على الخيط الرئيسي انتظار انتهاء خيطي المراقبة؛
- السطر 22: نقوم بتكوين القابل للمراقبة بحيث يتم تنفيذه على مؤشر ترابط تابع للمجدول [Schedulers.computation()]. سيكون المراقب على نفس مؤشر الترابط الذي يعمل عليه القابل للمراقبة؛
- الأسطر 25-27: يتم اشتراك مراقبين اثنين في المتغير القابل للمراقبة. سيؤدي ذلك إلى تشغيله بالكامل لكل من المراقبين: سيتم إرسال الأعداد الصحيحة 15 و16 و17؛
- السطر 30: ينتظر الخيط الرئيسي انتهاء عمل المراقبين؛
النتائج التي تم الحصول عليها هي كما يلي:
- السطر 2: الخيط الرئيسي معطل في انتظار انتهاء عمل المراقبين الاثنين؛
- السطران 3-4: نلاحظ أن المراقب 0 موجود على الخيط [RxComputationThreadPool-1] والمراقب 1 على الخيط [RxComputationThreadPool-2]؛
- الأسطر 3-10: نلاحظ أن المراقبين يتلقون العناصر نفسها تمامًا؛
سنستخدم الفئة Observateur المُعرَّفة على هذا النحو لتوضيح سلوك أنواع أخرى من العناصر القابلة للمراقبة.
7.3.2. المثال-09: طرق 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()));
}
- السطر 22: تُصدر المراقبة أعدادًا صحيحة طويلة كل 500 مللي ثانية. تبدأ السلسلة بالرقم 0؛
- السطر 22: يُصدر هذا المراقب عددًا لا نهائيًا من القيم. تُنشئ الطريقة [Observable.take(n)] مراقبًا جديدًا لا يحتفظ إلا بأول n عنصر من العناصر الصادرة؛
![]() |
لنعد إلى كود المراقب:
Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
.doOnNext(l -> showInfos.accept(l.toString()));
السطر 2، تُنفَّذ الطريقة [Observable.doOnNext] في كل مرة يصدر فيها المراقب عنصرًا جديدًا. وغالبًا ما تُستخدم هذه الطريقة لتسجيل المعلومات. هنا، نريد تسجيل تاريخ إصدار العناصر للتحقق من صحة الفاصل الزمني البالغ 500 مللي ثانية. لا تُغيّر الطريقة [Observable.doOnNext] المراقب الذي تُطبق عليه. وتعريفها كما يلي:
![]() |
يؤدي التنفيذ إلى النتائج التالية:
- الأسطر 3 و7 و11: نلاحظ أن الفاصل الزمني للإرسال يقترب تقريبًا من 500 مللي ثانية؛
- المراقبان موجودان بالطبع في خيطين مختلفين، على الرغم من أن المراقب لم يتم تهيئته للتنفيذ باستخدام جدولة محددة. هذا هو الأداء الافتراضي للمراقب [Observable.interval] الذي نراه هنا؛
7.3.3. أمثلة-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()));
}
وقد استُخدم هذا الكود بالفعل في المثال السابق. ولم تتغير سوى السطور 21-22. لذا سنقوم بتجميع معظم هذا الكود في الفئة التالية [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()));
}
- السطر 13: تقبل هذه الطريقة معلمتين:
- nbObservateurs: عدد المراقبين للعمليات التي تم تمريرها كمعلمة ثانية؛
- processes: العمليات (المراقبة المسماة) المراد مراقبتها. بفضل الترميز [IProcess<?>]، ستتمكن العمليات من إرسال عناصر من أنواع مختلفة؛
- السطر 16: يجب أن يتحول لون الإشارة إلى الأخضر عندما ينتهي جميع المراقبين من جميع عمليات المراقبة الخاصة بهم. وبالتالي، فإن القيمة الأولية للإشارة تساوي عدد المراقبين مضروبًا في عدد عمليات المراقبة؛
- الأسطر 20-25: يتم اشتراك كل مراقب في جميع العمليات التي يجب مراقبتها؛
- السطر 23: يتم استرداد المتغير القابل للمراقبة من العملية (انظر الفقرة 7.3.1)؛
- السطر 23: يتم تسجيل مشاهد على هذه المتغيرات القابلة للمراقبة. ويتم تمرير 4 معلومات إلى هذا المشاهد:
- اسمه؛
- السيمافور الذي يجب عليه تقليله عندما يتلقى إشعار انتهاء إرسال المتغير القابل للمراقبة الذي يراقبه؛
- الطريقة التي يجب استخدامها عندما يرغب في تسجيل المعلومات على وحدة التحكم؛
- اسم العملية التي سيراقبها؛
بعد تعريف هذه الفئات، سيكون المثال 10 كما يلي:
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));
}
}
في السطر 11، يتم تعريف الطريقة الثابتة [Observable.error] على النحو التالي:
![]() |
وبالتالي، فإن السطر 8 يقوم بتكوين عنصر قابل للمراقبة يكتفي بإصدار استثناء موجه إلى الطريقة [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]
في السطرين 3 و4، تلقت الطريقة [onError] لكل من المشتركين الاستثناء الذي أطلقته المتغيرات القابلة للمراقبة.
يتميز هذا التنفيذ بخاصية فريدة: لم يتم استدعاء الطرق [onCompleted] الخاصة بكل من المراقبين. ونتيجة لذلك، لم يتم خفض الحاجز وظل الخيط الرئيسي معطلاً في الطريقة الثابتة [ProcessUtils.subscribe] في السطر 3 التالي:
// في انتظار
showInfos.accept("main : attente fin observation");
latch.await();
// النهاية
showInfos.accept("main : fin observation");
نكتشف هنا أنه في حالة حدوث خطأ في المراقب، لا يتم استدعاء الدالة [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();
}
نضيف السطرين 7 و8 لإزالة الحاجز في حالة حدوث خطأ في المتغير القابل للمراقبة. باستخدام هذا الكود الجديد، يعطي التنفيذ النتائج التالية:
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]
نحصل على السطر 5 الذي لم نحصل عليه سابقًا.
سيكون المثال 11 كما يلي:
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));
}
}
في السطر 10، تُنشئ الطريقة الثابتة [Observable.empty] متغيرًا قابلًا للمراقبة لا يصدر أي عنصر. إنه يصدر فقط إشعار انتهاء الإرسال؛
![]() |
يؤدي تنفيذ كود المثال أعلاه إلى النتائج التالية:
- السطران 2 و3: نلاحظ أن كلا المراقبين يتلقيا إشعار انتهاء الإرسال دون أن يكونا قد تلقيا أي عناصر من قبل.
قد يتساءل المرء عن الفائدة من هذه الطريقة. يمكن استخدامها بطريقة مشابهة لمجموعة، تكون فارغة في البداية، ثم يتم تجميع العناصر فيها لاحقًا:
في السطر 3، يتم دمج المراقب الأولي obs (السطر 1) مع مراقبين آخرين.
يوضح المثال 12 الطريقة الثابتة [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] متغيرًا قابلًا للمراقبة لا يصدر أبدًا:
![]() |
يؤدي تنفيذ المثال إلى النتائج التالية:
في السطر 2، ينتظر الخيط الرئيسي إلى أجل غير مسمى. في الواقع، لا يصدر أي عنصر قابل للمراقبة الإشعار [onCompleted] الذي يسمح بتحويل إشارة المرور (حاجز المرور) إلى اللون الأخضر (خفض الحاجز).
7.4. Multi-threading
7.4.1. المثال 13: مؤشر ترابط الإجراء، مؤشر ترابط المراقبة
في الفقرة 7.1.3، أنشأنا عنصرًا قابلًا للمراقبة باستخدام الطريقة الثابتة [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();
}
- السطر 5: تتمتع الواجهة [IProcessAction<T>] بجميع خصائص الواجهة [Observable.OnSubscribe<T>]؛
- السطر 8: كما تحتوي على طريقة [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;
}
}
- السطر 8: الفئة [ProcessAction01<T>] تنفذ الواجهة [IProcessAction<T>] وبالتالي الواجهة [Observable.OnSubscribe<T>]؛
- السطر 11: اسم الإجراء؛
- السطر 12: عدد القيم المطلوب إصدارها؛
- السطر 13: مثيل من النوع [Func1<Integer, T>] الذي يقوم، انطلاقًا من عدد صحيح، بإنشاء نوع T الذي سيتم إصداره بواسطة المراقب (السطران 35 و37)؛
- الأسطر 16-20: يتم تمرير اسم الإجراء وعدد القيم المراد إصدارها ودالة الإصدار إلى المنشئ؛
- الأسطر 23-42: كود العملية؛
- السطر 23: تتلقى الطريقة [call] كمعلمة المشترك في المراقب المرتبط بالعملية؛
- السطر 28: تقوم العملية بإصدار عناصرها بعد فترة انتظار ذات مدة عشوائية؛
- السطر 32: إصدار خطأ؛
- السطر 37: إرسال عادي؛
- السطر 41: إرسال إشعار بنهاية الإرسال؛
- الأسطر 25-38: تقوم العملية بإرسال قيم حقيقية لـ nbValues بعد فترة انتظار عشوائية (السطر 30)؛
- السطر 35: يتم توفير القيمة المراد إرسالها بواسطة الدالة [func1] التي يتم تمريرها كمعلمة إلى المنشئ (السطر 16)؛
نقوم بإعادة هيكلة الفئة [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، يقبل المنشئ 3 معلمات:
- الإجراء المسمى الذي سيُستخدم لإنشاء العنصر القابل للمراقبة (السطر 5)؛
- جدولة العملية المراقبة (قد تكون null)؛
- جدول المراقب (قد يكون null)؛
- السطر 5: يتم إنشاء المتغير القابل للمراقبة بناءً على الإجراء الذي تم تمريره كمعلمة؛
يرصد الكود التالي [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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<String> process2 = new Process<>(
new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
// العملية 3
Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
Schedulers.computation());
// العملية 4
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);
}
}
- الأسطر 13-15: تنتج العملية process1 عددًا حقيقيًا واحدًا على خيط حسابي سيتم مراقبته على خيط حسابي آخر؛
- الأسطر 17-18: تنتج العملية process2 سلسلتين من الأحرف على خيط حسابي، ولا توجد إشارة إلى خيط المراقب. وتُظهر النتائج أن المراقبة تتم افتراضيًا على نفس الخيط الذي تُنفَّذ عليه العملية؛
- السطران 20-21: تنتج العملية process3 3 أعداد صحيحة على خيط غير محدد، وسيتم مراقبة هذه القيم على خيط حسابي. تظهر النتائج أن تنفيذ العملية يتم افتراضيًا على الخيط الرئيسي؛
- السطر 23: تنتج العملية process4 4 قيم منطقية على خيط غير محدد، وسيتم مراقبة هذه القيم على خيط غير محدد. تظهر النتائج أن تنفيذ العملية ومراقبتها يتمان افتراضيًا على الخيط الرئيسي؛
نتيجة تنفيذ هذا الكود هي كما يلي:
- تنتج العملية process1 عددًا حقيقيًا واحدًا (السطر 4) على مؤشر الترابط الحسابي [RxComputationThreadPool-4] الذي يتم رصده على مؤشر الترابط الحسابي [RxComputationThreadPool-3] (السطر 6)؛
- تنتج العملية process2 سلسلتين من الأحرف (السطران 12 و14) على مؤشر الترابط الحسابي [RxComputationThreadPool-5]، ويتم رصدهما على نفس مؤشر الترابط (السطران 13 و15)؛
- تنتج العملية process3 3 أعداد صحيحة (السطور 21، 23، 25) على الخيط الرئيسي، والتي يتم رصدها على خيط الحساب [RxComputationThreadPool-6] (السطور 22، 24، 28)؛
- تُنتج العملية process4 4 قيم منطقية (الأسطر 34 و36 و38 و40) على الخيط الرئيسي، والتي يتم رصدها على هذا الخيط الرئيسي نفسه (الأسطر 33 و35 و37 و39)؛
يُرجى من القارئ متابعة ما يلي:
- دورة حياة العملية المراقبة وخيطها؛
- دورة حياة المراقب الخاص بها وخيطه؛
يكمن جزء كبير من أهمية مكتبات Rx في هذا التعدد في الخيوط الذي لا يتعين على المطور إدارته بنفسه.
7.5. تركيبات لعدة عناصر قابلة للمراقبة
7.5.1. المثال 14: دمج عنصرين قابلين للمراقبة باستخدام [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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- الأسطر 15-17: ستقوم عملية تُسمى [process1] بإصدار 3 أعداد حقيقية على خيط حسابي. كما سيتم ملاحظتها على خيط حسابي؛
- السطور 19-20: ستقوم عملية تُسمى [process2] بإصدار سلسلتين من الأحرف على خيط حسابي. خيط المراقبة غير محدد. وقد رأينا سابقًا أنه في هذه الحالة، يكون خيط المراقبة هو خيط الحساب؛
- السطر 23: يتم دمج العمليتين، أي يتم إنشاء متغير قابل للمراقبة تأتي عناصره في وقت واحد من كلتا العمليتين. ويُستخدم لهذا الغرض الأسلوب الثابت [Observable.merge]:
![]() |
على عكس ما قد يوحي به الرسم التخطيطي أعلاه، عند الدمج، يمكن أن تتخلل عناصر التدفق 1 عناصر التدفق 2. وهذا ما تظهره نتائج التنفيذ:
- السطر 3: يتم تنفيذ العملية [process1] على مؤشر ترابط الحساب [RxComputationThreadPool-4]؛
- السطر 4: يتم تنفيذ العملية [process2] على مؤشر ترابط الحساب [RxComputationThreadPool-5]؛
- السطر 9: يتم مراقبة العملية [process12] على مؤشر الترابط الحسابي [RxComputationThreadPool-3]. لا أعرف القاعدة التي أدت إلى هذا الاختيار؛
- الأسطر 9-11: نرى أن المراقب يراقب عناصر من كل من العمليتين [process1] (السطر 5) و[process2] (السطران 6 و7) في حين أن أياً منهما لم ينتهِ بعد (هناك اختلاط)؛
- تنتهي العملية [process12] (السطر 17) عند انتهاء العمليتين process1 و process2؛
7.5.2. المثال 15: ربط متغيرين قابلين للمراقبة باستخدام [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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- الأسطر 15-17: ستقوم عملية تُسمى [process1] بإصدار 3 أعداد حقيقية على خيط حسابي. كما سيتم ملاحظتها على خيط حسابي؛
- السطور 19-20: ستقوم عملية تُسمى [process2] بإصدار سلسلتين من الأحرف على خيط غير مفروض، وهو هنا الخيط الرئيسي الافتراضي. وسيتم ملاحظتها على خيط حسابي؛
- السطر 23: يتم ربط العمليتين معًا، أي يتم إنشاء متغير قابل للمراقبة تتأتى عناصره من العمليتين. ولا يحدث خلط بين القيم المُصدرة. ستقوم العملية [process12] أولاً بإصدار جميع قيم العملية [process1] ثم قيم العملية [process2]. ويُستخدم لهذا الغرض الأسلوب الثابت [Observable.concat]:
![]() |
نتائج التنفيذ هي كما يلي:
- الأسطر 3-10: يتم تنفيذ العملية [process1] وتقوم العملية [process12] بإصدار القيم التي أصدرتها العملية [process1]؛
- السطر 9: انتهت العملية [process1]؛
- الأسطر 11-17: يتم تنفيذ العملية [process2] وتقوم العملية [process12] بإخراج القيم التي أخرجتها العملية [process2]؛
هناك أمر غريب يتعلق بالعملية process2: لم يتم تحديد خيط تنفيذ لها. لذا كان من المتوقع أن يكون خيط التنفيذ الافتراضي هو الخيط الرئيسي. لكن الأمر لم يكن كذلك. كان مؤشر الترابط التنفيذي هو مؤشر الترابط الحسابي [RxComputationThreadPool-3] (السطر 11). لذا، عندما لا يتم تحديد مؤشر ترابط تنفيذي أو مؤشر ترابط مراقبة، لا يمكن التكهن بمؤشر الترابط الذي سيتم اختياره.
7.5.3. المثال 16: دمج متغيرين قابلين للمراقبة باستخدام [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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- الأسطر 16-18: ستقوم عملية تُسمى [process1] بإخراج 3 أعداد حقيقية على خيط حسابي. كما سيتم ملاحظتها على خيط حسابي؛
- الأسطر 20-21: ستقوم عملية تُسمى [process2] بإخراج سلسلتين من الأحرف على خيط غير مفروض. كما أن خيط المراقبة غير مفروض أيضًا؛
- الأسطر 23-32: إنشاء مثيل لنوع [FuncN<String>] باستخدام فئة مجهولة. FuncN هي واجهة وظيفية:
![]() |
تتوقع الطريقة [FuncN.call] مصفوفة من الكائنات وتُرجع نوعًا R. ستُستخدم الدالة [funcn] لدمج العمليتين process1 و process2 بهذا الترتيب. في الدالة [FuncN.call]:
- ستكون args[0] هي Double؛
- سيكون args[1] هو String؛
هنا، ستكون نتيجة [funcn.call] هي سلسلة الأحرف الموجودة في السطر 27. ولا يتطلب تكوين هذه النتيجة معرفة أنواع معلمات الدالة 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]؛
يُنتج التنفيذ النتائج التالية:
- السطران 7 و11: تُصدر العملية process12 عنصرين؛
- السطر 8: العنصر الإضافي الذي أصدرته العملية process1، والتي ليس لها شريك في العملية process2، لم تصدره العملية الناتجة process12؛
ونلاحظ أن العملية process2، التي لم يُفرض عليها أي خيط تنفيذ أو خيط مراقبة، استخدمت الخيط الرئيسي لكليهما.
7.5.4. المثال 17: دمج متغيرين قابلين للمراقبة باستخدام [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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- الأسطر 14-16: ستقوم عملية تُسمى [process1] بإصدار 3 أعداد حقيقية على خيط حسابي. كما سيتم مراقبة هذه العملية على خيط حسابي؛
- الأسطر 18-20: ستقوم عملية تُسمى [process2] بإصدار عددين حقيقيين على خيط غير مفروض. وسيتم رصدهما على خيط حسابي؛
- السطر 23: يتم دمج المراقبتين مع الطريقة الثابتة التالية [Observable.combineLatest]:
![]() |
تعمل المتغيرات القابلة للمراقبة [combineLatest] بالطريقة التالية: عندما تصدر إحدى المتغيرتين القابلتين للمراقبة عنصرًا E1، يتم دمج هذا العنصر بواسطة [combineFunction] مع آخر عنصر أصدرته المتغير الأخرى القابلة للمراقبة.
يؤدي تنفيذ هذا الكود إلى النتيجة التالية:
- السطر 5: يتم دمج الإرسال من process2 (56) مع العنصر الأخير الذي أرسله process1 (54، السطر 4) وينتج عنه النتيجة الواردة في السطر 7؛
- السطر 6: يتم دمج الإرسال من process1 (51.6) مع العنصر الأخير الذي تم إرساله بواسطة process2 (56، السطر 5) وينتج عنه النتيجة الواردة في السطر 8؛
- السطر 9: يتم دمج الإرسال الصادر عن process2 (261.8) مع آخر عنصر أرسله process1 (51.6، السطر 6) وينتج عنه النتيجة الواردة في السطر 12؛
- السطر 13: يتم دمج الإرسال الصادر عن process1 (80.39) مع العنصر الأخير الذي أرسله process2 (261.8، السطر 9) وينتج عنه النتيجة الواردة في السطر 15؛
نحن هنا أمام صيغة مختلفة من المراقب [zip]، حيث لا تكون العناصر المدمجة هذه المرة بالضرورة هي العناصر التي تشغل نفس الموضع في التدفقات. ونلاحظ هنا أن العملية process2، التي لم يُفرض عليها خيط تنفيذ، قد نُفِّذت هنا على الخيط الرئيسي (السطر 2).
7.5.5. المثال 18: دمج عنصرين قابلين للمراقبة باستخدام [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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- الأسطر 14-16: ستصدر عملية تُسمى [process1] 3 أعداد حقيقية على خيط حسابي. كما سيتم ملاحظتها على خيط حسابي؛
- الأسطر 18-20: ستقوم عملية تُسمى [process2] بإصدار عددين حقيقيين على خيط غير مفروض. وسيتم رصدهما على خيط غير مفروض؛
- السطر 22: يتم دمج المراقبين مع الطريقة الثابتة التالية [Observable.amb]:
![]() |
كما يوضح الرسم التخطيطي أعلاه، فإن المتغير القابل للمراقبة [Observable.amb(Observable o1, Observable o2)] يُصدر عناصر المتغير القابل للمراقبة الذي يُصدر أولاً. وهذا ما تؤكده نتائج المثال المقدم:
- السطر 4، العملية process2 هي التي ترسل أولاً؛
- السطران 8 و12: العملية process12 تُرسل جميع العناصر التي أرسلتها العملية process2 (السطران 4 و11)؛
7.6. سلسلة معالجة عنصر قابل للمراقبة
7.6.1. المثال 19: تحويل متغير قابل للمراقبة باستخدام [Observable.map]
في الأمثلة السابقة، استعرضنا تركيبات متنوعة لدمج متغيرين قابلين للمراقبة في متغير قابل للمراقبة ثالث. نقدم الآن طرقًا ثابتة من الفئة [Observable] التي تتيح إجراء عمليات التحويل والتصفية والتجميع على متغير قابل للمراقبة. سنجد هنا طرقًا مشابهة لتلك الموجودة في الفئة [Stream] التي تمت دراستها في الفقرة 5.
سيكون مثالنا الأول كما يلي:
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 {
// العملية 1
Process<Double> process1 = new Process<>(
new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<String> process2 = new Process<>("process2",
process1.getObservable().map(d -> String.format("valeur-%s", d)));
// الاشتراكات
ProcessUtils.subscribe(1, process2);
}
}
- الأسطر 14-16: ستقوم عملية تُسمى process1 بإصدار 3 أعداد حقيقية على خيط حسابي. كما سيتم مراقبة هذه العملية على خيط حسابي؛
- السطور 17-18: سيتم تحويل الأرقام التي أصدرتها العملية process1 إلى سلاسل أحرف في العملية process2؛
- السطر 20: يتم مراقبة process2؛
الطريقة [Observable.map] في السطر 18 مشابهة للطريقة [Stream.map] التي تمت دراستها في الفقرة 5.5:
![]() |
نتائج المثال هي كما يلي:
- الأسطر 4 و5 و8: قيم process1. وهي أعداد حقيقية؛
- السطور 6 و7 و10: انبعاثات 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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
// الاشتراكات
ProcessUtils.subscribe(1, process2);
}
}
- السطران 11-12: ستقوم عملية تُسمى process1 بإصدار الأعداد الصحيحة من 0 إلى 2 على خيط حسابي. كما سيتم رصدها على خيط حسابي؛
- السطر 14: سيتم تصفية الأرقام التي تصدرها العملية process1 بحيث لا يتم الاحتفاظ في العملية process2 إلا بالأرقام الزوجية؛
- السطر 20: يتم مراقبة process2؛
الطريقة [Observable.filter] في السطر 18 مشابهة للطريقة [Stream.filter] التي تمت دراستها في الفقرة 5.4:
![]() |
نتائج المثال هي كما يلي:
- السطور 4 و5 و7: البث الصادر عن process1؛
- السطران 6 و9: الانبعاثات الملاحظة لـ process2. العناصر الزوجية هي تلك الخاصة بـ process1؛
7.6.3. المثال 21: تحويل متغير قابل للرصد باستخدام [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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- السطران 12-13: ستقوم عملية تُسمى process1 بإصدار الأعداد الصحيحة من 0 إلى 2 على خيط حسابي. كما سيتم ملاحظتها على خيط حسابي؛
- السطور 15-18: يتم تحويل كل عدد n يصدره process1 إلى عنصر قابل للمراقبة يصدر الأرقام الثلاثة (10*n، 10*n+1، 10*n+2). لو استخدمنا في السطر 15 الطريقة [map]، لكان process2 يُصدر نوعًا Observable<Integer> وليس نوعًا Integer. تسمح الطريقة المستخدمة [flatMap] بتسوية (flatten) سلسلة العناصر من النوع Observable<Integer> إلى سلسلة عناصر من النوع Integer تتكون من كل عنصر من عناصر كل نوع من أنواع Observable<Integer>؛
- السطر 20: نلاحظ وجود process2؛
الطريقة [Observable.flatMap] الواردة في السطر 15 مشابهة للطريقة [Stream.flatMap] التي تمت دراستها في الفقرة 5.6.12:
![]() |
نتائج المثال هي كما يلي:
- الأسطر 5-7: الإرسالات الثلاثة لـ process2 عقب إرسال السطر 4 من process1؛
- الأسطر 9-11: الإرسالات الثلاثة لـ process2 عقب إرسال السطر 8 من 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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
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);
}
}
- السطر 14: نستخدم الطريقة [Observable.map]؛
- السطر 16: الذي يُرجع نوعًا Integer[]؛
النتائج هي كما يلي:
- الأسطر 6 و7 و10: نرى نتائج 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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
int value = i * 10;
return Observable.just(value, value + 1, value + 2);
}).filter(i -> i % 2 == 0));
// الاشتراكات
ProcessUtils.subscribe(1, process2);
}
}
- الأسطر 15-18: يتبع flatMap filter؛
نتائج التنفيذ هي كما يلي:
- الأسطر 8-13: لم يصدر 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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMapIterable(i -> {
int value = i * 10;
return Arrays.asList(value, value + 1, value + 2);
}).filter(i -> i % 2 == 0));
// الاشتراكات
ProcessUtils.subscribe(1, process2);
}
}
في السطر 16، بدلاً من استخدام الطريقة [flatMap]، تُستخدم الطريقة [flatMapIterable]. في هذه الحالة، يجب أن تنتج دالة التحويل نوعًا Iterable<T> (السطر 18) بدلاً من نوع Observable<T>.
ونحصل على نفس النتائج التي حصلنا عليها سابقًا.
لنعد إلى تعريف الطريقة [flatMap]:
![]() |
نلاحظ أعلاه أن عنصرًا أزرق [3] قد أُدرج بين العنصرين الأخضرين [1-2]. وهذا يعني أن الطريقة [flatMap]، في عملية تسطيحها لعناصر Observable<T>، تحترم ترتيب إرسال هذه الملاحظات الداخلية المختلفة. ويتضح ذلك من خلال المثال التالي [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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
Schedulers.computation(), Schedulers.computation());
// العملية 3
Process<Integer> process3 = new Process<>("process3",
process1.getObservable().flatMap(i -> process2.getObservable()));
// الاشتراكات
ProcessUtils.subscribe(1, process3);
}
}
- السطران 11-12: تقوم العملية process1 بإصدار الأعداد الصحيحة [0,1]؛
- السطران 14-15: تُنتج العملية process2 الأعداد الصحيحة [10,11,12]؛
- السطران 17-18: يُربط كل عنصر يصدره process1 بالمُراقب الخاص بالعملية process2. وهذا يعني أن:
- سيتم ربط العنصر [0] في process1 بمتغير قابل للمراقبة يصدر [10,11,12]؛
- وينطبق الأمر نفسه على العنصر 1؛
في النهاية، سيتم إرسال الأرقام الستة [10, 11, 12, 10, 11, 12]. نريد معرفة الترتيب الذي ستصدر به.
نتائج التنفيذ هي كما يلي:
نلاحظ أن ترتيب الإرسال لعملية process3 كان: [10, 10, 11, 12, 11, 12] (الأسطر 11، 12، 14، 17، 19، 22). وبالتالي، فقد حدث بالفعل خلط بين العناصر التي أصدرتها العملية 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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
Schedulers.computation(), Schedulers.computation());
// العملية 3
Process<Integer> process3 = new Process<>("process3",
process1.getObservable().concatMap(i -> process2.getObservable()));
// الاشتراكات
ProcessUtils.subscribe(1, process3);
}
}
في السطر 18، تم استبدال [flatMap] بـ [concatMap]. وكانت نتائج التنفيذ كما يلي:
نلاحظ أن ترتيب إصدار العملية process3 كان: [10, 11, 12, 10, 11, 12] (الأسطر 12-14، 17، 19، 22). لم يتم خلط العناصر التي أصدرتها العملية process2.
هناك نسخة أخرى من الطريقة [map] وهي الطريقة [switchMap]:
![]() |
في المثال أعلاه، تنشأ من العنصر القابل للملاحظة [1] ثلاثة عناصر قابلة للملاحظة أخرى هي [2] المكونة من عنصرين، والتي يتم بعد ذلك تسطيحها كما في [flatMap] و[3]. يمكن ملاحظة أن النتيجة تحتوي على 5 عناصر وليس 6. ويرجع ذلك إلى أنه قبل أن يصدر المراقب الثاني عنصره رقم 2 [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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
Schedulers.computation(), Schedulers.computation());
// العملية 3
Process<Integer> process3 = new Process<>("process3",
process1.getObservable().switchMap(i -> process2.getObservable()));
// الاشتراكات
ProcessUtils.subscribe(1, process3);
}
}
يؤدي تنفيذ المثال إلى النتائج التالية:
- يُصدر process1 عنصرين ينتج عنهما عنصران قابلان للمراقبة process2 مكونان من 3 عناصر؛
- السطر 14: يتلقى المراقب العنصر رقم 0 الذي أرسله العنصر القابل للمراقبة الأول process2 في السطر 6؛
- السطر 15: يتلقى المراقب العنصر رقم 0 الذي أرسله المراقب الثاني process2 في السطر 13. لا توضح القصة سبب عدم تلقيه من قبل العنصرين 1 و2 اللذين أرسلهما العنصر الأول القابل للملاحظة process2 في السطرين 7 و8. وعلى أي حال، يتم التخلي عن العنصر الأول القابل للملاحظة process2؛
- وفي النهاية، لا يرى المراقب سوى 4 عناصر (السطور 14 و15 و17 و20) بدلاً من الـ6 التي تم إرسالها؛
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);
}
}
النتائج
[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);
}
}
النتائج
[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);
}
}
النتائج
[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);
}
}
- السطر 10: يحسب مجموع عناصر المتغير القابل للرصد. والنتيجة هي متغير قابل للرصد يُصدر هذا المجموع؛
النتائج
[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);
}
}
- السطر 10: يُرجع قيمة Observable<Boolean> التي تُصدر العنصر true، إذا كان شرط الدالة [all] صحيحًا بالنسبة لجميع العناصر، وإلا تُرجع false؛
النتائج
[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);
}
}
- السطر 10: تُنشئ [Observable.count] متغيرًا قابلًا للمراقبة مكونًا من عنصر واحد يمثل مجموع العناصر المراقبة؛
النتائج
[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);
}
}
النتائج
[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);
}
}
- السطر 11: تجمع الطريقة [groupBy] العناصر العشرة الصادرة في مجموعتين، الأعداد الزوجية والأعداد الفردية. والنتيجة هي نوع Observable<GroupedObservable<Boolean, Integer>>، أي كائن قابل للمراقبة (observable) عناصره من النوع GroupedObservable<Boolean, Integer>، حيث Boolean هو نوع مفتاح المجموعة (false، true هنا) وهو أيضًا نوع نتيجة دالة لامدا التي تم تمريرها كمعلمة إلى الطريقة [groupBy]، و Integer هو نوع عناصر المجموعة؛
- السطر 12: النوع GroupedObservable يحتوي على دالة [asObservable] التي تسمح بإنشاء متغير قابل للمراقبة من هذا النوع. وبالتالي سيكون لدينا نوعان هما Observable<Integer>، أحدهما للأعداد الزوجية والآخر للأعداد الفردية. ومن هذين الملاحظين، ستقوم الطريقة [concatMap] بإنشاء ملاحظ واحد فقط؛
النتائج
[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 {
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
Schedulers.computation());
// العملية 2
Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
// الاشتراكات
ProcessUtils.subscribe(1, process2);
}
}
- في السطر 15، تربط الطريقة [timestamp] ساعة بكل عنصر من العناصر المعالجة؛
النتائج
في هذا المثال، من الصعب تحديد ما تمثله المعلومات timestamp:
- السطران 4-5: نلاحظ أن العنصر 1 من process1 قد تم إرساله بعد 139 مللي ثانية من العنصر 0؛
- السطران 6 و7: نلاحظ أن العنصر 1 من process2 قد رُصد بعد 234 مللي ثانية من العنصر 0؛
- السطران 5 و8: نلاحظ أن العنصر 2 من process1 قد أُرسل بعد 33 مللي ثانية من العنصر 1؛
- السطران 7 و10: نلاحظ أن العنصر 2 من process2 قد تمت ملاحظته بعد 37 مللي ثانية من العنصر 1؛
تعود هذه الفروق الزمنية إلى أن خيوط المراقبة والتنفيذ الخاصة بالملاحظات ليست هي نفسها. إذا استبدلنا السطرين 12 و13 بالسطور التالية (المثال 22j):
// العملية 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
- السطران 2-3: لا يتم فرض مؤشر ترابط المراقبة. ومن المعروف أنه في هذه الحالة، يتم مراقبة المتغير القابل للمراقبة في المكان الذي يتم فيه تنفيذه؛
وهذا يعطي النتائج التالية:
- السطران 4 و6: ترسل العملية process1 العنصر رقم 1 بعد 1587 مللي ثانية من العنصر رقم 0؛
- السطران 5 و7: يراقب المراقب هذين العنصرين بفارق زمني قدره 586 مللي ثانية؛
- السطران 6 و8: تُرسل العملية process1 العنصر رقم 2 بعد 396 مللي ثانية من إرسال العنصر رقم 1؛
- السطران 7 و9: يلاحظ المراقب هذين العنصرين بفارق زمني قدره 396 مللي ثانية؛
هنا، تكون قيم عملية timestamp متسقة: فهي تمثل بالفعل تاريخ إرسال العنصر.
7.7. المجدولات
7.7.1. المثال 23: برنامج الجدولة [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 عمليات تُنفَّذ على خيط حسابي؛
- السطر 17: تُصدر كل عملية عددًا حقيقيًّا عشوائيًّا؛
- السطر 21: يتم الاشتراك في جميع هذه العمليات؛
النتائج هي كما يلي:
- الأسطر 2-10: تبدأ العمليات الثماني الأولى على 8 خيوط مختلفة (تحتوي الآلة المستخدمة على 8 نوى). يمكن ملاحظة أنها تبدأ جميعها في نفس اللحظة تقريبًا؛
- الأسطر 17-19: تنتهي 3 عمليات، وبذلك تُحرر 3 خيوط؛
- السطران 23-24: يمكن عندئذٍ للعمليتين الأخيرتين أن تبدآ باستخدام خيطين من الخيوط التي تم تحريرها؛
وبالتالي، يمكننا استنتاج أن المجدول [Schedulers.computation] يوفر مجموعة من n خيوط، حيث n هو عدد النوى في الجهاز. ويتم تنفيذ الخيوط بشكل متوازٍ على هذه النوى.
7.7.2. المثال 24: المجدول [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);
}
}
- السطر 18: تُنفَّذ العمليات باستخدام خيوط المجدول [Schedulers.io]؛
وهذا يعطي النتائج التالية:
- الأسطر 2-10: تبدأ كل عملية من العمليات العشر على خيط مختلف. وعلى عكس الحالة السابقة، تم تشغيل جميع العمليات بنجاح. ونلاحظ أن عمليات التشغيل هذه تستغرق 6 مللي ثانية، في حين كانت تستغرق 1 مللي ثانية في السابق؛
- الأسطر 13-18: تُصدر القيم القابلة للمراقبة واحدة تلو الأخرى وليس بشكل شبه متوازي كما كان الحال سابقًا؛
ما الفرق بين جدولي التشغيل [Schedulers.io] و [Schedulers.computation]؟ يمكن العثور على إجابة في URL و [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:
![]() |
7.7.3. المثال-25: برنامج الجدولة [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]:
في URL و [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io]، نوضح أن المجدول [Schedulers.io] يوفر مجموعة من الخيوط، وهو ما لا يفعله المجدول [Schedulers.newThread]. يقوم تجمع الخيوط تلقائيًا بإنشاء عدد n من الخيوط. ويقوم بتخصيصها للعمليات التي تحتاج إليها. وعندما تنتهي هذه العمليات، لا يتم حذف خيوطها بل تعود إلى المجموعة ويمكن إعادة استخدامها بعد ذلك من قبل عملية أخرى. وهذا أكثر اقتصادية من إنشاء/حذف الخيوط باستمرار. لذا يمكننا أن نعتقد أنه من الأفضل استخدام المجدول [Schedulers.io].
7.7.4. المثال 26: المجدولان [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()));
}
- السطر 17: جدولة. ستكون إما [Schedulers.immediate] كما هو الحال هنا أو [Schedulers.trampoline] لاحقًا؛
- السطر 19: يمكن تنفيذ إجراءات من النوع Action0 (السطران 21 و20) على عمال المجدول. تسمح الطريقة [Scheduler.createWorker] بإنشاء عامل. تسمح الطريقة [Worker.schedule(Action0)] بتنفيذ نوع Action0 بواسطة عامل؛
- الأسطر 21-27: إجراء أول يُسمى [action02] سيتم تنفيذه (السطر 40) بواسطة عامل المعالجة المذكور في السطر 19؛
- الأسطر 30-38: إجراء ثانٍ يُسمى [action01]. وتتميز هذه العملية بأنها تُنفذ الإجراء action02 على نفس العامل الذي تنفذ عليه هي نفسها (السطر 34). وهنا يكمن الفرق بين [Schedulers.immediate] و [Schedulers.trampoline]:
- إذا كان المجدول هو [Schedulers.immediate]، ففي السطر 34، سيتم تنفيذ الإجراء action02 على الفور (ومن هنا جاء اسم المجدول)، وسيتم إيقاف الإجراء action01 الجاري. وعندها ستظهر الرسالة الواردة في السطر 25. وبمجرد انتهاء الإجراء action02، سيستأنف الإجراء action01 وستظهر الرسالة الواردة في السطر 36؛
- إذا كان المجدول هو [Schedulers.trampoline]، ففي السطر 34، يتم تعليق العملية action02. ولن يتم تنفيذها إلا بعد انتهاء المهمة الجارية action01. عندها ستظهر الرسالة الموجودة في السطر 36. وبمجرد انتهاء الإجراء action01، سيتم تنفيذ الإجراء action02 وسنرى الرسالة الموجودة في السطر 25؛
يؤدي تنفيذ الكود أعلاه إلى النتائج التالية:
إذا استخدمنا في السطر 17 المجدول [Schedulers.trampoline]، فسنحصل على النتائج المعاكسة:
ومع ذلك، من الصعب ربط ذلك بالمراقبات. لم أجد مثالًا مقنعًا يوضح فائدة تنفيذ مراقبة على أحد هذين الخيطين. ومع ذلك، إليكم مثالًا، لكنني لا أجده طبيعيًا على الإطلاق:
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 = Schedulers.immediate().createWorker();
// العامل worker = Schedulers.trampoline().createWorker();
// المراقب 1 على العامل
worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {
@Override
public void call(Integer i) {
ProcessUtils.showInfos.accept(String.valueOf(i));
// المتغير القابل للمراقبة 2 على نفس العامل
worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
@Override
public void call(Integer i) {
ProcessUtils.showInfos.accept(String.valueOf(i));
}
}));
}
}));
}
}
- السطران 13-14: يتم إنشاء عامل (worker) من أحد المجدولين [Schedulers.immediate] و [Schedulers.trampoline]؛
- السطر 16: يتم جدولة عنصر قابل للمراقبة الأول obs1 على هذا العامل لإصدار الأرقام [1,2]
- السطر 22: في كل مرة يتم فيها رصد عنصر من هذا المتغير القابل للمراقبة obs1، يتم تشغيل عملية رصد متغير قابل للمراقبة ثانٍ obs2 على نفس العامل لإصدار الأرقام [100,101]؛
باستخدام المجدول [Schedulers.immediate]، نحصل على النتائج التالية:
بينما باستخدام المجدول [Schedulers.trampoline]، نحصل على النتائج التالية:
7.8. Conclusion
لا يزال هناك الكثير مما يتعين القيام به. للتعمق في مكتبة RxJava، يُدعى القارئ إلى مواصلة تدريبه بالرجوع إلى المراجع المذكورة في بداية هذا المستند. ومع ذلك، لدينا الأساس اللازم لاستخدام RxJava في بيئات Swing وAndroid. وهذا ما سنقوم بشرحه الآن.








































