Spec-Zone.ru › ReactiveX

Из

преобразуйте различные другие объекты и типы данных в Observable
From

Когда вы работаете с Observable, может быть удобнее, если все данные, с которыми вы хотите работать, могут быть представлены как Observable, а не как смесь Observable и других типов. Это позволяет использовать один набор операторов для управления всем жизненным циклом потока данных.

Например, Iterable можно рассматривать как своего рода синхронный Observable; Future — как своего рода Observable, который всегда излучает только один элемент. Явно преобразуя такие объекты в Observable, вы позволяете им взаимодействовать как равные с другими Observable.

По этой причине большинство реализаций ReactiveX имеют методы, которые позволяют преобразовать определенные объекты и структуры данных, специфичные для языка программирования, в Observable.

См. также

  • Just
  • Start
  • 101 Rx Примеры: Операторы наблюдения
  • RxJava Урок 03: Observable from, just, & create методы

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

RxGroovy decode from fromAction fromCallable fromFunc0 fromRunnable runAsync

from

В RxGroovy оператор from может преобразовать Future, Iterable или Array. В случае Iterable или Array результирующий Observable будет излучать каждый элемент, содержащийся в Iterable или Array.

В случае Future он будет излучать единственный результат вызова get. Вы можете необязательно передать вариант from, который принимает Future, два дополнительных параметра, указывающих на временной интервал ожидания и единицы времени, в которых этот интервал измеряется. Результирующий Observable завершится ошибкой, если этот временной интервал истечет до того, как Future ответит значением.

from по умолчанию не работает с конкретным Scheduler, однако вы можете передать вариант преобразования Future с Scheduler в качестве необязательного второго параметра, и он будет использовать этот Scheduler для управления Future.

  • Javadoc: from(array)
  • Javadoc: from(Iterable)
  • Javadoc: from(Future)
  • Javadoc: from(Future,Scheduler)
  • Javadoc: from(Future,timout,timeUnit)
fromFunc0

Кроме того, в пакете RxJavaAsyncUtil вам доступны следующие операторы, которые преобразуют действия, вызываемые объекты, функции и Runnable в Observable, излучающие результаты этих вещей:

  • fromAction
  • fromCallable
  • fromFunc0
  • fromRunnable

См. оператор Start для получения дополнительной информации об этих операторах.

from

Обратите внимание, что также существует оператор from, который является методом необязательного класса StringObservable. Он преобразует поток символов или Reader в Observable, который излучает массивы байтов или строки.

В отдельном пакете RxJavaAsyncUtil, который не включен по умолчанию в RxGroovy, также существует функция runAsync. Передайте runAsync Action и Scheduler, и она вернёт StoppableObservable, которая использует указанный Action для генерации излучаемых элементов.

Action принимает Observer и Subscription. Она использует Subscription для проверки условия isUnsubscribed, после чего прекратит излучение элементов. Вы также можете вручную остановить StoppableObservable в любое время, вызвав метод unsubscribe (который также отменит подписку Subscription, связанную с StoppableObservable).

Поскольку runAsync немедленно вызывает Action и начинает излучать элементы, возможно, что некоторые элементы могут быть потеряны в промежутке между установлением StoppableObservable с помощью этого метода и готовностью вашего Observer к получению элементов. Если это проблема, вы можете использовать вариант runAsync, который также принимает Subject, и передать ReplaySubject для извлечения пропущенных элементов.

decode

Класс StringObservable, который не является частью RxGroovy по умолчанию, также включает оператор decode, который преобразует поток многобайтовых символов в Observable, излучающий массивы байтов, учитывая границы символов.

RxJava 1․x decode from fromAction fromCallable fromFunc0 fromRunnable runAsync

from

В RxJava оператор from может преобразовать Future, Iterable или Array. В случае Iterable или Array результирующий Observable будет излучать каждый элемент, содержащийся в Iterable или Array.

Пример кода

Integer[] items = { 0, 1, 2, 3, 4, 5 };
Observable myObservable = Observable.from(items);

myObservable.subscribe(
    new Action1<Integer>() {
        @Override
        public void call(Integer item) {
            System.out.println(item);
        }
    },
    new Action1<Throwable>() {
        @Override
        public void call(Throwable error) {
            System.out.println("Error encountered: " + error.getMessage());
        }
    },
    new Action0() {
        @Override
        public void call() {
            System.out.println("Sequence complete");
        }
    }
);
0
1
2
3
4
5
Sequence complete

В случае Future он будет излучать единственный результат вызова get. Вы можете необязательно передать вариант from, который принимает Future, два дополнительных параметра, указывающих на временной интервал ожидания и единицы времени, в которых этот интервал измеряется. Результирующий Observable завершится ошибкой, если этот временной интервал истечет до того, как Future ответит значением.

