Spec-Zone.ru › ReactiveX

Наблюдаемый объект

В ReactiveX наблюдатель подписывается на наблюдаемый объект. Затем этот наблюдатель реагирует на любой элемент или последовательность элементов, которые излучает наблюдаемый объект. Эта модель облегчает одновременные операции, так как ей не нужно блокироваться, ожидая, пока наблюдаемый объект излучит объекты, а вместо этого она создаёт сторожа в виде наблюдателя, готового должным образом отреагировать в любой момент будущего, когда наблюдаемый объект это сделает.

Эта страница объясняет, что такое реактивная модель и что такое наблюдаемые объекты и наблюдатели (и как наблюдатели подписываются на наблюдаемые объекты). Другие страницы показывают, как вы используете разнообразные операторы для наблюдаемых объектов, чтобы связать наблюдаемые объекты вместе и изменить их поведение.

Данная документация сопровождает свои объяснения «диаграммами жемчужин». Вот как диаграммы жемчужин представляют наблюдаемые объекты и преобразования наблюдаемых объектов:

См. также

  • Single — специализированная версия наблюдаемого объекта, излучающая только один элемент
  • Rx Workshop: Введение
  • Введение в Rx: IObservable
  • Освоение наблюдаемых объектов (из документации Couchbase Server)
  • Двухминутное введение в Rx Андреа Стальтц («Представьте наблюдаемый объект как асинхронный неизменяемый массив»)
  • Введение в наблюдаемый объект Джафара Хусейна (видеоурок по JavaScript)
  • Объект наблюдаемый объект (RxJS) Денниса Стоянова
  • Преобразование обратного вызова в Rx-наблюдаемый объект @afterecho

Общие сведения

Во многих задачах программирования на основе программного обеспечения вы в той или иной степени ожидаете, что написанные вами инструкции будут выполняться и завершаться постепенно, по одной за раз, в том порядке, в котором вы их написали. Но в ReactiveX многие инструкции могут выполняться параллельно, а их результаты позже захватываются «наблюдателями» в произвольном порядке. Вместо того чтобы *вызывать* метод, вы определяете механизм для извлечения и преобразования данных в виде «наблюдаемого объекта», а затем *подписываете* наблюдателя на него, в этот момент ранее определённый механизм приходит в действие, и наблюдатель служит сторожем, чтобы захватывать и реагировать на его излучения всякий раз, когда они будут готовы.

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

Существует множество терминов, используемых для описания этой модели асинхронного программирования и проектирования. В этом документе будут использоваться следующие термины: наблюдатель подписывается на наблюдаемый объект. Наблюдаемый объект излучает элементы или отправляет уведомления своим наблюдателям, вызывая методы наблюдателей.

В других документах и в других контекстах то, что мы называем «наблюдателем», иногда называется «подписчиком», «наблюдателем» или «реактором». Эта модель в целом часто называется «паттерном реактора».

Создание наблюдателей

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

В обычном вызове метода — то есть *не* в типе асинхронных, параллельных вызовов, характерных для ReactiveX — поток выглядит примерно так:

  1. Вызовите метод.
  2. Сохраните возвращаемое значение этого метода в переменной.
  3. Используйте эту переменную и её новое значение для выполнения полезной работы.

Или примерно так:

// make the call, assign its return value to `returnVal`
returnVal = someMethod(itsParameters);
// do something useful with returnVal

В асинхронной модели поток выглядит примерно так:

  1. Определите метод, который выполняет полезную работу с возвращаемым значением асинхронного вызова; этот метод является частью *наблюдателя*.
  2. Определите сам асинхронный вызов как *наблюдаемый объект*.
  3. Присоедините наблюдателя к этому наблюдаемому объекту, *подписав* его (это также запускает действия наблюдаемого объекта).
  4. Продолжайте выполнять свою работу; всякий раз, когда вызов возвращается, метод наблюдателя начнёт работать со своим возвращаемым значением или значениями — *элементами*, излучаемыми наблюдаемым объектом.

Что-то вроде этого:

// defines, but does not invoke, the Subscriber's onNext handler
// (in this example, the observer is very simple and has only an onNext handler)
def myOnNext = { it -> do something useful with it };
// defines, but does not invoke, the Observable
def myObservable = someObservable(itsParameters);
// subscribes the Subscriber to the Observable, and invokes the Observable
myObservable.subscribe(myOnNext);
// go on about my business

onNext, onCompleted и onError

Метод Subscribe служит для подключения наблюдателя к наблюдаемому объекту. Ваш наблюдатель реализует подмножество следующих методов:

onNext
Наблюдаемый объект вызывает этот метод всякий раз, когда наблюдаемый объект излучает элемент. Этот метод принимает в качестве параметра излучаемый наблюдаемым объектом элемент.
onError
Наблюдаемый объект вызывает этот метод, чтобы указать, что ему не удалось сгенерировать ожидаемые данные или что он столкнулся с другой ошибкой. Он не будет делать дальнейших вызовов к onNext или onCompleted. Метод onError принимает в качестве параметра указание на причину ошибки.
onCompleted
Наблюдаемый объект вызывает этот метод после того, как он вызвал onNext в последний раз, если он не столкнулся с ошибками.

