ReactiveX
ReactiveX — это библиотека для составления асинхронных и основанных на событиях программ с помощью наблюдаемых последовательностей.
Она расширяет паттерн наблюдателя, чтобы поддерживать последовательности данных и/или событий, и добавляет операторы, которые позволяют вам объединять последовательности вместе декларативно, абстрагируясь от таких проблем, как низкоуровневая многопоточность, синхронизация, потокобезопасность, конкурирующие структуры данных и неблокирующее ввод-вывод.
| Наблюдаемые последовательности заполняют пробел, являясь идеальным способом доступа к асинхронным последовательностям из нескольких элементов | ||
|---|---|---|
| единичные элементы | множественные элементы | |
| синхронный | T getData() |
Iterable<T> getData() |
| асинхронный | Future<T> getData() |
Observable<T> getData() |
Иногда её называют «функционально-реактивным программированием», но это неверное название. ReactiveX может быть функциональной и реактивной, но «функционально-реактивное программирование» — это другое понятие. Главное различие состоит в том, что функционально-реактивное программирование оперирует значениями, которые изменяются непрерывно во времени, в то время как ReactiveX оперирует дискретными значениями, которые излучаются со временем. (См. работу Конала Эллиотта для более точной информации о функционально-реактивном программировании.)
Почему использовать наблюдаемые последовательности?
Модель ReactiveX Observable позволяет вам рассматривать потоки асинхронных событий с помощью тех же простых и объединяемых операций, что вы используете для коллекций данных, таких как массивы. Это освобождает вас от запутанных цепочек обратных вызовов и, следовательно, делает ваш код более читабельным и менее подверженным ошибкам.
Наблюдаемые последовательности объединяются
Такие методы, как Java Futures, просты в использовании для одноуровневого асинхронного выполнения, но они начинают добавлять значительную сложность при их вложенности.
Использование Futures для оптимального составления условных асинхронных потоков выполнения сложно (или невозможно, так как задержки каждого запроса меняются во время выполнения). Конечно, это можно сделать, но быстро становится сложным (и, следовательно, подверженным ошибкам) или преждевременно блокируется на Future.get(), что устраняет преимущества асинхронного выполнения.
С другой стороны, наблюдаемые последовательности ReactiveX предназначены для объединения потоков и последовательностей асинхронных данных.
Наблюдаемые последовательности гибкие
Наблюдаемые последовательности ReactiveX поддерживают не только выброс отдельных скалярных значений (как Futures), но также и последовательностей значений или даже бесконечных потоков. Observable — это единая абстракция, которая может быть использована для любого из этих вариантов использования. Наблюдаемая последовательность имеет всю гибкость и элегантность, связанную с её аналогом — Iterable.
| Наблюдаемая последовательность — это асинхронная/потоковая “дуальная” версия синхронной/вытягивающей Iterable | ||
|---|---|---|
| событие | Iterable (вытягивание) | Observable (поток) |
| получить данные | T next() |
onNext(T) |
| обнаружение ошибки | бросает Exception
|
onError(Exception) |
| завершение | !hasNext() |
onCompleted() |
Наблюдаемые последовательности менее предвзяты
ReactiveX не предвзято относится к определенному источнику конкурентности или асинхронности. Наблюдаемые последовательности могут быть реализованы с помощью пулов потоков, циклов событий, неблокирующего ввода-вывода, акторов (например, из Akka) или любой другой реализации, соответствующей вашим потребностям, стилю или опыту. Код клиента рассматривает все взаимодействия с наблюдаемыми последовательностями как асинхронные, независимо от того, является ли ваша основная реализация блокирующей или неблокирующей, и как вы её реализовали.
| Как реализована эта наблюдаемая последовательность? |
|---|
public Observable<data> getData(); |
| С точки зрения наблюдателя, это не имеет значения! |
|
И, что важно: с помощью ReactiveX вы можете позже изменить решение и радикально изменить базовую природу реализации вашей наблюдаемой последовательности, не нарушив потребителей вашей наблюдаемой последовательности.
Обратные вызовы имеют свои проблемы
Обратные вызовы решают проблему преждевременного блокирования на Future.get(), не позволяя ничего заблокировать. Они естественно эффективны, потому что выполняются, когда ответ готов.
Но, как и в случае с Futures, в то время как обратные вызовы легко использовать с одноуровневым асинхронным выполнением, при вложенном составлении они становятся громоздкими.
ReactiveX — это многоязыковая реализация
В настоящее время ReactiveX реализован на различных языках, соответствующих особенностям этих языков, и добавляются новые языки с большой скоростью.
Реактивное программирование
ReactiveX предоставляет набор операторов, с помощью которых можно фильтровать, выбирать, преобразовывать, комбинировать и составлять наблюдаемые последовательности. Это позволяет эффективно выполнять и составлять их.
Вы можете рассматривать класс Observable как «потоковую» эквивалент Iterable, которая является «вытягивающей». С Iterable потребитель вытягивает значения из производителя, и поток блокируется до тех пор, пока эти значения не придут. Напротив, с Observable производитель толкает значения потребителю всякий раз, когда значения доступны. Этот подход более гибкий, так как значения могут поступать синхронно или асинхронно.
| Пример кода, показывающий, как похожие высшего порядка функции могут применяться к Iterable и Observable | |
|---|---|
Iterable |
Observable |
getDataFromLocalMemory()
.skip(10)
.take(5)
.map({ s -> return s + " transformed" })
.forEach({ println "next => " + it }) | getDataFromNetwork()
.skip(10)
.take(5)
.map({ s -> return s + " transformed" })
.subscribe({ println "onNext => " + it }) |
Тип Observable добавляет два недостающих семантических понятия к патерну наблюдателя Gang of Four, чтобы соответствовать тем, которые доступны в типе Iterable:
- способность производителя сигнализировать потребителю о том, что больше нет данных (цикл foreach для Iterable завершается и возвращается нормально в таком случае; Observable вызывает метод
onCompletedсвоего наблюдателя) - способность производителя сигнализировать потребителю о возникновении ошибки (Iterable бросает исключение, если во время итерации происходит ошибка; Observable вызывает метод
onErrorсвоего наблюдателя)
С этими дополнениями ReactiveX согласовывает типы Iterable и Observable. Единственное различие между ними заключается в направлении потока данных. Это очень важно, потому что теперь любую операцию, которую вы можете выполнить над Iterable, вы также можете выполнить над Observable.
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/intro.html