from по умолчанию не работает с конкретным Scheduler, однако вы можете передать вариант преобразования Future с Scheduler в качестве необязательного второго параметра, и он будет использовать этот Scheduler для управления Future.

  • Javadoc: from(array)
  • Javadoc: from(Iterable)
  • Javadoc: from(Future)
  • Javadoc: from(Future,Scheduler)
  • Javadoc: from(Future,timout,timeUnit)
fromFunc0

Кроме того, в пакете RxJavaAsyncUtil вам доступны следующие операторы, которые преобразуют действия, вызываемые объекты, функции и Runnable в Observable, излучающие результаты этих вещей:

  • fromAction
  • fromCallable
  • fromFunc0
  • fromRunnable

См. оператор Start для получения дополнительной информации об этих операторах.

from

Обратите внимание, что также существует оператор from, который является методом необязательного класса StringObservable. Он преобразует поток символов или Reader в Observable, который излучает массивы байтов или строки.

В отдельном пакете RxJavaAsyncUtil, который не включен по умолчанию в RxJava, также существует функция runAsync. Передайте runAsync Action и Scheduler, и она вернёт StoppableObservable, которая использует указанный Action для генерации излучаемых элементов.

Action принимает Observer и Subscription. Она использует Subscription для проверки условия isUnsubscribed, после чего прекратит излучение элементов. Вы также можете вручную остановить StoppableObservable в любое время, вызвав метод unsubscribe (который также отменит подписку Subscription, связанную с StoppableObservable).

Поскольку runAsync немедленно вызывает Action и начинает излучать элементы, возможно, что некоторые элементы могут быть потеряны в промежутке между установлением StoppableObservable с помощью этого метода и готовностью вашего Observer к получению элементов. Если это проблема, вы можете использовать вариант runAsync, который также принимает Subject, и передать ReplaySubject для извлечения пропущенных элементов.

decode

Класс StringObservable, который не является частью RxGroovy по умолчанию, также включает оператор decode, который преобразует поток многобайтовых символов в Observable, излучающий массивы байтов, учитывая границы символов.

RxJS from fromCallback fromEvent fromEventPattern fromNodeCallback fromPromise of ofArrayChanges ofObjectChanges ofWithScheduler pairs

Существует несколько специализированных вариантов оператора From в RxJS:

from

В RxJS оператор from преобразует массив или итерируемый объект в Observable, которое испускает элементы этого массива или итерируемого объекта. Строка в данном контексте рассматривается как массив символов.

Этот оператор также принимает три дополнительных необязательных параметра:

  1. функцию преобразования, которая принимает элемент массива или итерируемого объекта в качестве входных данных и генерирует элемент, который будет испущен результирующим Observable в качестве выходных данных
  2. второй аргумент, который передается в функцию преобразования в качестве дополнительной контекстной информации
  3. Scheduler, на котором этот оператор должен работать

Пример кода

// Array-like object (arguments) to Observable
function f() {
  return Rx.Observable.from(arguments);
}

f(1, 2, 3).subscribe(
  function (x) { console.log('Next: ' + x); },
  function (err) { console.log('Error: ' + err); },
  function () { console.log('Completed'); });
Next: 1
Next: 2
Next: 3
Completed
// Any iterable object...
// Set
var s = new Set(['foo', window]);
Rx.Observable.from(s).subscribe(
  function (x) { console.log('Next: ' + x); },
  function (err) { console.log('Error: ' + err); },
  function () { console.log('Completed'); });
Next: foo
Next: window
Completed
// Map
var m = new Map([[1, 2], [2, 4], [4, 8]]);
Rx.Observable.from(m).subscribe(
  function (x) { console.log('Next: ' + x); },
  function (err) { console.log('Error: ' + err); },
  function () { console.log('Completed'); });
Next: [1, 2]
Next: [2, 4]
Next: [4, 8]
Completed
// String
Rx.Observable.from("foo").subscribe(
  function (x) { console.log('Next: ' + x); },
  function (err) { console.log('Error: ' + err); },
  function () { console.log('Completed'); });
Next: f
Next: o
Next: o
Completed
// Using an arrow function as the map function to manipulate the elements
Rx.Observable.from([1, 2, 3], function (x) { return x + x; }).subscribe(
  function (x) { console.log('Next: ' + x); },
  function (err) { console.log('Error: ' + err); },
  function () { console.log('Completed'); });
Next: 2
Next: 4
Next: 6
Completed
// Generate a sequence of numbers
Rx.Observable.from({length: 5}, function(v, k) { return k; }).subscribe(
  function (x) { console.log('Next: ' + x); },
  function (err) { console.log('Error: ' + err); },
  function () { console.log('Completed'); });
