Реализация собственных операторов
Вы можете реализовать собственные операторы 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