RxJava2 - अंतराल और अनुसूचक

Aug 24 2020

मान लीजिए कि मेरे पास एक अंतराल है, और मैंने इसे एक कम्प्यूटेशन दिया है। ऐशे ही:

Observable
    .interval(0, 1, TimeUnit.SECONDS, computationScheduler)
    .flatMap { ... }

फिर, क्या फ्लैटमेट {...} में होने वाली हर चीज को भी गणना सूत्र पर निर्धारित किया जाएगा?

ऑब्ज़र्वेबल.इन्टरवल (लंबी प्रारंभिक अवधि, लंबी अवधि, टाइम यूनिट यूनिट, शेड्यूलर शेड्यूलर) के स्रोतों में, यह कहता है:

 * @param scheduler
 * the Scheduler on which the waiting happens and items are emitted

RxJava के लिए एक शुरुआत के रूप में, मुझे इस टिप्पणी को समझने में मुश्किल समय आ रहा है। मैं समझता हूं कि गणना सूत्र पर अंतराल टाइमर / प्रतीक्षा तर्क होता है। लेकिन, क्या अंतिम भाग, उत्सर्जित होने वाली वस्तुओं के बारे में है, इसका मतलब यह भी है कि उत्सर्जित वस्तुओं का उपभोग उसी धागे पर किया जाएगा ? या इसके लिए एक अवलोकन की आवश्यकता है? ऐशे ही:

Observable
    .interval(0, 1, TimeUnit.SECONDS, computationScheduler)
    .observeOn(computationScheduler)
    .flatMap { ... }

यदि मैं चाहता हूं कि यदि मैं अभिकलन सूत्र पर कार्रवाई करना चाहता हूं तो क्या यह अवलोकन आवश्यक होगा?

जवाब

1 bubbles Aug 24 2020 at 21:34

यह सत्यापित करना आसान है: ऑपरेटर को निष्पादित करने वाले धागे को देखने के लिए बस वर्तमान थ्रेड प्रिंट करें:

Observable.just(1, 2, 3, 4, 5, 6, 7, 8, 9)
    .flatMap(e -> {
        System.out.println("on flatmap: " + Thread.currentThread().getName());
        return Observable.just(e).map(x -> "--> " + x);
    })
    .subscribe(s -> {
        System.out.println("on subscribe: " + Thread.currentThread().getName());
        System.out.println(s);
    });

यह हमेशा प्रिंट करेगा:

on subscribe: main
--> 1
on flatmap: main
on subscribe: main
--> 2
on flatmap: main
on subscribe: main
--> 3
on flatmap: main
on subscribe: main
--> 4
on flatmap: main
on subscribe: main
--> 5
on flatmap: main
on subscribe: main
--> 6
on flatmap: main
on subscribe: main
--> 7
on flatmap: main
on subscribe: main
--> 8
on flatmap: main
on subscribe: main
--> 9

क्रमिक रूप से संसाधित क्योंकि सभी एक ही धागे में होते हैं -> main।

observeOn डाउनस्ट्रीम निष्पादन थ्रेड को बदल देगा:

Observable.just(1, 2, 3, 4, 5, 6, 7, 8, 9)
    .observeOn(Schedulers.computation())
    .flatMap(e -> {
         System.out.println("on flatmap: " + Thread.currentThread().getName());
         return Observable.just(e).map(x -> "--> " + x);
     })
     .observeOn(Schedulers.io())
     .subscribe(s -> {
         System.out.println("on subscribe: " + Thread.currentThread().getName());
         System.out.println(s);
      });

परिणाम इस बार प्रत्येक निष्पादन के लिए अलग-अलग होगा, लेकिन flatmapऔर subscribeअलग धागे में कार्रवाई की जाएगी:

on flatmap: RxComputationThreadPool-1
on subscribe: RxCachedThreadScheduler-1

intervalके रूप में कार्य करेगा observeOnऔर डाउनस्ट्रीम निष्पादन थ्रेड (अनुसूचक) को बदल देगा:

Observable.interval(0, 1, TimeUnit.SECONDS, Schedulers.computation())
    .flatMap(e -> {
        System.out.println("on flatmap: " + Thread.currentThread().getName());
        return Observable.just(e).map(x -> "--> " + x);
    })
    .subscribe(s -> {
        System.out.println("on subscribe: " + Thread.currentThread().getName());
        System.out.println(s);
    });

इस बार निष्पादन संगणना अनुसूचक के एक धागे के अंदर अनुक्रमिक है:

on flatmap: RxComputationThreadPool-1
on subscribe: RxComputationThreadPool-1
--> 0
on flatmap: RxComputationThreadPool-1
on subscribe: RxComputationThreadPool-1
--> 1
on flatmap: RxComputationThreadPool-1
on subscribe: RxComputationThreadPool-1
--> 2
on flatmap: RxComputationThreadPool-1
on subscribe: RxComputationThreadPool-1
--> 3
...

intervalडिफ़ॉल्ट रूप से अभिकलन अनुसूचक का उपयोग करेगा, आपको इसे तर्क के रूप में पारित करने की आवश्यकता नहीं है और observeOnइसकी आवश्यकता नहीं है