Next: 0
Next: 1
Next: 2
Next: 3
Next: 4
Completed

Оператор from доступен в следующих дистрибутивах:

  • rx.js
  • rx.all.js
  • rx.all.compat.js
  • rx.compat.js
  • rx.lite.js
  • rx.lite.compat.js
fromCallback

Оператор fromCallback принимает функцию в качестве параметра, вызывает эту функцию и испускает возвращаемое ею значение в качестве единственного испускания.

Этот оператор также принимает два дополнительных необязательных параметра:

  1. параметр для передачи в функцию обратного вызова
  2. функция преобразования, которая принимает возвращаемое значение функции обратного вызова в качестве входных данных и возвращает элемент, который будет испущен результирующим Observable

Пример кода

var fs = require('fs'),
    Rx = require('rx');

// Wrap fs.exists
var exists = Rx.Observable.fromCallback(fs.exists);

// Check if file.txt exists
var source = exists('file.txt');

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (err) { console.log('Error: ' + err); },
    function () { console.log('Completed'); });
Next: true
Completed

Оператор fromCallback доступен в следующих дистрибутивах:

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

Также есть оператор fromNodeCallback, специализированный для типов функций обратного вызова, используемых в Node.js.

Этот оператор принимает три дополнительных необязательных параметра:

  1. Scheduler, на котором вы хотите запустить обратный вызов Node.js
  2. параметр, который нужно передать функции обратного вызова
  3. функция преобразования, которая принимает возвращаемое значение функции обратного вызова в качестве входных данных и возвращает элемент, который будет испущен результирующим Observable

Пример кода

var fs = require('fs'),
    Rx = require('rx');

// Wrap fs.exists
var rename = Rx.Observable.fromNodeCallback(fs.rename);

// Rename file which returns no parameters except an error
var source = rename('file1.txt', 'file2.txt');

var subscription = source.subscribe(
    function () { console.log('Next: success!'); },
    function (err) { console.log('Error: ' + err); },
    function () { console.log('Completed'); });
Next: success!
Completed

Оператор fromNodeCallback доступен в следующих дистрибутивах:

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

Оператор fromEvent принимает «элемент» и имя события в качестве параметров и затем прослушивает события с этим именем, происходящие на этом элементе. Он возвращает Observable, который испускает эти события. «Элемент» может быть простым элементом DOM, или NodeList, элементом jQuery, Zepto, Angular, Ember.js или EventEmitter.

Этот оператор также принимает необязательный третий параметр: функцию, которая принимает аргументы обработчика событий в качестве параметров и возвращает элемент, который будет испущен результирующим Observable вместо события.

Пример кода

// using a jQuery element
var input = $('#input');

var source = Rx.Observable.fromEvent(input, 'click');

var subscription = source.subscribe(
    function (x) { console.log('Next: Clicked!'); },
    function (err) { console.log('Error: ' + err); },
    function () { console.log('Completed'); });

input.trigger('click');
Next: Clicked!
// using a Node.js EventEmitter and the optional third parameter
var EventEmitter = require('events').EventEmitter,
    Rx = require('rx');

var eventEmitter = new EventEmitter();

var source = Rx.Observable.fromEvent(
    eventEmitter,
    'data',
    function (first, second) {
        return { foo: first, bar: second };
    });

var subscription = source.subscribe(
    function (x) {
        console.log('Next: foo -' + x.foo + ', bar -' + x.bar);
    },
    function (err) { console.log('Error: ' + err); },
    function () { console.log('Completed'); });

eventEmitter.emit('data', 'baz', 'quux');
Next: foo - baz, bar - quux

Оператор fromEvent доступен в следующих дистрибутивах:

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

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

Пример кода

var input = $('#input');

var source = Rx.Observable.fromEventPattern(
    function add (h) {
        input.bind('click', h);
    },
    function remove (h) {
        input.unbind('click', h);
    }
);

var subscription = source.subscribe(
    function (x) { console.log('Next: Clicked!'); },
    function (err) { console.log('Error: ' + err); },
    function () { console.log('Completed'); });

input.trigger('click');
Next: Clicked!
of

Оператор of принимает ряд элементов в качестве параметров и возвращает Observable, которое испускает каждый из этих параметров по порядку как последовательность испускаемых значений.

Пример кода

var source = Rx.Observable.of(1,2,3);

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (err) { console.log('Error: ' + err); },
    function () { console.log('Completed'); });
Next: 1
Next: 2
Next: 3
Completed

Оператор of доступен в следующих дистрибутивах:

  • rx.js
  • rx.all.js
  • rx.all.compat.js
  • rx.compat.js
  • rx.lite.js
  • rx.lite.compat.js

