Начало
создайте Observable, который испускает возвращаемое значение директивы, похожей на функцию
Существует множество способов получения значений в результате вычислений в языках программирования, с названиями, такими как функции, фьючерсы, действия, вызываемые объекты, выполняемые задачи и так далее. Операторы, сгруппированные здесь под категорией оператора Начало, позволяют этим вещам вести себя как Observable, чтобы их можно было объединять с другими Observable в каскаде Observable
См. также
Информация, специфичная для языка
RxGroovy asyncAction asyncFunc deferFuture forEachFuture fromAction fromCallable fromFunc0 fromRunnable start startFuture toAsync
Различные реализации Начало в RxGroovy находятся в дополнительном модуле rxjava-async.
Модуль rxjava-async включает оператор start, который принимает функцию в качестве параметра, вызывает эту функцию для получения значения и затем возвращает Observable, который будет испускать это значение каждому последующему наблюдателю.
Обратите внимание, что функция будет выполнена только один раз, даже если к полученному Observable будут подписаны более одного наблюдателя.
Модуль rxjava-async также включает операторы toAsync, asyncAction, и asyncFunc. Они принимают функцию или Action в качестве параметра. В случае функции, этот вариант оператора вызывает эту функцию для получения значения, а затем возвращает Observable, который будет испускать это значение каждому последующему наблюдателю (точно так же, как оператор start).
В случае Action процесс аналогичен, но значения не возвращается. В этом случае Observable, созданный этим оператором, испустит null перед завершением.
Обратите внимание, что функция или Action будут выполнены только один раз, даже если к полученному Observable будут подписаны более одного наблюдателя.
Модуль rxjava-async также включает оператор startFuture. Вы передаете ему функцию, которая возвращает Future. startFuture немедленно вызывает эту функцию для получения Future, и вызывает метод Future’s get для попытки получить его значение. Он возвращает Observable, которому он будет испускать это значение любым последующим наблюдателям.
Модуль rxjava-async также включает оператор deferFuture. Вы передаете ему функцию, которая возвращает Future, которая возвращает Observable. deferFuture возвращает Observable, но не вызывает переданную вами функцию до тех пор, пока наблюдатель не подпишется на возвращаемое Observable. Когда это происходит, он немедленно вызывает get на получившемся Future, а затем отражает испускания от Observable, возвращаемого Future, как свои собственные испускания.
Таким образом, вы можете включить Future, которая возвращает Observable, в каскад Observable как равную другим Observable.
Модуль rxjava-async также включает оператор fromAction. Он принимает Action в качестве параметра и возвращает Observable, который испускает элемент, который вы передаете в fromAction по завершении Action.
Модуль rxjava-async также включает оператор fromCallable. Он принимает Callable в качестве параметра и возвращает Observable, который испускает результат этого вызываемого объекта в качестве единственного испускания.
Модуль rxjava-async также включает оператор fromRunnable. Он принимает Runnable в качестве параметра и возвращает Observable, который испускает элемент, который вы передаете в fromRunnable по завершении Runnable.
Модуль rxjava-async также включает оператор forEachFuture. Это не совсем вариант оператора Начало, а нечто собственное. Вы передаете forEachFuture подмножество типичных методов наблюдателя (onNext, onError, и onCompleted) и Observable будет вызывать эти методы обычным способом. Но forEachFuture сам возвращает Future, который блокирует выполнение get до завершения исходного Observable, затем возвращает завершение или ошибку, в зависимости от того, как завершился Observable.
Вы можете использовать это, если вам нужна функция, которая блокирует выполнение до завершения Observable.
Модуль rxjava-async также включает оператор runAsync. Он интересен тем, что создает специализацию Observable, называемую StoppableObservable.
Передайте runAsync Action и Scheduler, и он вернёт StoppableObservable, который использует указанный Action для генерации испускаемых элементов. Action принимает Observer и Subscription. Он использует Subscription для проверки условия unsubscribed, по достижении которого он перестанет испускать элементы. Вы также можете вручную остановить StoppableObservable в любое время, вызвав его метод unsubscribe (который также отменит подписку Subscription, которую вы связали с StoppableObservable).
Поскольку runAsync немедленно вызывает Action и начинает испускать элементы (то есть, он производит активный Observable), возможно, что некоторые элементы могут быть утеряны в интервале между установлением подписки StoppableObservable с этим оператором и готовностью вашего Observer к получению элементов. Если это проблема, вы можете использовать вариант runAsync, который также принимает Subject, и передать ReplaySubject, с помощью которого можно получить утерянные элементы.
В RxGroovy также есть версия оператора From, которая преобразует Future в Observable, и таким образом напоминает оператор Начало.
RxJava 1․x asyncAction asyncFunc deferFuture forEachFuture fromAction fromCallable fromFunc0 fromRunnable start startFuture toAsync
Различные реализации RxJava оператора Start находятся в дополнительном rxjava-async модуле.
Модуль rxjava-async включает оператор start, который принимает функцию в качестве параметра, вызывает эту функцию для получения значения и возвращает Observable, которое будет отправлять это значение каждому последующему наблюдателю.
Обратите внимание, что функция будет выполнена только один раз, даже если к результату Observable подпишутся более одного наблюдателя.
Модуль rxjava-async также включает операторы toAsync, asyncAction, и asyncFunc. Они принимают функцию или Action в качестве параметра. В случае функции, этот вариант оператора вызывает функцию для получения значения и возвращает Observable, которое будет отправлять это значение каждому последующему наблюдателю (точно так же, как и оператор start).
В случае Action процесс аналогичен, но возвращаемого значения нет. В этом случае Observable, созданный этим оператором, отправит null перед завершением.
Обратите внимание, что функция или Action будут выполнены только один раз, даже если к результату Observable подпишутся более одного наблюдателя.
Модуль rxjava-async также включает оператор startFuture. Вы передаёте ему функцию, которая возвращает Future. startFuture сразу же вызывает эту функцию, чтобы получить Future, и вызывает метод Future’s get для попытки получить его значение. Он возвращает Observable, которому он будет отправлять это значение любым последующим наблюдателям.
Модуль rxjava-async также включает оператор deferFuture. Вы передаёте ему функцию, которая возвращает Future, который возвращает Observable. deferFuture возвращает Observable, но не вызывает переданную функцию до тех пор, пока наблюдатель не подпишется на возвращённое Observable. Когда это происходит, он немедленно вызывает get на полученном Future, а затем дублирует отправки из Observable, возвращённого Future, как свои собственные отправки.
Таким образом, вы можете включить Future, который возвращает Observable, в каскаде Observable как равный другим Observable.
Модуль rxjava-async также включает оператор fromAction. Он принимает Action в качестве параметра и возвращает Observable, который отправляет элемент, который вы передаёте fromAction при завершении Action
Модуль rxjava-async также включает оператор fromCallable. Он принимает Callable в качестве параметра и возвращает Observable, который отправляет результат этого вызываемого объекта в качестве единственной отправки.
Модуль rxjava-async также включает оператор fromRunnable. Он принимает Runnable в качестве параметра и возвращает Observable, который отправляет элемент, который вы передаёте fromRunnable при завершении Runnable
Модуль rxjava-async также включает оператор forEachFuture. Он не является вариантом оператора Start, а представляет собой нечто собственное. Вы передаёте forEachFuture некоторые подмножества типичных методов наблюдателя (onNext, onError, и onCompleted) и Observable будет вызывать эти методы обычным способом. Но сам forEachFuture возвращает Future, который блокируется на get до завершения исходного Observable, затем возвращает завершение или ошибку, в зависимости от того, как Observable завершился.
Вы можете использовать это, если вам нужна функция, которая блокируется до завершения Observable.
Модуль rxjava-async также включает оператор runAsync. Он отличается тем, что создаёт специализацию Observable, называемую StoppableObservable.
Передайте runAsync Action и Scheduler, и он вернёт StoppableObservable, который использует указанный Action для генерации отправляемых элементов. Action принимает Observer и Subscription. Он использует Subscription для проверки условия unsubscribed, по достижении которого он прекратит отправку элементов. Вы также можете вручную остановить StoppableObservable в любое время, вызвав его метод unsubscribe (который также отпишет Subscription , который вы связали с StoppableObservable) .
Поскольку runAsync немедленно вызывает Action и начинает отправлять элементы (то есть, он производит горячее Observable), возможно, что некоторые элементы могут быть потеряны в интервале между установкой StoppableObservable с этим оператором и готовностью вашего Observer принять элементы. Если это проблема, вы можете использовать вариант runAsync, который также принимает Subject и передать ReplaySubject, с помощью которого вы можете получить недостающие элементы.
В RxJava также существует версия оператора From, которая преобразует Future в Observable, и таким образом напоминает оператор Start.
RxJS start startAsync toAsync
RxJS реализует оператор start. Он принимает в качестве параметров функцию, значение которой будет отправкой от получившегося Observable, и, по желанию, любые дополнительные параметры для этой функции и Scheduler для запуска функции.
Пример кода
var context = { value: 42 };
var source = Rx.Observable.start(
function () {
return this.value;
},
context,
Rx.Scheduler.timeout
);
var subscription = source.subscribe(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}); Next: 42 Completed
start доступен в следующих дистрибутивах:
-
rx.async.js(требуетrx.binding.jsи либоrx.jsилиrx.compat.js) -
rx.async.compat.js(требуетrx.binding.jsи либоrx.jsилиrx.compat.js) rx.lite.jsrx.lite.compat.js
RxJS также реализует оператор startAsync. Он принимает в качестве параметров асинхронную функцию, значение которой будет отправкой от получившегося Observable.
Вы можете преобразовать функцию в асинхронную с помощью метода toAsync. Он принимает функцию, параметр функции и Scheduler в качестве параметров и возвращает асинхронную функцию, которая будет вызвана на указанном Scheduler. Последние два параметра необязательны; если вы не укажете Scheduler, по умолчанию будет использован timeout Scheduler.
Пример кода
var source = Rx.Observable.startAsync(function () {
return RSVP.Promise.resolve(42);
});
var subscription = source.subscribe(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}); Next: 42 Completed
startAsync доступен в следующих дистрибутивах:
-
rx.async.js(требуетrx.binding.jsи либоrx.jsилиrx.compat.js) -
rx.async.compat.js(требуетrx.binding.jsи либоrx.jsилиrx.compat.js) rx.lite.jsrx.lite.compat.js
toAsync доступен в следующих дистрибутивах:
-
rx.async.js(требуетrx.binding.jsи либоrx.jsилиrx.compat.js) -
rx.async.compat.js(требуетrx.binding.jsи либоrx.jsилиrx.compat.js)
RxPHP start
RxPHP реализует этот оператор как start.
Асинхронно вызывает указанную функцию на указанном планировщике, отображая результат через последовательность observable.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/start/start.php
$source = Rx\Observable::start(function () {
return 42;
});
$source->subscribe($stdoutObserver); Next value: 42 Complete!
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/start.html