Планировщик
Если вы хотите добавить многопоточность в вашу цепочку операторов Observable, вы можете это сделать, указав тем операторам (или конкретным Observable) работать на определенном Планировщике.
Некоторые операторы ReactiveX Observable имеют варианты, которые принимают планировщик в качестве параметра. Эти варианты инструктируют оператор выполнить часть или всю свою работу на конкретном планировщике.
По умолчанию Observable и цепочка операторов, которые к нему применяются, выполняют свою работу и уведомляют своих наблюдателей на той же нити, на которой вызывается его Subscribe метод. Оператор SubscribeOn изменяет это поведение, указывая другой планировщик, на котором Observable должен работать. Оператор ObserveOn указывает другой планировщик, который Observable будет использовать для отправки уведомлений своим наблюдателям.
Как показано на этой иллюстрации, оператор SubscribeOn определяет, на какой нити Observable начнет работать, независимо от того, в какой точке цепочки операторов этот оператор вызван. ObserveOn, с другой стороны, влияет на нить, которую Observable будет использовать ниже, где этот оператор появляется. По этой причине вы можете вызывать ObserveOn несколько раз в различных точках цепочки операторов Observable, чтобы изменить, на каких нитях работают определенные из этих операторов.
См. также
- Введение в Rx: Планирование и потоки выполнения
- Практическое руководство по Rx: Планировщики
- Использование планировщиков Дэннисом Стояновым
Информация, специфичная для языка
RxGroovy
Разновидности планировщиков
Планировщик получается из фабричных методов, описанных в классе Schedulers. В следующей таблице показаны разновидности планировщиков, доступные вам с помощью этих методов в RxGroovy:
| Планировщик | Назначение |
|---|---|
Schedulers.computation( ) |
предназначен для вычислительной работы, такой как обработка событий и обратных вызовов; не используйте этот планировщик для операций ввода-вывода (используйте Schedulers.io( ) вместо него); число потоков по умолчанию равно числу процессоров |
Schedulers.from(executor) |
использует указанный Executor в качестве планировщика |
Schedulers.immediate( ) |
планирует работу для начала немедленно в текущем потоке |
Schedulers.io( ) |
предназначен для операций ввода-вывода, таких как асинхронное выполнение блокирующих операций ввода-вывода; этот планировщик поддерживается пулом потоков, который будет увеличиваться по мере необходимости; для обычных вычислительных задач переключитесь на Schedulers.computation( ); Schedulers.io( ) по умолчанию — CachedThreadScheduler, что-то вроде планировщика новых потоков с кешированием потоков |
Schedulers.newThread( ) |
создаёт новый поток для каждой единицы работы |
Schedulers.trampoline( ) |
выстраивает работу в очередь для начала в текущем потоке после всех уже помещённых в очередь работ |
Планировщики по умолчанию для операторов RxGroovy Observable
Некоторые операторы Observable в RxGroovy имеют альтернативные формы, которые позволяют указать, какой планировщик будет использовать оператор для (по крайней мере, части) своей работы. Другие не работают с каким-либо конкретным планировщиком или работают с конкретным планировщиком по умолчанию. К тем, которые используют конкретный планировщик по умолчанию, относятся:
Тестовый планировщик
TestScheduler позволяет осуществлять тонкую ручную настройку поведения часов планировщика. Это может быть полезно для тестирования взаимодействий, которые зависят от точного расположения действий во времени. У этого планировщика есть три дополнительных метода:
advanceTimeTo(time,unit)- передвигает часы планировщика к определённой точке во времени
advanceTimeBy(time,unit)- передвигает часы планировщика вперёд на определённое количество времени
triggerActions( )- запускает все неначатые действия, запланированные на время, равное или меньшее текущего времени по часам планировщика
См. также
- Тестирование реактивных приложений Беном Кристенсеном
- Примеры потоков RxJava Гремом Ли
- Расширенный RxJava: планировщики (часть 1) (часть 2) (часть 3) (часть 4) Давидом Карноком
RxJava 1․x
Разновидности планировщиков (Schedulers)
Планировщик (Scheduler) вы получаете из фабричных методов, описанных в классе Schedulers. В следующей таблице показаны разновидности планировщиков (Schedulers), доступные вам в RxJava с помощью этих методов:
| Планировщик | назначение |
|---|---|
Schedulers.computation( ) |
предназначен для вычислений, таких как обработка событий и обратных вызовов; не используйте этот планировщик для ввода-вывода (используйте Schedulers.io( ) вместо этого); количество потоков по умолчанию равно количеству процессоров |
Schedulers.from(executor) |
использует указанный Executor как планировщик (Scheduler) |
Schedulers.immediate( ) |
планирует выполнение работы немедленно в текущем потоке |
Schedulers.io( ) |
предназначен для операций ввода-вывода, таких как асинхронное выполнение блокирующих операций ввода-вывода; этот планировщик (Scheduler) поддерживается пулом потоков, который будет расти по мере необходимости; для обычных вычислений переключитесь на Schedulers.computation( ); Schedulers.io( ) по умолчанию является CachedThreadScheduler, что-то вроде планировщика новых потоков с кэшированием потоков |
Schedulers.newThread( ) |
создаёт новый поток для каждого блока работы |
Schedulers.trampoline( ) |
помещает работу в очередь для запуска в текущем потоке после выполнения всех уже находящихся в очереди |
Значения по умолчанию для планировщиков (Schedulers) операторов RxJava 1.x Observable
Некоторые операторы Observable в RxJava имеют альтернативные формы, которые позволяют вам задать планировщик (Scheduler), который будет использовать оператор (по крайней мере, для части своей работы). Другие операторы не работают с каким-либо конкретным планировщиком или работают с определённым планировщиком по умолчанию. Те, что используют конкретный планировщик по умолчанию, включают:
Использование планировщиков (Schedulers)
Помимо передачи этих планировщиков (Schedulers) операторам RxJava Observable, вы также можете использовать их для планирования собственной работы в подписках (Subscriptions). В следующем примере используется метод schedule класса Scheduler.Worker для планирования работы в планировщике newThread:
worker = Schedulers.newThread().createWorker();
worker.schedule(new Action0() {
@Override
public void call() {
yourWork();
}
});
// some time later...
worker.unsubscribe(); Рекурсивные планировщики (Schedulers)
Для планирования рекурсивных вызовов можно использовать schedule и затем schedule(this) в объекте Worker:
worker = Schedulers.newThread().createWorker();
worker.schedule(new Action0() {
@Override
public void call() {
yourWork();
// recurse until unsubscribed (schedule will do nothing if unsubscribed)
worker.schedule(this);
}
});
// some time later...
worker.unsubscribe(); Проверка или установка состояния отписки
Объекты класса Worker реализуют интерфейс Subscription с методами isUnsubscribed и unsubscribe, поэтому вы можете остановить работу при отмене подписки или отменить подписку внутри запланированной задачи:
Worker worker = Schedulers.newThread().createWorker();
Subscription mySubscription = worker.schedule(new Action0() {
@Override
public void call() {
while(!worker.isUnsubscribed()) {
status = yourWork();
if(QUIT == status) { worker.unsubscribe(); }
}
}
}); Класс Worker также является Subscription, поэтому вы можете (и должны, в конечном итоге) вызвать метод unsubscribe для сигнализации о прекращении работы и освобождения ресурсов:
worker.unsubscribe();
Отложенные и периодические планировщики (Schedulers)
Вы также можете использовать версию schedule, которая откладывает ваше действие в данном планировщике (Scheduler) до истечения определённого временного интервала. В следующем примере планируется выполнение someAction в планировщике someScheduler через 500 мс в соответствии с его часами:
someScheduler.schedule(someAction, 500, TimeUnit.MILLISECONDS);
Другой метод Scheduler позволяет планировать действие для выполнения через регулярные интервалы. В следующем примере планируется выполнение someAction в планировщике someScheduler через 500 мс, а затем каждые 250 мс после этого:
someScheduler.schedulePeriodically(someAction, 500, 250, TimeUnit.MILLISECONDS);
Планировщик тестирования (Test Scheduler)
Планировщик тестирования (TestScheduler) позволяет вам осуществлять тонкую ручную настройку поведения часов планировщика. Это может быть полезно для тестирования взаимодействий, которые зависят от точного расположения действий во времени. У этого планировщика есть три дополнительных метода:
advanceTimeTo(time,unit)- передвигает часы планировщика (Scheduler) к определённой точке времени
advanceTimeBy(time,unit)- передвигает часы планировщика (Scheduler) вперёд на определённое количество времени
triggerActions( )- запускает любые незапущенные действия, которые были запланированы на время, равное или раньше текущего времени по часам планировщика (Scheduler)
См. также
- Примеры работы с потоками в RxJava от Graham Lea
- Тестирование реактивных приложений от Ben Christensen
- Расширенное RxJava: Планировщики (Schedulers) (часть 1) (часть 2) (часть 3) (часть 4) от Dávid Karnok
RxJS
В RxJS вы получаете планировщики из объекта Rx.Scheduler или как независимо реализованные объекты. В следующей таблице показаны доступные в RxJS типы планировщиков:.
| Планировщик | назначение |
|---|---|
Rx.Scheduler.currentThread |
планирует работу как можно скорее в текущей потоковой нити |
Rx.HistoricalScheduler |
планирует работу так, как будто она происходит в произвольное историческое время |
Rx.Scheduler.immediate |
планирует работу немедленно в текущей потоковой нити |
Rx.TestScheduler |
для модульных тестов; это позволяет вам вручную управлять перемещением времени |
Rx.Scheduler.timeout |
планирует работу с помощью вызова обратного вызова с таймером |
См. также
- StackOverflow: Что такое «Планировщик» в RxJS
- Планировщики Денниса Стоянова
- Примеры работы с потоками RxJava Грэма Ли
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/scheduler.html