Создание
создать Observable с нуля с помощью функции
Вы можете создать Observable с нуля, используя оператор Create. Вы передаёте этому оператору функцию, которая принимает наблюдателя в качестве параметра. Напишите эту функцию так, чтобы она вела себя как Observable — вызывая методы наблюдателя onNext, onError, и onCompleted соответствующим образом.
Хорошо сформированное конечное Observable должно попытаться вызвать метод наблюдателя onCompleted ровно один раз или метод onError ровно один раз и не должно затем пытаться вызывать другие методы наблюдателя.
См. также
- Введение в Rx: Создание
- Введение в Rx: Генерация
- 101 Примеры Rx: Генерация
- RxJava Tutorial 03: Observable from, just, & create methods
Информация, специфичная для языка
RxGroovy create
RxGroovy реализует этот оператор как create.
Пример кода
def myObservable = Observable.create({ aSubscriber ->
try {
for (int i = 1; i < 1000000; i++) {
if (aSubscriber.isUnsubscribed()) {
return;
}
aSubscriber.onNext(i);
}
if (!aSubscriber.isUnsubscribed()) {
aSubscriber.onCompleted();
}
} catch(Throwable t) {
if (!aSubscriber.isUnsubscribed()) {
aSubscriber.onError(t);
}
}
}) Хорошей практикой является проверка состояния наблюдателя isUnsubscribed, чтобы ваше Observable могло прекратить отправку элементов или выполнение дорогостоящих вычислений, когда больше нет заинтересованных наблюдателей.
create по умолчанию не работает с каким-либо конкретным Планировщиком.
- Javadoc:
create(OnSubscribe)
RxJava 1․x create
RxJava реализует этот оператор как create.
Хорошей практикой является проверка состояния наблюдателя isUnsubscribed в функции, которую вы передаёте в create, чтобы ваше Observable могло прекратить отправку элементов или выполнение дорогостоящих вычислений, когда больше нет заинтересованных наблюдателей.
Пример кода
Observable.create(new Observable.OnSubscribe<Integer>() {
@Override
public void call(Subscriber<? super Integer> observer) {
try {
if (!observer.isUnsubscribed()) {
for (int i = 1; i < 5; i++) {
observer.onNext(i);
}
observer.onCompleted();
}
} catch (Exception e) {
observer.onError(e);
}
}
} ).subscribe(new Subscriber<Integer>() {
@Override
public void onNext(Integer item) {
System.out.println("Next: " + item);
}
@Override
public void onError(Throwable error) {
System.err.println("Error: " + error.getMessage());
}
@Override
public void onCompleted() {
System.out.println("Sequence complete.");
}
}); Next: 1 Next: 2 Next: 3 Next: 4 Sequence complete.
create по умолчанию не работает с каким-либо конкретным Планировщиком.
- Javadoc:
create(OnSubscribe)
RxJS create createWithDisposable generate generateWithAbsoluteTime generateWithRelativeTime
RxJS реализует этот оператор как create (также есть альтернативное название для того же оператора: createWithDisposable).
Пример кода
/* Using a function */
var source = Rx.Observable.create(function (observer) {
observer.onNext(42);
observer.onCompleted();
// Note that this is optional, you do not have to return this if you require no cleanup
return function () { console.log('disposed'); };
});
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
Next: 42 Completed
/* Using a disposable */
var source = Rx.Observable.create(function (observer) {
observer.onNext(42);
observer.onCompleted();
// Note that this is optional, you do not have to return this if you require no cleanup
return Rx.Disposable.create(function () {
console.log('disposed');
});
});
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
Next: 42 Completed
create доступен в следующих дистрибутивах:
rx.jsrx.all.jsrx.all.compat.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
Вы можете использовать оператор generate для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generate принимает четыре параметра:
- Первый элемент для выдачи
- Функция для проверки элемента, чтобы определить, следует ли его выдать (
true) или завершить Observable (false) - Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
- Функция для преобразования элементов перед их выдачей
Вы также можете передать в качестве необязательного пятого параметра Планировщик, который generate будет использовать для создания и выдачи своей последовательности (он использует currentThread по умолчанию).
Пример кода
var source = Rx.Observable.generate(
0,
function (x) { return x < 3; },
function (x) { return x + 1; },
function (x) { return x; }
);
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: 0 Next: 1 Next: 2 Completed
generate доступен в следующих дистрибутивах:
rx.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
Вы можете использовать оператор generateWithRelativeTime для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generateWithRelativeTime принимает пять параметров:
- Первый элемент для выдачи
- Функция для проверки элемента, чтобы определить, следует ли его выдать (
true) или завершить Observable (false) - Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
- Функция для преобразования элементов перед их выдачей
- Функция для указания того, как долго, в миллисекундах, генератор должен ждать после выдачи предыдущего элемента, прежде чем выдать этот элемент
Вы также можете передать в качестве необязательного шестого параметра Планировщик, который generate будет использовать для создания и выдачи своей последовательности (он использует currentThread по умолчанию).
Пример кода
var source = Rx.Observable.generateWithRelativeTime(
1,
function (x) { return x < 4; },
function (x) { return x + 1; },
function (x) { return x; },
function (x) { return 100 * x; }
).timeInterval();
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: {value: 1, interval: 100}
Next: {value: 2, interval: 200}
Next: {value: 3, interval: 300}
Completed generateWithRelativeTime доступен в следующих дистрибутивах:
rx.lite.jsrx.lite.compat.js-
rx.time.js(требуетrx.jsилиrx.compat.js)
Вы можете использовать оператор generateWithAbsoluteTime для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generateWithAbsoluteTime принимает пять параметров:
- Первый элемент для выдачи
- Функция для проверки элемента, чтобы определить, следует ли его выдать (
true) или завершить Observable (false) - Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
- Функция для преобразования элементов перед их выдачей
- Функция для указания того, в какое время (выраженное как
Date) генератор должен выдать новый элемент
Вы также можете передать в качестве необязательного шестого параметра Планировщик, который generate будет использовать для создания и выдачи своей последовательности (он использует currentThread по умолчанию).
Пример кода
var source = Rx.Observable.generate(
1,
function (x) { return x < 4; },
function (x) { return x + 1; },
function (x) { return x; },
function (x) { return Date.now() + (100 * x); }
).timeInterval();
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: {value: 1, interval: 100}
Next: {value: 2, interval: 200}
Next: {value: 3, interval: 300}
Completed generateWithAbsoluteTime доступен в следующих дистрибутивах:
-
rx.time.js(требуетrx.jsилиrx.compat.js)
RxPHP create
RxPHP реализует этот оператор как create.
Создаёт последовательность observable из заданной подпискиAction callable реализации.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/create/create.php
//With static method
$source = \Rx\Observable::create(function (\Rx\ObserverInterface $observer) {
$observer->onNext(42);
$observer->onCompleted();
return new CallbackDisposable(function () {
echo "Disposed\n";
});
});
$subscription = $source->subscribe($createStdoutObserver()); Next value: 42 Complete! Disposed
RxSwift create generate
RxSwift реализует этот оператор как create.
Пример кода
let source : Observable = Observable.create { observer in
for i in 1...5 {
observer.on(.next(i))
}
observer.on(.completed)
// Note that this is optional. If you require no cleanup you can return
// `Disposables.create()` (which returns the `NopDisposable` singleton)
return Disposables.create {
print("disposed")
}
}
source.subscribe {
print($0)
} next(1) next(2) next(3) next(4) next(5) completed disposed
Вы можете использовать оператор generate для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generate принимает три параметра:
- Первый элемент для выдачи
- Функция для проверки элемента, чтобы определить, следует ли его выдать (
true) или завершить Observable (false) - Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
Вы также можете передать в качестве необязательного четвёртого параметра Планировщик, который generate будет использовать для создания и выдачи своей последовательности (он использует CurrentThreadScheduler по умолчанию).
Пример кода
let source = Observable.generate(
initialState: 0,
condition: { $0 < 3 },
iterate: { $0 + 1 }
)
source.subscribe {
print($0)
} next(0) next(1) next(2) completed
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/create.html