Skip to content

2. مثال تمهيدي

كانت أولى تجاربي مع RxJava من خلال الدورات التدريبية والبرامج التعليمية التي عثرت عليها على الإنترنت. وبصرف النظر عن أن النظرية كانت تستخدم مفاهيم لم أكن معتادًا عليها وكنت أجد صعوبة في فهمها، فإنني لم أكن أرى على الإطلاق ما الفائدة منها في الحياة الواقعية. لذلك سنبدأ بتقديم مثال (آمل أن يكون بسيطًا) يوضح كيف أن استخدام RxJava يؤدي إلى تبسيط حقيقي في كتابة الكود، ومن هناك سنحاول تحديد العناصر المهمة في هذه المكتبة.

تعتمد مكتبة RxJava على المفهوم التالي: يتم مراقبة تدفق العناصر من النوع T Observable<T> بواسطة مشترك واحد أو أكثر (مشتركين، مراقبين، مستهلكين) Subscriber<T>. تسمح مكتبة RxJava بتشغيل تدفق Observable<T> في مؤشر ترابط T1 ومراقبه Subscriber<T> في مؤشر ترابط T2 دون أن يضطر المطوريضطر إلى القلق بشأن إدارة دورة حياة هذه الخيوط والمشكلات الصعبة بطبيعتها، مثل مشاركة البيانات بين الخيوط ومزامنتها لتنفيذ مهمة شاملة. وبالتالي، فهي تسهل البرمجة غير المتزامنة.

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

2.1. بنية التطبيق النموذجي

تتميز بنية التطبيق النموذجي بما يلي:

Image

  • في [1]، تقوم طبقة الخدمة بتوفير قوائم من الأرقام العشوائية. يتم تنفيذ هذه الطبقة في نفس الخيط الذي تعمل فيه الطريقة [swing] التي تستخدمها. وبذلك، فإنها توفر أرقامها بشكل متزامن؛
  • في [2]، تسمح طبقة تكييف رقيقة مُنفَّذة باستخدام RxJava بتقديم تنفيذ غير متزامن لنفس الخدمة إلى الطبقة [swing]: ويمكن تنفيذ هذا التنفيذ في مؤشر ترابط مختلف عن مؤشر ترابط الطريقة [swing] التي تستخدمه؛
  • الاستدعاء [4] متزامن، في حين أن الاستدعاء [5-6] غير متزامن؛

ما نريد توضيحه هنا هو أن مكتبة Rx تتيح تحويل واجهة متزامنة إلى واجهة غير متزامنة بسهولة. لماذا يعد هذا مفيدًا؟ تُعالج أحداث واجهة Swing في مؤشر ترابط يُعرف عمومًا باسم «حلقة الأحداث» (event loop). تُوضع الأحداث في قائمة انتظار وتُعالج واحدة تلو الأخرى. لا يمكن معالجة الحدث Ei+1 إلا بعد الانتهاء تمامًا من معالجة الحدث السابق Ei. لذا، من المهم أن تكون مدة معالجة الحدث أقصر ما يمكن حتى تظل الواجهة الرسومية سريعة الاستجابة. في بعض الأحيان، قد تستغرق معالجة الحدث وقتًا طويلاً. ويحدث ذلك إذا كانت هذه المعالجة تتطلب الوصول إلى الشبكة. وإذا أردنا تجنب تجميد واجهة المستخدم الرسومية بطريقة غير مقبولة بالنسبة للمستخدم، فيجب أن تتم عمليات الوصول إلى الشبكة هذه في خيوط منفصلة عن حلقة الأحداث لتحريرها. وهنا ندخل في مجال البرمجة المتزامنة (حيث يتم تنفيذ عدة خيوط بالتوازي)، والتي تُعتبر بحق صعبة. وتقدم مكتبة Rx حلاً بسيطًا وأنيقًا لهذه المشكلة.

لمحاكاة العمليات التي تستغرق وقتًا طويلاً، تقوم الخدمة الواردة في المثال بتقديم أرقامها العشوائية بعد فترة انتظار معينة حتى نتمكن من ملاحظة سلوك واجهة المستخدم الرسومية.

2.2. L'exécutable

يوجد الملف القابل للتنفيذ لتطبيق المثال في المجلد [dvp/executables] ضمن الأمثلة:

هناك طرق متنوعة لتشغيل الملف المضغوط [swing-01] وفقًا لإعدادات الجهاز المستخدم لتشغيله. يمكن، على سبيل المثال، اتباع الإجراء [1-3]. عندئذٍ تظهر واجهة المستخدم الرسومية التالية:

 
  • تحتوي الواجهة على علامتي تبويب [1-2]، إحداهما [Request] مخصصة لطلب خدمة توليد الأرقام العشوائية، والأخرى [Response] مخصصة لعرض الأرقام المستلمة؛
  • في [3]، يتم تحديد عدد الطلبات التي نريد إرسالها إلى الخدمة؛
  • في [4]، يُحدد النطاق [a,b] لتوليد الأرقام المطلوبة؛
  • في [5]، سيكون عدد القيم التي ترسلها الخدمة رقمًا عشوائيًا ضمن النطاق [minCount, maxCount] الذي يحدده المستخدم؛
  • في [6]، قبل إرجاع الرد، ستنتظر الخدمة delay مللي ثانية، حيث delay هو عدد عشوائي ضمن النطاق [minDelay, maxDelay] الذي يحدده المستخدم؛
  • بشكل افتراضي، ستتصل الطبقة [swing] بالواجهة المتزامنة للخدمة. وللتواصل مع الطبقة غير المتزامنة، سيقوم المستخدم بتحديد [7]. في هذه الحالة، ستُنفَّذ خدمة التوليد في خيوط منفصلة عن حلقة الأحداث الخاصة بالواجهة الرسومية. توفر مكتبة Rx استراتيجيات متنوعة لتوليد هذه الخيوط. يمكن للمستخدم اختيار استراتيجيته في [8]؛
  • يتم إنشاء الأرقام باستخدام الزر [9]؛
 
  • في [10]، عرض النتائج. سنشرح هيكل هذه النتائج؛
  • في [11]، عدد النتائج التي تم الحصول عليها؛
  • في [12]، مدة التنفيذ بالمللي ثانية؛
  • في [13]، يمكن للمستخدم إلغاء التنفيذ؛

تتخذ كل نتيجة الشكل التالي:

{"idClient":0,"serviceResponse":{"delay":412,"aleas":[146,115,128,174,159,112,162,127],"executedOn":"RxComputationThreadPool-6"},"observedOn":"AWT-EventQueue-0","requestAt":"02:42:47:708","responseAt":"02:42:52:931"}
  • [idClient]: رقم الطلب. تجدر الإشارة إلى أنه يتم إرسال عدة طلبات إلى خدمة الإنشاء؛
  • [delay]: وقت الانتظار بالمللي ثانية الذي سجلته الخدمة قبل إرسال النتيجة؛
  • [aleas]: الأرقام العشوائية التي أرسلتها الخدمة؛
  • [executedOn]: اسم الخيط الذي تم فيه تنفيذ الخدمة؛
  • [observedOn]: اسم الخيط الذي عرض النتيجة. مع واجهة Swing، لا يمكن أن يكون هذا سوى خيط حلقة الأحداث، وهو هنا [AWT-EventQueue-0
  • [requestAt]: وقت الطلب بالصيغة [heures:minutes:secondes:millisecondes
  • [responseAt]: وقت استلام النتائج بنفس الصيغة؛

سنعرض الآن أجزاء الكود المفيدة لفهم المثال.

2.3. الواجهة المتزامنة

Image

تقدم طبقة الخدمة [1] الواجهة التالية:


public interface IService {
  // أرقام عشوائية في [a,b]
  // يتم توليد n أرقام عشوائية في النطاق [minCount, maxCount]
  // يتم توليد الأعداد بعد انتظار مدته delay مللي ثانية،
  // حيث [delay] هو عدد عشوائي في النطاق [minDelay, maxDelay]
  public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}

الاستجابة [ServiceResponse] هي كما يلي:


public class ServiceResponse {

  // فترة انتظار الخدمة
  private int delay;
  // أرقام عشوائية
  private List<Integer> aleas;
  // خيط التنفيذ
  private String executedOn;

  // منشئات

  public ServiceResponse(int delay, List<Integer> aleas) {
    executedOn = Thread.currentThread().getName();
    this.delay = delay;
    this.aleas = aleas;
  }

  // المنشورات والمُعيّنات
...
}

تتكون الإجابة من ثلاثة عناصر:

  • السطر 6: الأرقام العشوائية التي تم توليدها؛
  • السطر 4: فترة الانتظار التي يلتزم بها الخدمة قبل إرجاع النتيجة؛
  • السطر 8: مؤشر ترابط تنفيذ الخدمة؛

2.4. الاستدعاء المتزامن

Image

نستعرض الآن بالتفصيل الاستدعاء المتزامن [4] الذي تقوم به الطبقة [swing] إلى الخدمة [1]:


  private void doGenerateWithService() {
    // بداية الانتظار
    beginWaiting();
    try {
      for (int i = 0; i < nbRequests; i++) {
        UiResponse uiResponse = new UiResponse();
        uiResponse.setIdClient(i);
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        uiResponse.setResponseAt();
        model.add(0, jsonMapper.writeValueAsString(uiResponse));
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
      }
    } catch (JsonProcessingException | RuntimeException e) {
      System.out.println(e);
    }
    // نهاية الانتظار
    endWaiting();
}
  • الأسطر 5-12: حلقة تنفيذ طلبات [nbRequests] التي يطلبها المستخدم؛
  • السطر 8: [service] هو تنفيذ الواجهة المتزامنة [IService] المذكورة في الفقرة 2.3؛
  • السطر 10: [model] هو النموذج الذي يعرضه المكون JList في علامة التبويب [Response]. عناصر هذا النموذج هي سلاسل jSON من العناصر من النوع [UiResponse] التالية:

public class UiResponse {

  // معرف العميل
  private int idClient;
  // استجابة الخدمة
  private ServiceResponse serviceResponse;
  // اسم مؤشر الترابط المراقب
  private String observedOn;
  // وقت الطلب
  private String requestAt;
  // وقت الرد
  private String responseAt;

  // المنشئات

  public UiResponse() {
    observedOn = Thread.currentThread().getName();
    requestAt = getTimeStamp();
  }
  // الطرق الخاصة

  private String getTimeStamp() {
    return new SimpleDateFormat("hh:mm:ss:SSS").format(Calendar.getInstance().getTime());
  }

  // المنشورات والمُعيّنات
...
}
  • السطر 6: رد خدمة توليد الأرقام؛
  • السطر 4: رقم الطلب الذي تم الرد عليه؛
  • السطر 8: مؤشر ترابط عرض هذا الرد. وكما ذكرنا سابقًا، سيكون هذا دائمًا مؤشر ترابط حلقة الأحداث؛
  • السطران 10 و12: وقت الطلب ووقت الرد؛

2.5. اختبارات المكالمات المتزامنة

نقوم بتنفيذ التكوين التالي:

 

ونحصل على النتائج التالية في علامة التبويب [Response]:

 
  • في [1-2]، حصلنا بالفعل على 10 استجابات كما كان مطلوبًا. وقد تم إدراجها في المرتبة الأولى حسب ترتيب وصولها. ونلاحظ أنها تم الحصول عليها حسب ترتيب الطلبات؛
  • وقد تم تنفيذها جميعًا وعرضها في مؤشر ترابط حلقة الأحداث [AWT-EventQueue-0]. وبالتالي، تم تنفيذ الطلبات واحدًا تلو الآخر في هذا المؤشر. لم تكن هناك أي طلبات متزامنة؛
  • ما لا يظهر هنا هو أن واجهة المستخدم الرسومية تتجمد أثناء التنفيذ. على سبيل المثال، لا توجد طريقة للوصول إلى علامة التبويب [Response] لمشاهدة وصول الردود أو إيقاف التنفيذ باستخدام الزر [Annuler]. وحتى لو كان هذا الزر موجودًا في علامة التبويب [Request]، لكان غير قابل للاستخدام. ففي هذه الحالة، سيكون هناك حدثان:
    • النقر على الزر [Générer
    • النقر على الزر [Annuler

ولا تتم معالجة النقر على الزر [Annuler] إلا بعد انتهاء العملية التي تم تشغيلها بالنقر على الزر [Générer]. لقد رأينا للتو أن هذه العملية كانت تشغل مؤشر ترابط حلقة الأحداث طوال مدة التنفيذ، مما منع معالجة النقر على الزر [Annuler]. هذا هو النوع النموذجي من المواقف التي يمكن أن يحقق فيها Rx تحسناً ملحوظاً؛

2.6. الواجهة غير المتزامنة وتنفيذها

ننتقل الآن إلى واجهة الطبقة [2] وكذلك إلى تنفيذها باستخدام Rx. لن تكون هذه الواجهة مفهومة على الفور. نريد ببساطة إبراز بساطة كود هذا التنفيذ.

الواجهة غير المتزامنة هي كما يلي:


public interface IRxService {
  // أرقام عشوائية في [a,b]
  // يتم توليد n أعداد عشوائية في النطاق [minCount, maxCount]
  // يتم توليد الأعداد بعد انتظار مدته delay مللي ثانية،
  // حيث [delay] هو عدد عشوائي في النطاق [minDelay, maxDelay]
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}

فيما يلي الاختلافات عن الواجهة المتزامنة التي تم عرضها في الفقرة 2.3:

  • أصبحت الفئة [UiResponse] المذكورة في الفقرة 2.3 الآن جزءًا من معلمات الطريقة [getAleas] (السطر 6). والسبب في ذلك هو أنه، نظرًا لأن الطلبات تُنفَّذ الآن بشكل متوازٍ وأن الخدمة تنتظر فترة زمنية عشوائية قبل إرجاع نتيجتها، فإن الردود لن تصل إلينا بترتيب الطلبات. لذا، نقوم بتمرير الكائن [UiResponse] الذي يحتوي، من بين معلومات أخرى، على رقم الطلب:

  // معرف العميل (الطلب)
  private int idClient;
  // استجابة الخدمة
  private ServiceResponse serviceResponse;
  // اسم مؤشر الترابط المراقب
  private String observedOn;
  // وقت الطلب
  private String requestAt;
  // وقت الرد
  private String responseAt;
  • نوع استجابة الخدمة غير المتزامنة هو [Observable<UiResponse>]. أما النوع [Observable<>] فيتم توفيره بواسطة مكتبة Rx. تشير النتيجة من النوع [Observable<UiResponse>] إلى أن الطريقة [getAleas] توفر تدفقًا من القيم من النوع [UiResponse]، وهي قيم يتم دفعها (pushed) واحدة تلو الأخرى إلى المراقب الخاص بها؛

لنلقِ نظرة الآن على تنفيذ هذه الواجهة:


public class RxService implements IRxService {

  // الخدمة
  private IService service;

  // المُنشئ
  public RxService(IService service) {
    this.service = service;
  }

  @Override
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
    return Observable.create(subscriber -> {
      try {
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
        subscriber.onError(e);
      } finally {
        subscriber.onCompleted();
      }
    });
  }
}
  • الأسطر 7-9: يتم تزويد المنشئ بإشارة إلى الواجهة المتزامنة [IService]. وهي التي ستتولى مهمة توليد الأرقام العشوائية؛
  • يتم إنشاء المراقب الذي ترجعها الطريقة [getAleas] بواسطة الطريقة الثابتة [Observable.create]. وهذه الطريقة هي التي تسمح بإنشاء تنفيذ غير متزامن انطلاقًا من تنفيذ متزامن؛
  • السطر 13: المعلمة الخاصة بالطريقة الثابتة [Observable.create] هي هنا دالة لامدا تتلقى كمعلمة نوعًا [Subscriber]، وهو أيضًا نوع Rx. [Subscriber] هو كائن يشترك في تدفق من العناصر القابلة للمراقبة، أي تدفق بيانات يتم تسليمها بشكل غير متزامن. نستخدم هنا ثلاث طرق لهذا المشترك:
    • [Subscriber.onNext] لإرسال بيانات إليه (السطر 16)؛
    • [Subscriber.onError] لإرسال استثناء إليه (السطر 18)؛
    • [Subscriber.onCompleted] لإعلام المشترك بأن تدفق البيانات قد انتهى (السطر 20)؛

قد يكون هناك عدة مشتركين في نفس العنصر القابل للمراقبة. هنا، سيكون لدينا مشترك واحد فقط يشترك في تدفق لبيانات واحدة، وهي تلك التي يتم إنتاجها في السطرين 15-16. يتم إنتاج البيانات من خلال التنفيذ المتزامن للخدمة (السطر 15) وتسليمها إلى المشترك (السطر 16).

على الرغم من أن كل هذا قد يبدو غامضًا، إلا أننا لا نستطيع إلا أن نندهش من الإيجاز الشديد لهذا التنفيذ غير المتزامن للخدمة.

2.7. الاستدعاء غير المتزامن

Image

سنقوم الآن بتفصيل الاستدعاء المتزامن [5] الذي تقوم به الطبقة [swing] إلى الخدمة [2]:


private void doGenerateWithRxService() {
        // بداية الانتظار
        beginWaiting();
        // يُطلب توفير الأرقام العشوائية
        Observable<UiResponse> observables = Observable.empty();
        for (int i = 0; i < nbRequests; i++) {
            UiResponse uiResponse = new UiResponse();
            uiResponse.setIdClient(i);
            // المجدول
            int schedulerIndex = jComboBoxSchedulers.getSelectedIndex();
            switch (schedulerIndex) {
            case 0:
                observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.io()));
                break;
...
            }
        }
...
    }
  • الأسطر 6-10: تنفيذ طلبات [nbRequests] التي طلبها المستخدم؛
  • الأسطر 7-8: إعداد الكائن [UiResponse] الذي تحتاجه الطريقة [getAleas] للخدمة غير المتزامنة (السطر 13). ويتمثل ذلك أساسًا في تسجيل رقم الطلب [idClient
  • السطر 13: يتم استدعاء الأسلوب [getAleas] الخاص بالخدمة غير المتزامنة. وهو يُرجع كائن [Observable<UiResponse>]. لا يستدعي هذا الاستدعاء الخدمة المتزامنة بعد. لنعد إلى كود الأسلوب غير المتزامن [getAleas]:

  @Override
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
    return Observable.create(subscriber -> {
      try {
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
        subscriber.onError(e);
      } finally {
        subscriber.onCompleted();
      }
    });
}

لا يتم تنفيذ الأسطر 4-11 من الكود التي ستستدعي الخدمة المتزامنة إلا عندما يقوم أحد المشتركين بالتسجيل. وطالما لا يوجد مشتركون، فإن هذا الكود لا يتم تنفيذه.

لنعد إلى كود الطريقة [doGenerateWithRxService]:

  • السطر 5: يتم إنشاء متغير قابل للمراقبة فارغ (لا يتم مراقبة أي شيء)؛
  • السطر 13: يتم إنشاء كائن قابل للمراقبة يكون تدفقه عبارة عن دمج التدفقات غير المتزامنة [nbRequests] المرتبطة بالطلبات [nbRequests]. ويتم تحقيق ذلك باستخدام الطريقة [Observable.mergeWith] التي تسمح بدمج تدفقين غير متزامنين. في مصطلحات Rx، يُطلق على [mergeWith] اسم «مشغل التدفق». وتتميز هذه المشغلات بأن نتيجة العملية تكون في أغلب الأحيان [Observable] مرة أخرى. وفي النهاية، بعد السطر 17، تشير المتغير [observables] إلى تدفق واحد يتكون من الردود غير المتزامنة [nbRequests] الصادرة عن الخدمة غير المتزامنة؛
  • السطر 13: كان من الممكن كتابة عملية الدمج على النحو التالي:

observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse));