Вариант этого оператора, ofWithScheduler, принимает Scheduler в качестве первого параметра и выполняет результирующий Observable на этом Scheduler.

Также существует оператор fromPromise, который преобразует Promise в Observable, преобразуя вызовы `then` в испускания `next`, и вызовы `catch` в испускания `error`.

Оператор fromPromise доступен в следующих дистрибутивах:

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

Пример кода

var promise = new RSVP.Promise(function (resolve, reject) {
   resolve(42);
});

var source = Rx.Observable.fromPromise(promise);

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (e) { console.log('Error: ' + e); },
    function ( ) { console.log('Completed'); });
Next: 42:
Completed
var promise = new RSVP.Promise(function (resolve, reject) {
   reject(new Error('reason'));
});

var source = Rx.Observable.fromPromise(promise);

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (e) { console.log('Error: ' + e); },
    function ( ) { console.log('Completed'); });
Error: Error: reject

Также есть оператор ofArrayChanges который отслеживает массив с помощью метода Array.observe, и возвращает Observable, который испускает любые изменения, которые происходят в массиве. Этот оператор доступен только в дистрибутиве rx.all.js.

Пример кода

var arr = [1,2,3];
var source = Rx.Observable.ofArrayChanges(arr);

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (e) { console.log('Error: ' + e); },
    function ( ) { console.log('Completed'); });

arr.push(4)
Next: {type: "splice", object: Array[4], index: 3, removed: Array[0], addedCount: 1}

Аналогичный оператор — ofObjectChanges. Он возвращает Observable, который испускает любые изменения, внесенные в определенный объект, как сообщается методом Object.observe. Он также доступен только в дистрибутиве rx.all.js.

Пример кода

var obj = {x: 1};
var source = Rx.Observable.ofObjectChanges(obj);

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (e) { console.log('Error: ' + e); },
    function ( ) { console.log('Completed'); });

obj.x = 42;
Next: {type: "update", object: Object, name: "x", oldValue: 1}

Также есть оператор pairs. Этот оператор принимает объект и возвращает Observable, который испускает пары ключ/значение, атрибуты этого объекта.

Пример кода

var obj = {
  foo: 42,
  bar: 56,
  baz: 78
};

var source = Rx.Observable.pairs(obj);

var subscription = source.subscribe(
    function (x) { console.log('Next: ' + x); },
    function (e) { console.log('Error: ' + e); },
    function ( ) { console.log('Completed'); });
Next: ['foo', 42]
Next: ['bar', 56]
Next: ['baz', 78]
Completed

Оператор pairs доступен в следующих дистрибутивах:

  • rx.js
  • rx.all.js
  • rx.all.compat.js
  • rx.compat.js
  • rx.lite.js
  • rx.lite.compat.js

RxPHP fromArray fromIterator asObservable fromPromise

RxPHP реализует этот оператор как fromArray.

Преобразует массив в последовательность observable

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/fromArray/fromArray.php

$source = \Rx\Observable::fromArray([1, 2, 3, 4]);

$subscription = $source->subscribe($stdoutObserver);

//Next value: 1
//Next value: 2
//Next value: 3
//Next value: 4
//Complete!
Next value: 1
Next value: 2
Next value: 3
Next value: 4
Complete!

RxPHP также имеет оператор fromIterator.

Преобразует Iterator в observable последовательность

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/iterator/iterator.php

$generator = function () {
    for ($i = 1; $i <= 3; $i++) {
        yield $i;
    }

    return 4;
};

$source = Rx\Observable::fromIterator($generator());

$source->subscribe($stdoutObserver);
Next value: 1
Next value: 2
Next value: 3
Next value: 4
Complete!

RxPHP также имеет оператор asObservable.

Скрывает идентичность последовательности observable.

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/asObservable/asObservable.php

// Create subject
$subject = new \Rx\Subject\AsyncSubject();

// Send a value
$subject->onNext(42);
$subject->onCompleted();

// Hide its type
$source = $subject->asObservable();

$source->subscribe($stdoutObserver);
Next value: 42
Complete!

RxPHP также имеет оператор fromPromise.

Преобразует promise в observable

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/promise/fromPromise.php

$promise = \React\Promise\resolve(42);

$source = \Rx\Observable::fromPromise($promise);

$subscription = $source->subscribe($stdoutObserver);
Next value: 42
Complete!

RxSwift from toObservable

В Swift это реализуется с помощью метода класса Observable.from.

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

Пример кода

let numbers = [1,2,3,4,5]

let source = Observable.from(numbers)

source.subscribe {
    print($0)
}
next(1)
next(2)
next(3)
next(4)
next(5)
completed

© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/from.html

Spec-Zone.ru

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