Spec-Zone.ru › ReactiveX

Планировщик

Если вы хотите добавить многопоточность в вашу цепочку операторов Observable, вы можете это сделать, указав тем операторам (или конкретным Observable) работать на определенном Планировщике.

Некоторые операторы ReactiveX Observable имеют варианты, которые принимают планировщик в качестве параметра. Эти варианты инструктируют оператор выполнить часть или всю свою работу на конкретном планировщике.

По умолчанию Observable и цепочка операторов, которые к нему применяются, выполняют свою работу и уведомляют своих наблюдателей на той же нити, на которой вызывается его Subscribe метод. Оператор SubscribeOn изменяет это поведение, указывая другой планировщик, на котором Observable должен работать. Оператор ObserveOn указывает другой планировщик, который Observable будет использовать для отправки уведомлений своим наблюдателям.

Как показано на этой иллюстрации, оператор SubscribeOn определяет, на какой нити Observable начнет работать, независимо от того, в какой точке цепочки операторов этот оператор вызван. ObserveOn, с другой стороны, влияет на нить, которую Observable будет использовать ниже, где этот оператор появляется. По этой причине вы можете вызывать ObserveOn несколько раз в различных точках цепочки операторов Observable, чтобы изменить, на каких нитях работают определенные из этих операторов.

ObserveOn and SubscribeOn

См. также

  • Введение в 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 имеют альтернативные формы, которые позволяют указать, какой планировщик будет использовать оператор для (по крайней мере, части) своей работы. Другие не работают с каким-либо конкретным планировщиком или работают с конкретным планировщиком по умолчанию. К тем, которые используют конкретный планировщик по умолчанию, относятся:

Оператор Планировщик
buffer(timespan) computation
buffer(timespan, count) computation
buffer(timespan, timeshift) computation
debounce(timeout, unit) computation
delay(delay, unit) computation
delaySubscription(delay, unit) computation
interval computation
repeat trampoline
replay(time, unit) computation
replay(buffersize, time, unit) computation
replay(selector, time, unit) computation
replay(selector, buffersize, time, unit) computation
retry trampoline
sample(period, unit) computation
skip(time, unit) computation
skipLast(time, unit) computation
take(time, unit) computation
takeLast(time, unit) computation
takeLast(count, time, unit) computation
takeLastBuffer(time, unit) computation
takeLastBuffer(count, time, unit) computation
throttleFirst computation
throttleLast computation
throttleWithTimeout computation
timeInterval immediate
timeout(timeoutSelector) immediate
timeout(firstTimeoutSelector, timeoutSelector) immediate
timeout(timeoutSelector, other) immediate
timeout(timeout, timeUnit) computation
timeout(firstTimeoutSelector, timeoutSelector, other) immediate
timeout(timeout, timeUnit, other) computation
timer computation
timestamp immediate
window(timespan) computation
window(timespan, count) computation
window(timespan, timeshift) computation

Тестовый планировщик

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), который будет использовать оператор (по крайней мере, для части своей работы). Другие операторы не работают с каким-либо конкретным планировщиком или работают с определённым планировщиком по умолчанию. Те, что используют конкретный планировщик по умолчанию, включают:

оператор планировщик
buffer(timespan) computation
buffer(timespan, count) computation
buffer(timespan, timeshift) computation
debounce(timeout, unit) computation
delay(delay, unit) computation
delaySubscription(delay, unit) computation
interval computation
repeat trampoline
replay(time, unit) computation
replay(buffersize, time, unit) computation
replay(selector, time, unit) computation
replay(selector, buffersize, time, unit) computation
retry trampoline
sample(period, unit) computation
skip(time, unit) computation
skipLast(time, unit) computation
take(time, unit) computation
takeLast(time, unit) computation
takeLast(count, time, unit) computation
takeLastBuffer(time, unit) computation
takeLastBuffer(count, time, unit) computation
throttleFirst computation
throttleLast computation
throttleWithTimeout computation
timeInterval immediate
timeout(timeoutSelector) immediate
timeout(firstTimeoutSelector, timeoutSelector) immediate
timeout(timeoutSelector, other) immediate
timeout(timeout, timeUnit) computation
timeout(firstTimeoutSelector, timeoutSelector, other) immediate
timeout(timeout, timeUnit, other) computation
timer computation
timestamp immediate
window(timespan) computation
window(timespan, count) computation
window(timespan, timeshift) computation

Использование планировщиков (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

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API