لكننا كتبنا:


observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.io()));

استخدمنا هنا المشغل [subscribeOn] على المتغير القابل للمراقبة [rxService.getAleas]. وكما هو الحال غالبًا، تكون النتيجة متغيرًا قابلًا للمراقبة مرة أخرى. يسمح المشغل [subscribeOn] بتحديد أن القيمة القابلة للمراقبة يجب أن تُنفَّذ في مؤشر ترابط (thread) مقدم من قبل [Scheduler]. هناك العديد من [Scheduler] الممكنة والمناسبة لمواقف مختلفة. في واجهة المستخدم الرسومية، قمنا بعرض العديد منها لمشاهدة تأثيرات كل منها:

  

وينتج عن ذلك الكود التالي:


    private void doGenerateWithRxService() {
        // بدء الانتظار
        beginWaiting();
        // يتم طلب الأرقام العشوائية
        Observable<UiResponse> observables = Observable.empty();
        for (int i = 0; i < nbRequests; i++) {
            UiResponse uiResponse = new UiResponse();
            uiResponse.setIdClient(i);
            // المجدول
            int schedulerIndex = jComboBoxSchedulers.getSelectedIndex();
            switch (schedulerIndex) {
            case 0:
                observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.io()));
                break;
            case 1:
                observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.computation()));
                break;
            case 2:
                observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.newThread()));
                break;
            case 3:
                observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.trampoline()));
                break;
            case 4:
                observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.immediate()));
                break;
            }
        }