Согласно контракту Observable, он может вызвать onNext ноль или более раз, а затем может последовать за этими вызовами вызовом либо onCompleted, либо onError, но не обоих, что будет его последним вызовом. По соглашению, в этом документе вызовы к onNext обычно называют «излучениями» элементов, тогда как вызовы к onCompleted или onError называют «уведомлениями».

Более полный пример вызова subscribe выглядит так:

def myOnNext     = { item -> /* do something useful with item */ };
def myError      = { throwable -> /* react sensibly to a failed call */ };
def myComplete   = { /* clean up after the final response */ };
def myObservable = someMethod(itsParameters);
myObservable.subscribe(myOnNext, myError, myComplete);
// go on about my business

См. также

  • Введение в Rx: IObserver

Отмена подписки

В некоторых реализациях ReactiveX есть специализированный интерфейс наблюдателя, Subscriber, который реализует метод unsubscribe. Вы можете вызвать этот метод, чтобы указать, что подписчик больше не заинтересован ни в одном из наблюдаемых объектов, на которые он сейчас подписан. Эти наблюдаемые объекты могут (если у них нет других заинтересованных наблюдателей) выбрать прекращение генерации новых элементов для излучения.

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

Некоторые заметки по соглашениям об именовании

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

Кроме того, некоторые из этих имён имеют разное значение в других контекстах или кажутся неудобными в идиоме конкретного языка реализации.

Например, существует паттерн именования onEvent (например, onNext, onCompleted, onError). В некоторых контекстах такие имена указывали бы на методы, с помощью которых регистрируются обработчики событий. Однако в ReactiveX они называют сами обработчики событий.

«Горячие» и «холодные» наблюдаемые объекты

Когда наблюдаемый объект начинает излучать свою последовательность элементов? Это зависит от наблюдаемого объекта. «Горячий» наблюдаемый объект может начать излучать элементы сразу же после создания, и поэтому любой наблюдатель, который позже подпишется на этот наблюдаемый объект, может начать наблюдать последовательность где-то в середине. С другой стороны, «холодный» наблюдаемый объект ждёт, пока наблюдатель подпишется на него, прежде чем начать излучать элементы, и поэтому такой наблюдатель гарантированно увидит всю последовательность с самого начала.

В некоторых реализациях ReactiveX также есть нечто, называемое «соединяемым» наблюдаемым объектом. Такой наблюдаемый объект не начинает излучать элементы до тех пор, пока не будет вызван его метод Connect, независимо от того, подписались ли на него какие-либо наблюдатели.

Составление через операторы наблюдаемых объектов

Наблюдаемые объекты и наблюдатели — это только начало ReactiveX. Сами по себе они были бы лишь небольшим расширением стандартного паттерна наблюдателя, лучше подходящего для обработки последовательности событий, а не единственного обратного вызова.

Настоящая мощь заключается в «реактивных расширениях» (отсюда и «ReactiveX») — операторах, которые позволяют преобразовывать, комбинировать, манипулировать и работать с последовательностями элементов, излучаемых наблюдаемыми объектами.

Эти операторы Rx позволяют вам составлять асинхронные последовательности объявляющим способом со всеми преимуществами эффективности обратных вызовов, но без недостатков вложенных обработчиков обратных вызовов, обычно связанных с асинхронными системами.

Данная документация группирует информацию о различных операторах и примерах их использования на следующих страницах:

Создание наблюдаемых
Create, Defer, Empty/Never/Throw, From, Interval, Just, Range, Repeat, Start, и Timer
Преобразование элементов наблюдаемых
Buffer, FlatMap, GroupBy, Map, Scan, и Window
Фильтрация наблюдаемых
Debounce, Distinct, ElementAt, Filter, First, IgnoreElements, Last, Sample, Skip, SkipLast, Take, и TakeLast
Комбинирование наблюдаемых
And/Then/When, CombineLatest, Join, Merge, StartWith, Switch, и Zip
Операторы обработки ошибок
Catch и Retry
Утилитарные операторы
Delay, Do, Materialize/Dematerialize, ObserveOn, Serialize, Subscribe, SubscribeOn, TimeInterval, Timeout, Timestamp, и Using
Условные и булевы операторы
All, Amb, Contains, DefaultIfEmpty, SequenceEqual, SkipUntil, SkipWhile, TakeUntil, и TakeWhile
Математические и агрегатные операторы
Average, Concat, Count, Max, Min, Reduce, и Sum
Преобразование наблюдаемых
To
Подключаемые операторы наблюдаемых
Connect, Publish, RefCount, и Replay
Операторы обратной загрузки
разнообразные операторы, обеспечивающие определенные политики управления потоком

Эти страницы содержат информацию о некоторых операторах, которые не являются частью ядра ReactiveX, но реализованы в одной или нескольких языковых реализациях и/или дополнительных модулях.

Цепочки операторов

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

Существуют и другие шаблоны, такие как шаблон Builder, в котором различные методы определённого класса работают с элементом того же класса, изменяя этот объект через выполнение метода. Эти шаблоны также позволяют вам создавать цепочки методов аналогичным образом. Но в то время как в шаблоне Builder порядок следования методов в цепочке обычно не имеет значения, с операторами Observable порядок важен.

Цепочка операторов Observable не работает независимо с исходным наблюдаемым, который инициирует цепочку, а работает по очереди, каждый оператор работает с наблюдаемым, сгенерированным оператором, непосредственно предшествующим ему в цепочке.

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

Spec-Zone.ru

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