RxJava2 - khoảng thời gian và bộ lập lịch
Giả sử tôi có một khoảng thời gian và tôi đã cấp cho nó một computationScheduler. Như thế này:
Observable
.interval(0, 1, TimeUnit.SECONDS, computationScheduler)
.flatMap { ... }
Sau đó, mọi thứ xảy ra trong bản đồ phẳng {...} cũng được lên lịch trên một chuỗi tính toán chứ?
Trong các nguồn cho Observable.interval (thời gian ban đầu dài, thời gian dài, đơn vị TimeUnit, bộ lập lịch Scheduler), nó cho biết:
* @param scheduler
* the Scheduler on which the waiting happens and items are emitted
Là một người mới bắt đầu sử dụng RxJava, tôi đang rất khó hiểu nhận xét này. Tôi hiểu rằng bộ đếm thời gian / logic chờ khoảng thời gian xảy ra trên luồng tính toán. Nhưng, phần cuối cùng, về các vật phẩm được phát ra, cũng có nghĩa là các vật phẩm được phát ra sẽ được tiêu thụ trên cùng một chủ đề? Hay là một ObserOn được yêu cầu cho điều đó? Như thế này:
Observable
.interval(0, 1, TimeUnit.SECONDS, computationScheduler)
.observeOn(computationScheduler)
.flatMap { ... }
Điều đó có cần thiết không nếu tôi muốn phát ra được xử lý trên luồng tính toán?
Trả lời
Điều này rất đơn giản để xác minh: chỉ cần in luồng hiện tại để xem toán tử được thực thi trên luồng nào:
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);
});
Điều này sẽ luôn in:
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
Được xử lý tuần tự vì tất cả đều diễn ra trong một chuỗi duy nhất -> main.
observeOn sẽ thay đổi luồng thực thi xuôi dòng:
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);
});
Kết quả lần này sẽ khác nhau đối với mỗi lần thực thi nhưng flatmapvà subscribesẽ được xử lý trong các chuỗi khác nhau:
on flatmap: RxComputationThreadPool-1
on subscribe: RxCachedThreadScheduler-1
intervalsẽ hoạt động như observeOnvà thay đổi luồng thực thi xuôi dòng (bộ lập lịch):
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);
});
Lần này việc thực thi diễn ra tuần tự bên trong một luồng của bộ lập lịch tính toán:
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
...
intervaltheo mặc định sẽ sử dụng bộ lập lịch tính toán, bạn không cần chuyển nó làm đối số và observeOnkhông cần thiết