Hola a todos. Soy Sergey Omelnitsky. No hace mucho, realicé una transmisión en vivo sobre programación reactiva, donde hablé sobre la asincronía en JavaScript. Hoy me gustaría resumir este material.

Pero antes de comenzar con el material principal, necesitamos hacer una introducción. Así que, empecemos con definiciones: ¿qué es una pila y una cola?
Pila es una colección cuyos elementos se obtienen según el principio de «último en entrar, primero en salir» LIFO.
Cola es una colección cuyos elementos se obtienen según el principio de «primero en entrar, primero en salir» FIFO.
Está bien, continuemos.

JavaScript es un lenguaje de programación de un solo hilo. Esto significa que tiene un solo hilo de ejecución y una sola pila, en la cual se colocan las funciones en cola para su ejecución. Por lo tanto, en un momento dado, JavaScript solo puede ejecutar una operación, las otras operaciones esperarán su turno en la pila hasta que sean llamadas.
Pila de llamadas es una estructura de datos que, simplificando, registra información sobre el lugar en el programa donde nos encontramos. Si entramos en una función, colocamos un registro de ella en la parte superior de la pila. Cuando regresamos de la función, sacamos el elemento más alto de la pila y retornamos al punto donde se llamó a esa función. Eso es todo lo que puede hacer una pila. Y ahora surge una pregunta muy interesante. ¿Cómo funciona entonces la asincronía en JavaScript?

En realidad, además de la pila, en los navegadores hay una cola especial para trabajar con lo que se llama WebAPI. Las funciones de esta cola se ejecutan en orden solo después de que la pila esté completamente vacía. Solo después de eso se trasladan de la cola a la pila para su ejecución. Si en la pila hay al menos un elemento, no pueden entrar en la pila. Precisamente por eso, la llamada a funciones con timeout a menudo no es precisa en el tiempo, ya que una función no puede pasar de la cola a la pila mientras esta esté ocupada.
Veamos el siguiente ejemplo y abordemos su «ejecución» paso a paso. También analicemos qué sucede en el sistema en ese momento.
console.log('Hi);
setTimeout(function cb1() {
console.log('cb1');
}, 5000);
console.log('Bye');
1) Por ahora, no sucede nada. La consola del navegador está limpia, y la pila de llamadas está vacía.

2) Luego, el comando console.log('Hi') se añade a la pila de llamadas.

3) Y se ejecuta.

4) Luego, console.log(‘Hi’) se elimina de la pila de llamadas.

5) Ahora pasamos al comando setTimeout(function cb1() {… }). Este se añade a la pila de llamadas.

6) El comando setTimeout(function cb1() {… }) se ejecuta. El navegador crea un temporizador, que es parte de la API web. Este hará una cuenta regresiva.

7) El comando setTimeout(function cb1() {… }) ha terminado y se elimina de la pila de llamadas.

8) El comando console.log(‘Bye’) se añade a la pila de llamadas.

9) El comando console.log(‘Bye’) se ejecuta.

10) El comando console.log(‘Bye’) se elimina de la pila de llamadas.

11) Después de al menos 5000 ms, el temporizador termina y coloca el callback cb1 en la cola de callbacks.

12) El ciclo de eventos toma la función cb1 de la cola de callbacks y la coloca en la pila de llamadas.

13) La función cb1 se ejecuta y añade console.log(‘cb1’) a la pila de llamadas.

14) El comando console.log(‘cb1’) se ejecuta.

15) El comando console.log(‘cb1’) se elimina de la pila de llamadas.

16) La función cb1 se elimina de la pila de llamadas.
Veamos un ejemplo en dinámica:

Ya hemos discutido cómo se implementa la asincronía en JavaScript. Ahora hablemos brevemente sobre la evolución del código asíncrono.
Evolución del código asíncrono.
a(function (resultsFromA) {
b(resultsFromA, function (resultsFromB) {
c(resultsFromB, function (resultsFromC) {
d(resultsFromC, function (resultsFromD) {
e(resultsFromD, function (resultsFromE) {
f(resultsFromE, function (resultsFromF) {
console.log(resultsFromF);
})
})
})
})
})
});La programación asíncrona, tal como la conocemos en JavaScript, solo puede ser implementada mediante funciones. Estas pueden ser pasadas como cualquier otra variable a otras funciones. Así nacieron los callbacks. Y esto es genial, divertido y entretenido, hasta que se convierte en tristeza, melancolía y pena. ¿Por qué? Pues por algo simple:
- Con el aumento de la complejidad del código, el proyecto se convierte rápidamente en bloques anidados difíciles de entender: el infierno de los callbacks.
- El manejo de errores se puede pasar por alto fácilmente.
- No se pueden devolver expresiones con return.
Con la llegada de Promise, la situación mejoró un poco.
new Promise(function(resolve, reject) {
setTimeout(() => resolve(1), 2000);
}).then((result) => {
alert(result);
return result + 2;
}).then((result) => {
throw new Error('FAILED HERE');
alert(result);
return result + 2;
}).then((result) => {
alert(result);
return result + 2;
}).catch((e) => {
console.log('error: ', e);
});- Se introdujeron cadenas de promesas, lo que mejoró la legibilidad del código.
- Apareció un método separado para la captura de errores.
- Ahora hay una opción para la ejecución paralela mediante Promise.all
- La asincronía anidada se puede resolver mediante async/await
Pero las promesas tienen sus propias limitaciones. Por ejemplo, una promesa no se puede cancelar sin complicaciones, y lo más importante es que trabaja con un solo valor.
Así que hemos llegado suavemente a la programación reactiva. ¿Cansado? Bueno, puedes ir a preparar un té, pensar un poco y regresar a seguir leyendo. Yo continuaré.

Programación reactiva — una paradigma de programación enfocada en flujos de datos y en la propagación de cambios. Vamos a analizar más a fondo qué es un flujo de datos.
// Получаем ссылку на элемент
const input = ducument.querySelector('input');
const eventsArray = [];
// Пушим каждое событие в массив eventsArray
input.addEventListener('keyup',
event => eventsArray.push(event)
);Imaginemos que tenemos un campo de entrada. Creamos un array y en cada evento keyup del input guardaremos el evento en nuestro array. Me gustaría señalar que nuestro array está ordenado por tiempo, es decir, el índice de los eventos más recientes es mayor que el índice de los más antiguos. Este array representa un modelo simplificado de un flujo de datos, pero todavía no es un flujo. Para que este array se pueda considerar un flujo, debe ser capaz de comunicar a los suscriptores que ha recibido nuevos datos. Así hemos llegado a la definición de un flujo.
Flujo de datos
const { interval1 } = Rx;
const { take } = RxOperators;
interval(1000).pipe(
take(4)
)
Flujo — es un array de datos ordenados por tiempo que puede informar sobre cambios en los datos. Ahora imagina lo conveniente que se vuelve escribir código en el cual una acción requiere disparar varios eventos en diferentes partes del código. Simplemente nos suscribimos al flujo y él nos informará cuándo ocurren cambios. Y eso es lo que la biblioteca RxJs puede hacer.

RxJS — es una biblioteca para trabajar con programas asíncronos y basados en eventos utilizando secuencias observables. La biblioteca proporciona un tipo básico Observable, varios tipos auxiliares (Observer, Schedulers, Subjects) y operadores para trabajar con eventos como colecciones (map, filter, reduce, every y similares de JavaScript Array).
Vamos a entender los conceptos básicos de esta biblioteca.
Observable, Observer, Producer
Observable — el primer tipo básico que vamos a revisar. Esta clase contiene la parte principal de la implementación de RxJs. Está relacionada con el flujo observable, al cual se puede suscribir mediante el método subscribe.
En Observable se implementa un mecanismo auxiliar para crear actualizaciones, conocido como Observer. La fuente de valores para el Observer se llama Producer. Puede ser un array, un iterador, un socket web, algún evento, etc. Así que podemos decir que el observable actúa como un intermediario entre Producer y Observer.
Observable maneja tres tipos de eventos del Observer:
- next – nuevos datos
- error – un error, si la secuencia se termina debido a una situación excepcional. Este evento también implica el final de la secuencia.
- complete — una señal de que la secuencia ha terminado. Esto significa que no habrá más nuevos datos
Veamos una demostración:

Al principio, procesaremos los valores 1, 2, 3, y después de 1 segundo, obtendremos 4 y finalizaremos nuestro flujo.
Pensando en voz alta
Y aquí me di cuenta de que era más interesante contar que escribir sobre esto. 😀
Subscription
Cuando nos suscribimos a un flujo, creamos una nueva clase subscription, que nos permite cancelar la suscripción mediante el método unsubscribe. También podemos agrupar suscripciones usando el método add. Y, por supuesto, podemos desagrupar flujos usando eliminar. Los métodos add y remove toman otra suscripción como entrada. Me gustaría destacar que cuando cancelamos una suscripción, estamos cancelando todas las suscripciones secundarias como si hubiéramos llamado al método unsubscribe en ellas. Sigamos.
Tipos de flujos
| HOT | COLD |
|---|---|
| El Producer se crea fuera del observable | El Producer se crea dentro del observable |
| Los datos se transmiten en el momento de la creación del observable | Los datos se comunican en el momento de la suscripción |
| Se necesita lógica adicional para cancelar la suscripción | El flujo se completa por sí mismo |
| Utiliza una relación uno a muchos | Utiliza una relación uno a uno |
| Todas las suscripciones tienen un valor único | Las suscripciones son independientes |
| Los datos se pueden perder si no hay suscripción | Reemite todos los valores del flujo para una nueva suscripción |
Si hiciera una analogía, representaría un flujo caliente como una película en el cine. En el momento en que llegas, comienzas a verla desde ese punto. Un flujo frío lo compararía con una llamada al soporte técnico. Cualquiera que llame escucha la grabación del contestador automático de principio a fin, pero puedes colgar usando unsubscribe.
Me gustaría señalar que también existen lo que se llaman flujos cálidos (definición que he encontrado raramente y solo en comunidades extranjeras) — es un flujo que se transforma de un flujo frío a uno caliente. Surge la pregunta: ¿dónde usarlo?)) Daré un ejemplo práctico.
Trabajo con Angular. Utiliza activamente RxJS. Para obtener datos del servidor, espero un flujo frío y uso este flujo en la plantilla con asyncPipe. Si uso este pipe varias veces, volviendo a la definición del flujo frío, cada pipe solicitará datos al servidor, lo cual, por decirlo suavemente, es extraño. Pero si convierto el flujo frío en uno cálido, la solicitud se realizará una única vez.
En general, entender los tipos de flujos es bastante complicado para los principiantes, pero es importante.
Operadores
return this.http.get(`${environment.apiUrl}/${this.apiUrl}/trade_companies`)
.pipe(
tap(({ data }: TradeCompanyList) => this.companies$$.next(cloneDeep(data))),
map(({ data }: TradeCompanyList) => data)
);Los operadores nos brindan la capacidad de trabajar con flujos. Ayudan a controlar los eventos que ocurren en un Observable. Vamos a revisar un par de los más populares, y se puede encontrar más información sobre los operadores en los enlaces de información útil.
Operadores — of
Comencemos con el operador auxiliar of. Crea un Observable a partir de un valor simple.

Operadores — filter

El operador de filtrado filter, como se puede entender por su nombre, filtra la señal del flujo. Si el operador devuelve verdadero, permite que continúe.
Operadores — take

take — Acepta el número de emisiones, después de las cuales finaliza el flujo.
Operadores — debounceTime

debounceTime — descarta los valores emitidos que ocurren en el intervalo de tiempo especificado entre las salidas; tras el transcurso del intervalo de tiempo, emite el último valor.
const { Observable } = Rx;
const { debounceTime, take } = RxOperators;
Observable.create((observer) => {
let i = 1;
observer.next(i++);
// Emitir valor cada 1000ms
setInterval(() => {
observer.next(i++)
}, 1000);
// Emitir valor cada 1500ms
setInterval(() => {
observer.next(i++)
}, 1500);
}).pipe(
debounceTime(700), // Esperar 700ms antes de procesar el valor
take(3)
); 
Operadores — takeWhile

Emite valores hasta que takeWhile devuelva falso, después de lo cual se desuscribe del flujo.
const { Observable } = Rx;
const { debounceTime, takeWhile } = RxOperators;
Observable.create((observer) => {
let i = 1;
observer.next(i++);
// Emite un valor cada 1000 ms
setInterval(() => {
observer.next(i++)
}, 1000);
}).pipe(
takeWhile( producer => producer < 5 )
); 
Operadores — combineLatest
El operador combinado combineLatest se asemeja a promise.all. Combina múltiples flujos en uno. Una vez que cada flujo emite al menos una vez, obtenemos los últimos valores de cada uno como un array. Luego, después de cualquier emisión de los flujos combinados, comenzará a devolver nuevos valores.

