Наблюдаемые
Наблюдаемые — это ленивые Push-коллекции нескольких значений. Они заполняют недостающее место в следующей таблице:
| Единичный | Множественный | |
|---|---|---|
| Pull | Function | Iterator |
| Push | Promise | Observable |
Пример. Следующее — это наблюдаемое, которое выталкивает значения 1, 2, 3 сразу (синхронно) при подписке и значение 4 после того, как пройдёт секунда с момента вызова подписки, а затем завершается:
import { Observable } from 'rxjs';
const observable = new Observable((subscriber) => {
subscriber.next(1);
subscriber.next(2);
subscriber.next(3);
setTimeout(() => {
subscriber.next(4);
subscriber.complete();
}, 1000);
}); Для вызова наблюдаемого и просмотра этих значений нам нужно подписаться на него:
import { Observable } from 'rxjs';
const observable = new Observable((subscriber) => {
subscriber.next(1);
subscriber.next(2);
subscriber.next(3);
setTimeout(() => {
subscriber.next(4);
subscriber.complete();
}, 1000);
});
console.log('just before subscribe');
observable.subscribe({
next(x) {
console.log('got value ' + x);
},
error(err) {
console.error('something wrong occurred: ' + err);
},
complete() {
console.log('done');
},
});
console.log('just after subscribe'); Что выполняется следующим образом в консоли:
just before subscribe got value 1 got value 2 got value 3 just after subscribe got value 4 done
Pull против Push
Pull и Push — два разных протокола, описывающие, как производитель данных может взаимодействовать с потребителем данных.
Что такое Pull? В системах Pull потребитель определяет, когда он получает данные от производителя. Сам производитель не знает, когда данные будут доставлены потребителю.
Каждая функция JavaScript является системой Pull. Функция — это производитель данных, а код, вызывающий функцию, потребляет её, «вытягивая» одно возвращаемое значение из её вызова.
ES2015 представил генераторные функции и итераторы (function*), другой тип системы Pull. Код, который вызывает iterator.next() — это потребитель, «вытягивающий» несколько значений из итератора (производителя).
| Производитель | Потребитель | |
|---|---|---|
| Pull | Пассивный: производит данные по запросу. | Активный: решает, когда запрашиваются данные. |
| Push | Активный: производит данные со своим темпом. | Пассивный: реагирует на полученные данные. |
Что такое Push? В системах Push производитель определяет, когда отправлять данные потребителю. Потребитель не знает, когда он получит эти данные.
Обещания — это наиболее распространённый тип системы Push в JavaScript сегодня. Обещание (производитель) предоставляет разрешённое значение зарегистрированным обратным вызовам (потребителям), но в отличие от функций, именно обещание отвечает за определение точного момента, когда это значение «выталкивается» к обратным вызовам.
RxJS вводит наблюдаемые, новую систему Push для JavaScript. Наблюдаемое — это производитель нескольких значений, «выталкивающих» их к наблюдателям (потребителям).
- Функция — это лениво вычисляемое вычисление, которое синхронно возвращает одно значение при вызове.
- Генератор — это лениво вычисляемое вычисление, которое синхронно возвращает ноль или (потенциально) бесконечное количество значений при итерации.
- Обещание — это вычисление, которое может (или не может) в конечном итоге вернуть одно значение.
- Наблюдаемое — это лениво вычисляемое вычисление, которое может синхронно или асинхронно возвращать ноль или (потенциально) бесконечное количество значений с момента вызова.
Для получения дополнительной информации о том, что использовать при преобразовании наблюдаемых в обещания, см. это руководство.
Наблюдаемые как обобщения функций
Вопреки распространённому мнению, наблюдаемые не похожи на EventEmitters и не похожи на обещания для нескольких значений. Наблюдаемые могут действовать как EventEmitters в некоторых случаях, в частности, когда они многоадресны с помощью RxJS Subjects, но обычно они не действуют как EventEmitters.
Наблюдаемые похожи на функции без аргументов, но обобщают их, разрешая несколько значений.
Рассмотрим следующее:
function foo() {
console.log('Hello');
return 42;
}
const x = foo.call(); // same as foo()
console.log(x);
const y = foo.call(); // same as foo()
console.log(y); Ожидаемый вывод:
"Hello" 42 "Hello" 42
Вы можете написать такое же поведение, но с помощью наблюдаемых:
import { Observable } from 'rxjs';
const foo = new Observable((subscriber) => {
console.log('Hello');
subscriber.next(42);
});
foo.subscribe((x) => {
console.log(x);
});
foo.subscribe((y) => {
console.log(y);
}); И вывод такой же:
"Hello" 42 "Hello" 42
Это происходит потому, что и функции, и наблюдаемые — это ленивые вычисления. Если вы не вызываете функцию, console.log('Hello') не произойдёт. Также с наблюдаемыми, если вы их не «вызываете» (с subscribe), console.log('Hello') не произойдёт. Кроме того, «вызов» или «подписка» — это изолированная операция: два вызова функций вызывают два отдельных побочных эффекта, а две подписки на наблюдаемые вызывают два отдельных побочных эффекта. В отличие от EventEmitters, которые совместно используют побочные эффекты и имеют жадное выполнение независимо от наличия подписчиков, наблюдаемые не имеют общего выполнения и ленивы.
Подписка на наблюдаемое аналогична вызову функции.
Некоторые утверждают, что наблюдаемые асинхронны. Это не так. Если вы окружете вызов функции логами, как это:
console.log('before');
console.log(foo.call());
console.log('after'); Вы увидите вывод:
"before" "Hello" 42 "after"
И это то же поведение с наблюдаемыми:
console.log('before');
foo.subscribe((x) => {
console.log(x);
});
console.log('after'); И вывод:
"before" "Hello" 42 "after"
Что доказывает, что подписка на foo была полностью синхронной, как и функция.
Наблюдаемые могут доставлять значения как синхронно, так и асинхронно.
В чём разница между наблюдаемым и функцией? Наблюдаемые могут «возвращать» несколько значений со временем, чего функции не могут. Вы не можете сделать это:
function foo() {
console.log('Hello');
return 42;
return 100; // dead code. will never happen
} Функции могут возвращать только одно значение. Наблюдаемые же могут сделать это:
import { Observable } from 'rxjs';
const foo = new Observable((subscriber) => {
console.log('Hello');
subscriber.next(42);
subscriber.next(100); // "return" another value
subscriber.next(200); // "return" yet another
});
console.log('before');
foo.subscribe((x) => {
console.log(x);
});
console.log('after'); С синхронным выводом:
"before" "Hello" 42 100 200 "after"
Но вы также можете «возвращать» значения асинхронно:
import { Observable } from 'rxjs';
const foo = new Observable((subscriber) => {
console.log('Hello');
subscriber.next(42);
subscriber.next(100);
subscriber.next(200);
setTimeout(() => {
subscriber.next(300); // happens asynchronously
}, 1000);
});
console.log('before');
foo.subscribe((x) => {
console.log(x);
});
console.log('after'); С выводом:
"before" "Hello" 42 100 200 "after" 300
Заключение:
-
func.call()означает «дайте мне одно значение синхронно» -
observable.subscribe()означает «дайте мне любое количество значений, синхронно или асинхронно»
Анатомия наблюдаемого
Наблюдаемые создаются с помощью new Observable или оператора создания, подписываются с помощью наблюдателя, выполняются для доставки next / error / complete уведомлений наблюдателю, и их выполнение может быть прервано. Все эти четыре аспекта закодированы в экземпляре наблюдаемого, но некоторые из этих аспектов связаны с другими типами, например, наблюдателем и подпиской.
Основные вопросы, связанные с наблюдаемыми:
- Создание наблюдаемых
- Подписка на наблюдаемые
- Выполнение наблюдаемого
- Прерывание наблюдаемых
Создание наблюдаемых
Конструктор Observable принимает один аргумент: функцию subscribe.
В следующем примере создаётся наблюдаемое, которое посылает строку 'hi' каждую секунду подписчику.
import { Observable } from 'rxjs';
const observable = new Observable(function subscribe(subscriber) {
const id = setInterval(() => {
subscriber.next('hi');
}, 1000);
}); Наблюдаемые могут быть созданы с new Observable. Чаще всего наблюдаемые создаются с помощью функций создания, таких как of, from, interval и т. д.
В приведённом выше примере функция subscribe является наиболее важной для описания наблюдаемого. Давайте посмотрим, что означает подписка.
Подписка на наблюдаемые
Наблюдаемое observable в примере может быть подписано на него так:
observable.subscribe((x) => console.log(x));
Это не случайно, что observable.subscribe и subscribe в new Observable(function subscribe(subscriber) {...}) имеют одинаковые имена. В библиотеке они разные, но в практических целях вы можете рассматривать их как концептуально равные.
Это показывает, как вызовы subscribe не разделяются между несколькими наблюдателями одного и того же наблюдаемого. При вызове observable.subscribe с наблюдателем функция subscribe в new Observable(function subscribe(subscriber) {...}) выполняется для данного подписчика. Каждый вызов observable.subscribe запускает свою собственную независимую настройку для данного подписчика.
Подписка на наблюдаемое — это как вызов функции, предоставляющей обратные вызовы, куда будут переданы данные.
Это кардинально отличается от API обработчиков событий, таких как addEventListener / removeEventListener. С observable.subscribe, данный наблюдатель не регистрируется как слушатель в наблюдаемом. Наблюдаемое даже не поддерживает список прикреплённых наблюдателей.
Вызов subscribe — это просто способ начать «выполнение наблюдаемого» и передавать значения или события наблюдателю этого выполнения.
Выполнение наблюдаемых
Код внутри new Observable(function subscribe(subscriber) {...}) представляет собой «выполнение наблюдаемого», ленивое вычисление, которое происходит только для каждого наблюдателя, который подписывается. Выполнение производит несколько значений со временем, либо синхронно, либо асинхронно.
Существует три типа значений, которые может передавать выполнение наблюдаемого:
- Уведомление «Следующий»: отправляет значение, например, число, строку, объект и т. д.
- Уведомление «Ошибка»: отправляет ошибку JavaScript или исключение.
- Уведомление «Завершено»: не отправляет значение.
Уведомления «Следующий» являются наиболее важными и распространёнными: они представляют фактические данные, которые передаются подписчику. Уведомления «Ошибка» и «Завершено» могут произойти только один раз во время выполнения наблюдаемого, и может быть только одно из них.
Эти ограничения лучше всего выражаются в так называемой грамматике наблюдаемых или контракте, написанном как регулярное выражение:
next*(error|complete)?
Во время выполнения наблюдаемого может быть передано от нуля до бесконечного числа уведомлений «Следующий». Если передано уведомление «Ошибка» или «Завершено», то ничего больше передаваться не может.
Следующее — это пример выполнения наблюдаемого, которое отправляет три уведомления «Следующий», а затем завершается:
import { Observable } from 'rxjs';
const observable = new Observable(function subscribe(subscriber) {
subscriber.next(1);
subscriber.next(2);
subscriber.next(3);
subscriber.complete();
}); Наблюдаемые строго следуют контракту наблюдаемых, поэтому следующий код не передаст уведомление «Следующий» 4.
import { Observable } from 'rxjs';
const observable = new Observable(function subscribe(subscriber) {
subscriber.next(1);
subscriber.next(2);
subscriber.next(3);
subscriber.complete();
subscriber.next(4); // Is not delivered because it would violate the contract
}); Рекомендуется заключать любой код в subscribe с блоком try/catch, который отправит уведомление об ошибке, если он перехватит исключение:
import { Observable } from 'rxjs';
const observable = new Observable(function subscribe(subscriber) {
try {
subscriber.next(1);
subscriber.next(2);
subscriber.next(3);
subscriber.complete();
} catch (err) {
subscriber.error(err); // delivers an error if it caught one
}
}); Прерывание выполнения наблюдаемых
Поскольку выполнение Observable может быть бесконечным, и часто наблюдателю требуется прервать выполнение за конечное время, нам нужен API для отмены выполнения. Так как каждое выполнение предназначено только для одного наблюдателя, после того как наблюдатель закончит получение значений, у него должен быть способ остановить выполнение, чтобы избежать потерь вычислительной мощности или ресурсов памяти.
Когда observable.subscribe вызывается, наблюдатель присоединяется к новому выполнению Observable. Этот вызов также возвращает объект, Subscription:
const subscription = observable.subscribe((x) => console.log(x));
Подписка представляет собой текущее выполнение и имеет минимальный API, который позволяет отменить это выполнение. Узнайте больше о типе Subscription здесь. С помощью subscription.unsubscribe() вы можете отменить текущее выполнение:
import { from } from 'rxjs';
const observable = from([10, 20, 30]);
const subscription = observable.subscribe((x) => console.log(x));
// Later:
subscription.unsubscribe(); При подписке вы получаете подписку, которая представляет текущее выполнение. Просто вызовите unsubscribe() для отмены выполнения.
Каждый Observable должен определять, как освободить ресурсы этого выполнения, когда мы создаём Observable, используя create(). Вы можете сделать это, вернув пользовательскую функцию unsubscribe внутри function subscribe().
Например, вот как мы очищаем выполнение интервала, заданного с помощью setInterval:
import { Observable } from 'rxjs';
const observable = new Observable(function subscribe(subscriber) {
// Keep track of the interval resource
const intervalId = setInterval(() => {
subscriber.next('hi');
}, 1000);
// Provide a way of canceling and disposing the interval resource
return function unsubscribe() {
clearInterval(intervalId);
};
}); Так же, как observable.subscribe напоминает new Observable(function subscribe() {...}), unsubscribe, которое мы возвращаем из subscribe концептуально равно subscription.unsubscribe. На самом деле, если мы удалим типы ReactiveX, окружающие эти понятия, у нас останется довольно простой JavaScript.
function subscribe(subscriber) {
const intervalId = setInterval(() => {
subscriber.next('hi');
}, 1000);
return function unsubscribe() {
clearInterval(intervalId);
};
}
const unsubscribe = subscribe({ next: (x) => console.log(x) });
// Later:
unsubscribe(); // dispose the resources Причина, по которой мы используем типы Rx, такие как Observable, Observer и Subscription, заключается в обеспечении безопасности (например, контракта Observable) и возможности комбинирования с операторами.
© 2015–2022 Google, Inc., Netflix, Inc., Microsoft Corp. and contributors.
Code licensed under an Apache-2.0 License. Documentation licensed under CC BY 4.0.
https://rxjs.dev/guide/observable