Spec-Zone.ru › ReactiveX

Создание

создать Observable с нуля с помощью функции
Create

Вы можете создать Observable с нуля, используя оператор Create. Вы передаёте этому оператору функцию, которая принимает наблюдателя в качестве параметра. Напишите эту функцию так, чтобы она вела себя как Observable — вызывая методы наблюдателя onNext, onError, и onCompleted соответствующим образом.

Хорошо сформированное конечное Observable должно попытаться вызвать метод наблюдателя onCompleted ровно один раз или метод onError ровно один раз и не должно затем пытаться вызывать другие методы наблюдателя.

См. также

  • Введение в Rx: Создание
  • Введение в Rx: Генерация
  • 101 Примеры Rx: Генерация
  • RxJava Tutorial 03: Observable from, just, & create methods

Информация, специфичная для языка

RxGroovy create

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

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

create

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.js
  • rx.all.js
  • rx.all.compat.js
  • rx.compat.js
  • rx.lite.js
  • rx.lite.compat.js
generate

Вы можете использовать оператор generate для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generate принимает четыре параметра:

  1. Первый элемент для выдачи
  2. Функция для проверки элемента, чтобы определить, следует ли его выдать (true) или завершить Observable (false)
  3. Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
  4. Функция для преобразования элементов перед их выдачей

Вы также можете передать в качестве необязательного пятого параметра Планировщик, который 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.js
  • rx.compat.js
  • rx.lite.js
  • rx.lite.compat.js
generateWithRelativeTime

Вы можете использовать оператор generateWithRelativeTime для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generateWithRelativeTime принимает пять параметров:

  1. Первый элемент для выдачи
  2. Функция для проверки элемента, чтобы определить, следует ли его выдать (true) или завершить Observable (false)
  3. Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
  4. Функция для преобразования элементов перед их выдачей
  5. Функция для указания того, как долго, в миллисекундах, генератор должен ждать после выдачи предыдущего элемента, прежде чем выдать этот элемент

Вы также можете передать в качестве необязательного шестого параметра Планировщик, который 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.js
  • rx.lite.compat.js
  • rx.time.js (требует rx.js или rx.compat.js)
generateWithAbsoluteTime

Вы можете использовать оператор generateWithAbsoluteTime для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generateWithAbsoluteTime принимает пять параметров:

  1. Первый элемент для выдачи
  2. Функция для проверки элемента, чтобы определить, следует ли его выдать (true) или завершить Observable (false)
  3. Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента
  4. Функция для преобразования элементов перед их выдачей
  5. Функция для указания того, в какое время (выраженное как 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

create

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

Вы можете использовать оператор generate для создания простых Observable, которые могут генерировать свои последующие выдачи и определять, когда нужно завершиться, на основе значения предыдущей выдачи. Основная форма generate принимает три параметра:

  1. Первый элемент для выдачи
  2. Функция для проверки элемента, чтобы определить, следует ли его выдать (true) или завершить Observable (false)
  3. Функция для генерации следующего элемента для проверки и выдачи на основе значения предыдущего элемента

Вы также можете передать в качестве необязательного четвёртого параметра Планировщик, который 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

Spec-Zone.ru

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