const { combineLatest, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 1000 ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
});
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 750 ms
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
});
combineLatest(observer_1, observer_2).pipe(take(5));
Operadores — zip
Zip espera un valor de cada flujo y forma un array basado en esos valores. Si un valor no llega de algún flujo, el grupo no será formado.

const { zip, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 1000 ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
});
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 750 ms
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
});
const observer_3 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 500 ms
setInterval(() => {
observer.next('c: ' + i++);
}, 500);
});
zip(observer_1, observer_2, observer_3).pipe(take(5));
Operadores — forkJoin
forkJoin también combina flujos, pero emite un valor solo cuando todos los flujos se han completado.

const { forkJoin, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 1000 ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
}).pipe(take(3));
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 750 ms
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
}).pipe(take(5));
const observer_3 = Observable.create((observer) => {
let i = 1;
// Emite un valor cada 500 ms
setInterval(() => {
observer.next('c: ' + i++);
}, 500);
}).pipe(take(4));
forkJoin(observer_1, observer_2, observer_3);
Operadores — map
El operador de transformación map convierte el valor emitido en uno nuevo.

const { Observable } = Rx;
const { take, map } = RxOperators;
Observable.create((observer) => {
let i = 1;
// Emite un valor cada 1000 ms
setInterval(() => {
observer.next(i++);
}, 1000);
}).pipe(
map(x => x * 10),
take(3)
);

Operadores – share, tap
El operador tap permite realizar efectos secundarios, es decir, acciones que no afectan a la secuencia.
El operador utilitario share puede convertir un flujo frío en uno caliente.

Hemos terminado con los operadores. Pasemos al Subject.
Pensando en voz alta
Y aquí fui a tomarme un té. Me han cansado estos ejemplos 😀
Familia de subjects
La familia de subjects es un claro ejemplo de flujos calientes. Estas clases son una especie de híbrido que actúa simultáneamente como observable y observer. Dado que el subject es un flujo caliente, es necesario desuscribirse de él. Los métodos principales son:
- next – transmisión de nuevos datos al flujo
- error – error y finalización del flujo
- complete – finalización del flujo
- subscribe – suscribirse al flujo
- unsubscribe – desuscribirse del flujo
- asObservable – transformarlo en un observador
- toPromise – lo transforma en una promesa
Se identifican 4 o 5 tipos de subjects.
Pensando en voz alta
En la transmisión hablé de 4, y resultó que han agregado uno más. Como se dice, vive y aprende.
Subject simple new Subject()es el tipo más simple de subjects. Se crea sin parámetros. Solo transmite valores que llegan después de la suscripción.
BehaviorSubject new BehaviorSubject(defaultData) es, en mi opinión, el tipo de subject más común. Acepta un valor por defecto. Siempre guarda los datos del último emitido, que se transmite en la suscripción. Esta clase también tiene un método útil, value, que devuelve el valor actual del flujo.
ReplaySubject new ReplaySubject(bufferSize?: number, windowTime?: number) Puede aceptar opcionalmente como primer argumento el tamaño del búfer de valores que almacenará y como segundo el tiempo durante el cual necesitamos las modificaciones.
AsyncSubject new AsyncSubject() al suscribirse no ocurre nada y el valor se devolverá solo al completo. Solo se devolverá el último valor del flujo.
WebSocketSubject new WebSocketSubject(urlConfigOrSource: string | WebSocketSubjectConfig | Observable, destination?: Observer) la documentación no dice nada al respecto y yo lo veo por primera vez. Quien sepa qué hace, escriba y lo complementamos.
Uf. Bueno, hemos revisado todo lo que quería contar hoy. Espero que esta información haya sido útil. Puedes revisar la bibliografía en la pestaña de información útil.
Información útil
- — enlace a la transmisión
- — Cómo funciona JS: revisión del motor, mecanismos de ejecución, pila de llamadas
- — Cómo funciona JS: el ciclo de eventos, la asincronía y cinco formas de mejorar el código con async / await
- — Cómo funciona el Event Loop en JavaScript
- — La evolución de JavaScript asíncrono
- — Qué es RxJS y por qué es útil conocerlo
- — Aplicaciones prácticas de RxJS
- — Tutorial de RxJS Observables — Creación y suscripción a Observables
- — RXJS: Observables Fríos y Calientes
- — Clases, funciones para crear Observable. Operadores.
- — Operadores de RxJS por ejemplo
- — Lista de API
- — Variedades de Subject y Programaciones en RxJS
Fuente: habr.com