...
}

لنعد إلى الكود في الأسطر 12-14. يقوم المجدول [Schedulers.io()] بتخصيص مؤشر ترابط جديد لكل متغير قابل للمراقبة. إذا تابعنا الكود:

  • السطر 5: لدينا متغير قابل للمراقبة فارغ؛
  • السطر 13، التكرار 1: observables هي القائمة [observable0/thread0] (المتغير القابل للمراقبة observable0 الذي يتم تنفيذه على الخيط thread0
  • السطر 13، التكرار 2: observables هي القائمة [observable0/thread0, observable1/thread1
  • إلخ...

في النهاية، بعد السطر 28، نحصل على متغير قابل للمراقبة ناتج عن دمج المتغيرات القابلة للمراقبة [nbRequests] التي تُنفَّذ على خيوط مختلفة [nbRequests]. لا تعمل جميع برامج الجدولة بهذه الطريقة، كما سنرى خلال الاختبارات.

لنواصل دراسة كود استدعاء الخدمة غير المتزامنة:


private void doGenerateWithRxService() {
        // بدء الانتظار
        beginWaiting();
        // يتم طلب الأرقام العشوائية
        Observable<UiResponse> observables = Observable.empty();
        for (int i = 0; i < nbRequests; i++) {
        ...
        }
        // المراقب
        observables = observables.observeOn(SwingScheduler.getInstance());
        // يتم تنفيذ هذه المتغيرات القابلة للمراقبة
        subscriptions.add(observables.subscribe(uiResponse -> {
            updateUi(uiResponse);
        } , th -> {
            System.out.println(th);
            doCancel();
        } , this::doCancel));
    }
  • لقد رأينا أنه عند الوصول إلى السطر 10، يكون لدينا عنصر قابل للمراقبة واحد، وهو دمج لـ [nbRequests] من العناصر القابلة للمراقبة التي يمكن أن تُنفَّذ على [nbRequests] خيوط مختلفة أو لا، وفقًا لجدولة العمليات التي يختارها المستخدم؛
  • السطر 10: يتيح المشغل [observeOn] تحديد الخيط الذي نريد استرداد البيانات منه من العنصر القابل للمراقبة، وهنا كائنات من النوع [UiResponse]. في واجهة Swing، لا يوجد خيار آخر. يجب أن يتم أي تحديث للواجهة في مؤشر الترابط الخاص بحلقة الأحداث. هنا، سيتم عرض بيانات العنصر القابل للمراقبة في مكون Swing من النوع JList. يمثل الخيط [SwingScheduler.getInstance()] خيط حلقة الأحداث. لا تنتمي الفئة [SwingScheduler] إلى المكتبة RxJava بل إلى المكتبة المشتقة RxSwing؛
  • عند الوصول إلى السطر 12، لم يتم استدعاء الخدمة المتزامنة بعد لأن العنصر القابل للمراقبة في السطر 10 لا يزال بدون مشترك. وتقوم الأسطر 12-17 بتزويده بمشترك، بفضل المشغل [subscribe]. ومعلمات هذا المشغل هنا هي ثلاث دوال لامدا:
    • الأولى [uiResponse -> {updateUi(uiResponse);}] تقبل كمعلمة أحد الكائنات [UiResponse] التي تنتجها المراقبة. تجدر الإشارة إلى أننا سنحصل هنا على كائنات من هذا النوع. ويجب أن تستفيد الطريقة المرتبطة بها، وهي updateUi في هذه الحالة، من هذه النتيجة؛
    • تقبل الدالة الثانية [th -> {System.out.println(th);doCancel();}] كمعلمة نوعًا من النوع [Throwable]، وهو في هذه الحالة استثناء حدث أثناء تنفيذ المراقب. يجب أن تستفيد الطريقة المرتبطة بهذه المعلومة. هنا، يتم عرضها على وحدة التحكم (السطر 15) ويتم إلغاء التنفيذ، مما سيؤدي إلى تحديث بعض عناصر واجهة المستخدم الرسومية؛
    • يتم استدعاء الدالة الثالثة [this::doCancel] عندما يشير المراقب إلى أنه لم يعد لديه بيانات لإرسالها. هنا، المراقب هو تجميع لمراقبين من نوع [nbRequests]. ستشير المراقبة الناتجة إلى أنها انتهت عندما تكون جميع المراقبات المكونة لها قد أشارت بدورها إلى أنها أنهت عملها. لذا، عند تنفيذ دالة لامدا الثالثة هذه، يكون قد تم استلام جميع البيانات. تقوم الطريقة المحلية [doCancel] بتحديث واجهة المستخدم الرسومية لتعكس انتهاء التنفيذ؛

يتم تعريف المتغير [subscriptions] على النحو التالي:


    // الاشتراكات في القيم القابلة للمراقبة
protected List<Subscription> subscriptions = new ArrayList<Subscription>();

يمثل النوع [Subscription] اشتراكًا، أي الارتباط بين المشترك [Subscriber] وما يراقبه [Observable]. استخدمنا هنا قائمة بالاشتراكات على الرغم من وجود اشتراك واحد فقط في هذا المثال. الطريقة المحلية [doCancel] التي يتم تنفيذها عندما يشير العنصر القابل للمراقبة إلى أنه لم يعد لديه بيانات لإرسالها هي كما يلي:


    @Override
    protected void doCancel() {
        // نهاية الانتظار
        endWaiting();
        // في حالة الاشتراكات
        if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
            subscriptions.forEach(Subscription::unsubscribe);
        }
}
  • السطر 7 يلغي اشتراك جميع المشتركين في العنصر القابل للمراقبة؛

من هذا الشرح الموجز، يمكن استخلاص النقاط الرئيسية التالية:

  • يشير النوع [Observable] إلى تدفق للقيم، وهي قيم يتم دفعها واحدة تلو الأخرى إلى المشتركين أو المراقبين؛
  • يشير النوع [Subscriber] إلى مشترك من النوع [Observable
  • يشير النوع [Subscription] إلى اشتراك، أي الرابط بين [Subscriber] و[Observable
  • النوع [Observable] يقبل مشغلات من النوع [mergeWith, empty, subscribeOn, observeOn, ...] التي تنتج في الغالب قيمًا قابلة للمراقبة. تُستخدم هذه المشغلات لتكوين القيمة القابلة للمراقبة قبل تنفيذها:
    • ما نريد ملاحظته؛
    • الخيط الذي يتم تنفيذ المتغير القابل للمراقبة عليه؛
    • الخيط الذي يتلقى عليه المشترك بيانات القابل للمراقبة؛
  • يتم التمييز بين نوعين من الملاحظات، وهما [froid / cold] و [chaud / hot]. يتم تنفيذ الملاحظة الباردة بالكامل عند كل مشترك جديد. إذا كان كل تنفيذ ينتج نفس البيانات، فإن كل مشترك جديد يتلقى نفس البيانات التي تلقّاها المشترك السابق. أما المراقب «الساخن» فينتج البيانات بشكل مستمر عمومًا. وعندما يشترك مشترك ما، يتلقى البيانات التي تم إرسالها بدءًا من وقت اشتراكه. ولا يتلقى البيانات التي ربما تم إرسالها سابقًا. في مثالنا، المراقب بارد: يتم إعادة تنفيذه بالكامل عند كل مشترك جديد. ما الذي يتم تنفيذه فعليًّا في مثالنا؟ لمعرفة ذلك، علينا الرجوع إلى تعريف المراقب الذي نراقبه:

  @Override
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
    return Observable.create(subscriber -> {
      try {
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
        subscriber.onError(e);
      } finally {
        subscriber.onCompleted();
      }
    });
}

مع كل مشترك جديد، يتم إعادة تنفيذ دالة لامدا، وهي معلمة الطريقة [Observable.create] (السطر 3). وبالتالي، يتم تنفيذ الأسطر من 4 إلى 11 لكل مشترك جديد في [subscriber

2.8. اختبارات المكالمات غير المتزامنة

نبدأ بعرض تأثير المجدولات المختلفة المقترحة. ونستخدم لهذا الغرض المعلمات التالية:

 

نضع قيمًا صغيرة في [1-2] حتى لا ننتظر طويلاً في حالة تنفيذ الطلبات على نفس الخيط.

2.8.1. مع برنامج الجدولة [Schedulers.io]

 

يمكن ملاحظة النقاط التالية:

  • نحصل على الردود بترتيب يختلف عن ترتيب الطلبات (انظر idClient
  • تم تنفيذ كل طلب في مؤشر ترابط مختلف؛
  • لم تعد الواجهة الرسومية ثابتة هذه المرة:
    • يمكن الانتقال من علامة تبويب إلى أخرى؛
    • نرى البيانات وهي تصل؛
    • لا يتسنى لنا رؤية الزر [Annuler] لأن التنفيذ سريع جدًّا. سنسلط الضوء عليه في اختبار آخر؛

2.8.2. مع المجدول [Schedulers.computation]

 

يمكن ملاحظة النقاط التالية:

  • نحصل على الردود بترتيب يختلف عن ترتيب الطلبات (انظر idClient
  • تم تنفيذ الطلبات في 8 خيوط؛
  • تم استخدام الخيط رقم 3 للطلبين 8 و0؛
  • تم استخدام الخيط رقم 4 للطلبين 9 و1؛
  • كل طلب من الطلبات الأخرى استخدم خيطًا مختلفًا؛

يستخدم المجدول [Schedulers.computation] عددًا من الخيوط يساوي عدد النوى في الجهاز المستخدم. يتم الحصول على هذه المعلومة من خلال التعبير [Runtime.getRuntime().availableProcessors()].

2.8.3. مع المجدول [Schedulers.newThread]

 

يتم تشغيله بطريقة مشابهة لتشغيل المجدول [Schedulers.io].

2.8.4. مع المجدولات [Schedulers.trampoline, Schedulers.immediate]

 

نلاحظ أن العمل يتم بشكل متزامن. يتم تنفيذ جميع الاستعلامات على مؤشر ترابط حلقة الأحداث. لا ينبغي تعميم هذه النتيجة، بل يجب القول ببساطة إن المجدولين في هذا المثال المحدد بالذات عملا بشكل متزامن.

2.9. الحالات الحدية

سنعمل في هذا المثال مع المجدولات التي تسمح بالتشغيل غير المتزامن. أولاً، نزيد عدد الطلبات إلى 100 باستخدام المجدول [Schedulers.computation] الذي يعمل هنا بـ 8 خيوط. ونحصل على النتيجة التالية:

 
  • في [1]، الزر [Annuler] موجود وقابل للاستخدام (تشغيل غير متزامن)؛

الآن، دعونا نكمل التنفيذ حتى النهاية:

 

نلاحظ في [2] أن تنفيذ الـ 100 استعلام استغرق حوالي 4 ثوانٍ (على 8 خيوط).

والآن، دعونا نجري نفس الـ 100 استعلام باستخدام المجدول [Schedulers.newThread] الذي ينفذ كل استعلام على خيط منفصل:

 

في [1]، نلاحظ أن تنفيذ الـ 100 استعلام (على 100 خيط) استغرق نصف ثانية. وهذا أسرع بكثير من المجدول [Schedulers.computation].

الآن، لنقم بإجراء 800 طلب في نفس الظروف باستخدام المجدول [Schedulers.newThread]. نحصل على النتائج التالية:

 

يتم تنفيذ الـ 800 طلب في حوالي ثانية واحدة.

وعند زيادة هذا العدد (إلى ما يزيد عن 2500 استعلام على جهازي - يتم تنفيذها في 1.5 ثانية - وهذا العدد يعتمد بالطبع بشكل كبير على بيئة العمل وقت التنفيذ)، نحصل في النهاية على الاستثناء التالي:

  

إذن، لدينا تجاوز سعة المكدس. تُظهر الاختبارات أن عمل المجدول [Schedulers.newThread] ليس حتميًا. فقد تظهر الاستثناء السابق، ثم نجري تجارب جديدة، ثم نعود إلى التكوين الذي تسبب في الاستثناء ولا يظهر مرة أخرى.

2.10. Conclusion

لقد عرضنا مثالاً على استخدام مكتبة Rx. فلنلخص ما تعلمناه:

لقد انطلقنا من البنية التالية:

Image

  • في [4]، كانت الطبقة [swing] تُجري استدعاءات متزامنة إلى الطبقة [service
  • في [5]، كانت الطبقة [swing] تُجري استدعاءات غير متزامنة إلى الطبقة [rxService] التي كانت بدورها تستدعي الطبقة [service] بشكل متزامن من خلال الطبقة [6]؛

أول ما لاحظناه هو أن مكتبة Rx تتيح إنشاء الواجهة غير المتزامنة [rxService] بسهولة انطلاقًا من الواجهة المتزامنة [service] (انظر الفقرة 2.4). وهذا درس مهم لأنه يعني أنه يمكن بسهولة تحويل تطبيق متزامن إلى تطبيق غير متزامن.

في الطبقة [swing]، تمت كتابة طريقتين منفصلتين:

  • إحداهما لإجراء استدعاءات متزامنة للخدمة (انظر الفقرة 2.4
  • والأخرى لإجراء استدعاءات غير متزامنة لها (انظر الفقرة 2.7

وقد تبين أن كتابة الاستدعاءات غير المتزامنة أكثر تعقيدًا بكثير من كتابة الاستدعاءات المتزامنة. ومع ذلك، فإن من سبق لهم العمل في البرمجة المتزامنة مع عدة خيوط (threads) تحتاج إلى التزامن، سيجدون أن حل Rx أسهل في الكتابة ويتجنب جميع مشاكل التزامن والتواصل بين الخيوط، وهي مشاكل صعبة. وأثناء كتابة هذا النص، ميزنا النقاط المهمة التالية:

  • يشير النوع [Observable] إلى تدفق للأحداث (القيم) التي قد تكون (ولكن ليس بالضرورة) غير متزامنة ويمكن ملاحظتها؛
  • يشير النوع [Subscriber] إلى مشترك في النوع [Observable
  • يشير النوع [Subscription] إلى اشتراك، أي الرابط بين [Subscriber] و [Observable
  • النوع [Observable] يقبل مشغلات من النوع [mergeWith, empty, subscribeOn, observeOn, ...] التي تنتج في الغالب قيمًا قابلة للمراقبة. تُستخدم هذه المشغلات لتكوين القيمة القابلة للمراقبة قبل تنفيذها:
    • ما نريد ملاحظته؛
    • الخيط الذي يتم تنفيذ المتغير القابل للمراقبة عليه؛
    • الخيط الذي يتلقى عليه المشترك البيانات من القابل للمراقبة؛
  • يتم التمييز بين نوعين من الملاحظات، وهما [froid / cold] و [chaud / hot]. يتم تنفيذ الملاحظة الباردة بالكامل عند كل مشترك جديد. إذا كان كل تنفيذ ينتج نفس البيانات، فإن كل مشترك جديد يتلقى نفس البيانات التي تلقّاها المشترك السابق. أما المراقب «الساخن» فينتج البيانات بشكل مستمر عمومًا. عندما يشترك مشترك ما، يتلقى البيانات التي تم إرسالها بدءًا من وقت اشتراكه. ولا يتلقى البيانات التي ربما تم إرسالها سابقًا. في مثالنا، المراقب «بارد»: يتم إعادة تنفيذه بالكامل مع كل مشترك جديد.

والآن بعد أن اطلعنا على مثال يوضح فائدة مكتبة Rx، سنقدمها بمزيد من التفصيل.

تحتوي مكتبة Rx على العديد من الطرق التي تحتوي على معلمات عامة في توقيعها. سنقدم لمحة موجزة عن هذه التوقيعات (الفقرة 3). غالبًا ما تكون معلمات هذه الطرق واجهات وظيفية (Java 8)، أي واجهات لا تحتوي إلا على طريقة واحدة. لذا يجب أن تكون المعلمات الفعلية مثيلات لهذه الواجهات. قبل Java 8، كان من المعتاد تنفيذ واجهة ما بواسطة فئة مجهولة. مع Java 8، وإذا كانت الواجهة واجهة وظيفية، فمن الأكثر إيجازًا تنفيذها باستخدام دالة لامدا. لذا سنقدم هذه الدوال (الفقرة 4). وبعد الانتهاء من ذلك، سنقدم الفئة [Stream] (الفقرة 5) التي تتيح معالجة مجموعات Java باستخدام دوال لامدا. وتعد هذه الفئة مثيرة للاهتمام لأن الفئة [Observable] من RxJava تستعير منها:

  • بعض الطرق؛
  • نفس طريقة ربط الطرق ببعضها البعض لمعالجة نفس المراقب؛

سنقدم بعد ذلك الواجهات الوظيفية الخاصة بمكتبة RxJava (الفقرة 6). وسنواصل مع العناصر الرئيسية لمكتبة Rx [Observable, Subscriber, Subscription, opérateurs] (الفقرة 7). تحتوي الفئة [Observable] على عشرات من العوامل التي يتم إعادة تحميلها عدة مرات. وهذا يخلق في البداية تعقيدًا كبيرًا لأن هذه العوامل وإعادة تحميلاتها لا تختلف أحيانًا إلا في تفصيل واحد، ومن الصعب، دون خبرة، معرفة العامل الذي يجب استخدامه. لن نعرض سوى عدد محدود من العوامل، وسنتجاهل في معظم الأحيان عمليات إعادة تعريفها.

سيتم تنفيذ الجزء السابق بالكامل باستخدام المكتبة RxJava في تطبيقات بسيطة تعمل على وحدة التحكم. وبمجرد إتقان المكتبة RxJava، سنستخدمها في نوعين من التطبيقات الرسومية:

  • في الفقرة سنعود إلى تطبيق Swing النموذجي لتفصيله بشكل أكبر. وسنستخدم حينها المكتبة RxSwing؛
  • في الفقرة سننشئ تطبيقًا لنظام Android باستخدام المكتبة RxAndroid؛

عندما ينتهي القارئ من كل هذا، سيكون لديه الأدوات اللازمة ليعتمد على نفسه. ومن المحتمل أن يستغرق الأمر بعض الوقت قبل أن يتمكن من استخدام مكتبة Rx بطريقة بديهية. لقد وجدت هذه المكتبة مثيرة للاهتمام بشكل خاص. ومع ذلك، فقد وجدتُها معقدة في الفهم، واستغرق تعلمها وقتًا طويلاً. آمل أن يساعد هذا المستند في تقصير مدة التعلم بالنسبة للقارئ. يبدو لي أن الأمر يستحق العناء.