Spec-Zone.ru › ReactiveX

Операторы обратной зависимости

стратегии для работы с Observable, которые производят элементы быстрее, чем их наблюдатели их потребляют

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

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

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

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

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

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

Холодные Observable идеально подходят для модели реактивной обратной зависимости с подтягиванием, реализованной в некоторых реализациях ReactiveX (которое описано в другом месте). Горячие Observable обычно плохо справляются с моделью реактивной обратной зависимости с подтягиванием и являются лучшими кандидатами для других стратегий управления потоком, таких как использование операторов, описанных на этой странице, или операторов, таких как Buffer, Sample, Debounce или Window.

См. также

  • Buffer
  • Sample
  • Debounce
  • Window

Информация о языке

RxGroovy onBackpressureBuffer onBackpressureDrop onBackpressureLatest

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

onBackpressureBuffer

onBackpressureBuffer сохраняет буфер всех неотслеженных эмиссий из исходного Observable и отправляет их в наблюдатели вниз по потоку в соответствии с запросами, которые они генерируют.

Версия этого оператора, представленная в RxGroovy 1.1, позволяет задать емкость буфера; применение этого оператора приведет к завершению результирующего Observable с ошибкой, если этот буфер будет переполнен. Вторая версия, представленная в том же релизе, позволяет задать Action для вызова onBackpressureBuffer в случае переполнения буфера.

  • Javadoc: onBackpressureBuffer()
  • Javadoc: onBackpressureBuffer(long) (RxGroovy 1.1)
  • Javadoc: onBackpressureBuffer(long, Action0) (RxGroovy 1.1)
onBackpressureDrop

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

Версия этого оператора, представленная в релизе 1.1, уведомляет вас посредством Action, которое вы передаёте в качестве параметра, когда элемент был отброшен и какой именно элемент был отброшен.

  • Javadoc: onBackpressureDrop()
  • Javadoc: onBackpressureDrop(Action1) (RxGroovy 1.1)
onBackpressureLatest

onBackpressureLatest (новая в RxJava 1.1) удерживает последний испущенный элемент из исходного Observable и немедленно отправляет этот элемент своему наблюдателю по запросу. Он отбрасывает любые другие элементы, которые он наблюдает между запросами от своего наблюдателя.

  • Javadoc: onBackpressureLatest()

RxJava 1․x onBackpressureBuffer onBackpressureDrop onBackpressureLatest

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

onBackpressureBuffer

onBackpressureBuffer сохраняет буфер всех неотслеженных эмиссий из исходного Observable и отправляет их в наблюдатели вниз по потоку в соответствии с запросами, которые они генерируют.

Версия этого оператора, представленная в RxJava 1.1, позволяет задать емкость буфера; применение этого оператора приведет к завершению результирующего Observable с ошибкой, если этот буфер будет переполнен. Вторая версия, представленная в том же релизе, позволяет задать Action для вызова onBackpressureBuffer в случае переполнения буфера.

  • Javadoc: onBackpressureBuffer()
  • Javadoc: onBackpressureBuffer(long) (RxJava 1.1)
  • Javadoc: onBackpressureBuffer(long, Action0) (RxJava 1.1)
onBackpressureDrop

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

Версия этого оператора, представленная в релизе 1.1, уведомляет вас посредством Action, которое вы передаёте в качестве параметра, когда элемент был отброшен и какой именно элемент был отброшен.

  • Javadoc: onBackpressureDrop()
  • Javadoc: onBackpressureDrop(Action1) (RxJava 1.1)
onBackpressureLatest

onBackpressureLatest (новая в RxJava 1.1) удерживает последний испущенный элемент из исходного Observable и немедленно отправляет этот элемент своему наблюдателю по запросу. Он отбрасывает любые другие элементы, которые он наблюдает между запросами от своего наблюдателя.

  • Javadoc: onBackpressureLatest()

RxJS controlled pausable pausableBuffered stopAndWait windowed

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

stopAndWait

В качестве альтернативы использованию request для извлечения элементов из ControlledObservable, можно применить оператор stopAndWait. Этот оператор запросит новый элемент из Observable каждый раз, когда процедура наблюдателя onNext получит последний элемент.

windowed

Ещё один вариант — использовать оператор windowed(n). Он работает аналогично stopAndWait, но имеет внутренний буфер из n элементов, что позволяет ControlledObservable работать немного впереди наблюдателя время от времени. windowed(1) эквивалентен stopAndWait.

Также есть два оператора, которые преобразуют обычное Observable в паузную форму.

pausable

Если вызвать метод pause созданного с помощью оператора pausable PausableObservable, он отбросит (проигнорирует) любые элементы, испущенные базовым исходным Observable, до тех пор, пока вы не вызовете его метод resume, после чего он продолжит передавать испущенные элементы своим наблюдателям.

См. также

  • RxMarbles: pausable
pausableBuffered

Если вызвать метод pause созданного с помощью оператора pausableBuffered PausableObservable, он будет буферизовать любые элементы, испущенные базовым исходным Observable, до тех пор, пока вы не вызовете его метод resume, после чего он отправит эти буферизованные элементы, а затем продолжит передавать любые дополнительные испущенные элементы своим наблюдателям.

См. также

  • RxMarbles: pausableBuffered

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

Spec-Zone.ru

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