Spec-Zone.ru › ReactiveX

Реализация собственных операторов

Вы можете реализовать собственные операторы Observable. На этой странице показано, как это сделать.

Если ваш оператор предназначен для создания Observable, а не для преобразования или реакции на исходный Observable, используйте метод create( ), а не пытайтесь реализовать Observable вручную. В противном случае следуйте инструкциям ниже.

Цепочечное применение пользовательских операторов со стандартными операторами RxJava

Следующий пример демонстрирует, как вы можете объединить пользовательский оператор (в данном примере: myOperator) со стандартными операторами RxJava, используя оператор lift( ):

Observable foo = barObservable.ofType(Integer).map({it*2}).lift(new myOperator<T>()).map({"transformed by myOperator: " + it});

В следующем разделе будет показано, как создать каркас вашего оператора, чтобы он корректно работал с lift( ).

Реализация оператора

Определите свой оператор как публичный класс, который реализует интерфейс Operator, как показано ниже:

public class myOperator<T> implements Operator<T> {
  public myOperator( /* any necessary params here */ ) {
    /* any necessary initialization here */
  }

  @Override
  public Subscriber<? super T> call(final Subscriber<? super T> s) {
    return new Subscriber<t>(s) {
      @Override
      public void onCompleted() {
        /* add your own onCompleted behavior here, or just pass the completed notification through: */
        if(!s.isUnsubscribed()) {
          s.onCompleted();
        }
      }

      @Override
      public void onError(Throwable t) {
        /* add your own onError behavior here, or just pass the error notification through: */
        if(!s.isUnsubscribed()) {
          s.onError(t);
        }
      }

      @Override
      public void onNext(T item) {
        /* this example performs some sort of simple transformation on each incoming item and then passes it along */
        if(!s.isUnsubscribed()) {
          transformedItem = myOperatorTransformOperation(item);
          s.onNext(transformedItem);
        }
      }
    };
  }
}

Другие соображения

  • Ваш оператор должен проверять статус подписчика isUnsubscribed( ) перед отправкой любого элемента или уведомления подписчику. Не тратьте время на генерацию элементов, которыми подписчик не заинтересован.
  • Ваш оператор должен следовать основным принципам контракта Observable:
    • Он может вызывать метод подписчика onNext( ) любое количество раз, но эти вызовы не должны перекрываться.
    • Он может вызвать либо метод подписчика onCompleted( ), либо onError( ), но не оба, ровно один раз, и он не должен впоследствии вызывать метод подписчика onNext( ).
    • Если вы не можете гарантировать, что ваш оператор соответствует вышеуказанным двум принципам, вы можете добавить к нему оператор serialize( ), чтобы принудительно обеспечить правильное поведение.
  • Не блокируйте выполнение внутри вашего оператора.
  • В идеале, новые операторы следует создавать путём комбинирования уже существующих, в тех случаях, когда это возможно, а не изобретать велосипед. RxJava делает это с некоторыми своими стандартными операторами, например:
    • first( ) определен как take(1).single( )
    • ignoreElements( ) определен как filter(alwaysFalse( ))
    • reduce(a) определен как scan(a).last( )
  • Если ваш оператор использует функции или лямбда-выражения, которые передаются в качестве параметров (например, предикаты), имейте в виду, что они могут быть источником исключений, и будьте готовы перехватывать эти исключения и уведомлять подписчиков с помощью вызовов onError( ).
  • В общем случае уведомляйте подписчиков об ошибках немедленно, а не делайте попытку сначала отправить больше элементов.
  • В некоторых реализациях ReactiveX ваш оператор может нуждаться в чувствительности к стратегиям «backpressure» этой реализации. (См., например: Проблемы при реализации операторов (часть 2) Дэвида Карнока.)

См. также

  • Проблемы при реализации операторов (часть 1) и (часть 2) Дэвида Карнока.
  • Реализация собственных операторов Observable (в RxJS) Денниса Стойанова

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

Spec-Zone.ru

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