Hallo zusammen. Hier ist Sergey Omelnitsky. Vor kurzem habe ich einen Livestream zum reaktiven Programmieren gehalten, in dem ich über Asynchronität in JavaScript gesprochen habe. Heute möchte ich dieses Material zusammenfassen.

Aber bevor wir mit dem Hauptmaterial beginnen, müssen wir eine Einführung machen. Lass uns also mit Definitionen beginnen: Was ist ein Stack und eine Queue?
Stack — das ist eine Sammlung, deren Elemente nach dem Prinzip „Last In, First Out“ (LIFO) verarbeitet werden.
Warteschlange — das ist eine Sammlung, deren Elemente nach dem Prinzip „First In, First Out“ (FIFO) verarbeitet werden.
Okay, lass uns fortfahren.

JavaScript ist eine einzeln laufende Programmiersprache. Das bedeutet, dass es nur einen Ausführungsstrang und einen Stack gibt, in den Funktionen zur Ausführung eingereiht werden. Folglich kann JavaScript zu einem bestimmten Zeitpunkt nur eine Operation ausführen, während andere Operationen im Stack warten müssen, bis sie aufgerufen werden.
Call Stack — das ist eine Datenstruktur, die vereinfacht gesagt Informationen darüber speichert, wo wir uns im Programm befinden. Wenn wir in eine Funktion eintreten, fügen wir einen Eintrag dafür an die Spitze des Stacks hinzu. Wenn wir aus der Funktion zurückkehren, ziehen wir das oberste Element aus dem Stack und gelangen dorthin, wo wir diese Funktion aufgerufen haben. Das ist alles, was der Stack kann. Und jetzt eine sehr interessante Frage: Wie funktioniert dann die Asynchronität in JavaScript?

Tatsächlich gibt es neben dem Stack in Browsern eine spezielle Queue für die Arbeit mit dem sogenannten WebAPI. Funktionen aus dieser Queue werden der Reihe nach erst nach vollständiger Leerung des Stacks ausgeführt. Erst dann werden sie von der Queue in den Stack zur Ausführung verschoben. Wenn sich zu diesem Zeitpunkt mindestens ein Element im Stack befindet, können sie nicht in den Stack gelangen. Genau wegen dieser Tatsache ist der Aufruf von Funktionen mit einem Timeout oft nicht zeitgenau, da die Funktion nicht in den Stack gelangen kann, solange dieser voll ist.
Betrachten wir das folgende Beispiel und schauen uns seine schrittweise „Ausführung“ an. Lass uns auch untersuchen, was dabei im System passiert.
console.log('Hi);
setTimeout(function cb1() {
console.log('cb1');
}, 5000);
console.log('Bye');
1) Bis jetzt passiert nichts. Die Konsole des Browsers ist leer, der Call Stack ist leer.

2) Dann wird der Befehl console.log(‘Hi’) in den Call Stack eingefügt.

3) Und er wird ausgeführt.

4) Dann wird console.log(‘Hi’) aus dem Call Stack entfernt.

5) Jetzt gehen wir zum Befehl setTimeout(function cb1() {… }) über. Er wird zum Aufrufstack hinzugefügt.

6) Der Befehl setTimeout(function cb1() {… }) wird ausgeführt. Der Browser erstellt einen Timer, der Teil der Web-API ist. Er wird den Countdown starten.

7) Der Befehl setTimeout(function cb1() {… }) hat seine Arbeit beendet und wird aus dem Aufrufstack entfernt.

8) Der Befehl console.log('Bye') wird zum Aufrufstack hinzugefügt.

9) Der Befehl console.log('Bye') wird ausgeführt.

10) Der Befehl console.log('Bye') wird aus dem Aufrufstack entfernt.

11) Nachdem mindestens 5000 ms vergangen sind, beendet der Timer seine Arbeit und fügt den Callback cb1 in die Callback-Warteschlange ein.

12) Der Event Loop nimmt die Funktion cb1 aus der Callback-Warteschlange und fügt sie in den Aufrufstack ein.

13) Die Funktion cb1 wird ausgeführt und fügt console.log('cb1') zum Aufrufstack hinzu.

14) Der Befehl console.log('cb1') wird ausgeführt.

15) Der Befehl console.log('cb1') wird aus dem Aufrufstack entfernt.

16) Die Funktion cb1 wird aus dem Aufrufstack entfernt.
Sehen wir uns ein Beispiel in der Dynamik an:

Nun haben wir gesehen, wie Asynchronität in JavaScript implementiert ist. Lassen Sie uns nun kurz über die Entwicklung asynchronen Codes sprechen.
Die Entwicklung asynchronen Codes.
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);
})
})
})
})
})
});Asynchrones Programmieren, wie wir es in JavaScript kennen, kann nur durch Funktionen umgesetzt werden. Diese können wie jede andere Variable an andere Funktionen übergeben werden. So entstanden Callbacks. Und das macht Spaß, solange es nicht in Traurigkeit und Kummer umschlägt. Warum? Weil:
- Mit wachsender Komplexität des Codes verwandelt sich das Projekt schnell in schwer verständliche mehrfach geschachtelte Blöcke – „Callback-Hölle“.
- Fehlerbehandlung kann leicht übersehen werden.
- Die Ausdrücke können nicht mit return zurückgegeben werden.
Mit dem Erscheinen von Promise wurde die Situation etwas besser.
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);
});- Es entstanden Ketten von Promises, was die Lesbarkeit des Codes verbesserte.
- Eine separate Methode zur Fehlerbehandlung wurde eingeführt.
- Es wurde die Möglichkeit der parallelen Ausführung mit Promise.all geschaffen.
- Wir können die eingebettete Asynchronität mit async/await lösen.
Aber ein Promise hat seine Beschränkungen. Zum Beispiel kann ein Promise ohne viel Aufwand nicht abgebrochen werden, und das Wichtigste ist — es funktioniert mit einem einzigen Wert.
Nun, wir sind also sanft zum reaktiven Programmieren übergegangen. Müde? Nun, zum Glück kann man sich einen Tee machen, nachdenken und dann weiterlesen. Ich fahre fort.

Reaktives Programmieren — eine Programmierparadigma, die sich auf Datenströme und das Verbreiten von Änderungen konzentriert. Lassen Sie uns genauer betrachten, was ein Datenstrom ist.
// Получаем ссылку на элемент
const input = ducument.querySelector('input');
const eventsArray = [];
// Пушим каждое событие в массив eventsArray
input.addEventListener('keyup',
event => eventsArray.push(event)
);Stellen wir uns vor, wir haben ein Eingabefeld. Wir erstellen ein Array und bei jedem keyup-Ereignis des Inputs werden wir das Ereignis in unserem Array speichern. Dabei sollte erwähnt werden, dass unser Array nach Zeit sortiert ist, d.h. der Index späterer Ereignisse ist größer als der Index früherer Ereignisse. Ein solches Array stellt ein vereinfachtes Modell eines Datenstroms dar, aber es ist noch kein Strom. Damit man dieses Array als Strom bezeichnen kann, muss es in der Lage sein, den Abonnenten mitzuteilen, dass neue Daten eingegeben wurden. So sind wir zur Definition eines Stroms gekommen.
Datenfluss
const { interval1 } = Rx;
const { take } = RxOperators;
interval(1000).pipe(
take(4)
)
Ein Strom — ist ein Array von Daten, das nach Zeit sortiert ist und mitteilen kann, dass sich die Daten geändert haben. Stellen Sie sich vor, wie bequem es wird, Code zu schreiben, bei dem eine Aktion mehrere Ereignisse an verschiedenen Stellen im Code erfordert. Wir abonnieren einfach den Strom, und er sagt uns selbst, wann Änderungen stattfinden. Und das kann die Bibliothek RxJs.

RxJS — ist eine Bibliothek zur Arbeit mit asynchronen und ereignisgesteuerten Programmen unter Verwendung von beobachtbaren Sequenzen. Die Bibliothek bietet den Haupttyp Observable, mehrere Hilfstypen (Observer, Schedulers, Subjects) und Operatoren zur Bearbeitung von Ereignissen wie Sammlungen (map, filter, reduce, every und ähnliches aus JavaScript Array).
Lassen Sie uns die grundlegenden Begriffe dieser Bibliothek klären.
Observable, Observer, Producer
Observable — der erste grundlegende Typ, den wir betrachten werden. Diese Klasse enthält den Hauptteil der Umsetzung von RxJs. Sie ist mit dem beobachtbaren Strom verbunden, auf den man sich mit der Methode subscribe anmelden kann.
In Observable wird ein Hilfsmechanismus zur Erstellung von Updates implementiert, der sogenannte Observer. Die Quelle der Werte für den Observer wird genannt Producer. Dies kann ein Array, ein Iterator, ein Websocket, ein Ereignis usw. sein. Man kann also sagen, dass observable ein Vermittler zwischen Producer und Observer ist.
Observable verarbeitet drei Arten von Ereignissen für den Observer:
- next – neue Daten
- error – ein Fehler, wenn die Sequenz aufgrund einer Ausnahme beendet wurde. Dieses Ereignis signalisiert ebenfalls das Ende der Sequenz.
- complete – ein Signal für das Ende der Sequenz. Das bedeutet, dass keine neuen Daten mehr verfügbar sind.
Sehen wir uns eine Demo an:

Zu Beginn verarbeiten wir die Werte 1, 2, 3, und nach 1 Sekunde erhalten wir 4 und beenden unseren Stream.
Gedanken laut
Und hier wurde mir klar, dass es interessanter war, darüber zu erzählen als es aufzuschreiben. 😀
Subscription
Wenn wir uns für einen Stream anmelden, erstellen wir eine neue Klasse subscription, die es uns ermöglicht, die Anmeldung mit der Methode unsubscribezu widerrufen. Wir können auch Abonnements mit der Methode addgruppieren. Und es ist logisch, dass wir Streams mit removeaufteilen können. Die Methoden add und remove nehmen eine andere Subscription als Eingabe. Ich möchte hervorheben, dass wir bei einer Abmeldung auch von allen untergeordneten Abonnements abgemeldet werden, als ob wir bei ihnen die Methode unsubscribe aufgerufen hätten. Lassen Sie uns fortfahren.
Arten von Streams
| HOT | COLD |
|---|---|
| Producer wird außerhalb von observable erstellt | Producer wird innerhalb von observable erstellt |
| Daten werden zum Zeitpunkt der Erstellung von observable übergeben | Daten werden zum Zeitpunkt der Anmeldung übermittelt |
| Zusätzliche Logik ist erforderlich, um sich abzumelden | Der Stream endet von selbst |
| Verwendet eine Eins-zu-viele-Beziehung | Verwendet eine Eins-zu-eins-Beziehung |
| Alle Abonnements haben denselben Wert | Abonnements sind unabhängig |
| Daten können verloren gehen, wenn kein Abonnement besteht | Gibt alle Werte des Streams für ein neues Abonnement erneut aus |
Wenn ich eine Analogie ziehen müsste, würde ich den heißen Stream mit einem Film im Kino vergleichen. Zu welchem Zeitpunkt du kommst, beginnt auch dein Sehen. Den kalten Stream würde ich mit einem Anruf beim technischen Support vergleichen. Jeder Anrufer hört die Ansage des Anrufbeantworters von Anfang bis Ende, aber du kannst das Gespräch mit der Methode unsubscribe beenden.
Ich möchte darauf hinweisen, dass es auch sogenannte warmen Ströme gibt (diese Definition habe ich sehr selten gesehen, und zwar nur in ausländischen Gemeinschaften) — es handelt sich um einen Strom, der sich von einem kalten Strom in einen heißen verwandelt. Die Frage stellt sich — wo zu verwenden)) Ich gebe ein Beispiel aus der Praxis.
Ich arbeite mit Angular. Es verwendet aktiv RxJS. Um Daten auf dem Server zu erhalten, erwarte ich einen kalten Strom und verwende diesen Strom im Template mit asyncPipe. Wenn ich dieses Pipe mehrmals verwende, wird, bezogen auf die Definition des kalten Stroms, jede Pipe Daten vom Server anfordern, was gelinde gesagt seltsam ist. Wenn ich jedoch den kalten Strom in einen warmen umwandle, erfolgt die Anforderung nur einmal.
Das Verständnis der Arten von Strömen ist für Anfänger ziemlich komplex, aber wichtig.
Operatoren
return this.http.get(`${environment.apiUrl}/${this.apiUrl}/trade_companies`)
.pipe(
tap(({ data }: TradeCompanyList) => this.companies$$.next(cloneDeep(data))),
map(({ data }: TradeCompanyList) => data)
);Operatoren bieten erweiterte Möglichkeiten zur Arbeit mit Strömen. Sie helfen, Ereignisse zu kontrollieren, die in Observable ablaufen. Wir werden ein paar der beliebtesten betrachten; detailliertere Informationen zu den Operatoren finden Sie in den Links zu nützlichen Informationen.
Operatoren — of
Beginnen wir mit dem Hilfsoperator of. Er erstellt ein Observable basierend auf einem einfachen Wert.

Operatoren — filter

Der Filteroperator filtert, wie man dem Namen entnehmen kann, das Signal des Stroms. Wenn der Operator true zurückgibt, wird das Signal weitergeleitet.
Operatoren — take

take — Nimmt die Anzahl der Emissionen, nach denen der Strom beendet wird.
Operatoren — debounceTime

debounceTime — verwirft emittierte Werte, die in den angegebenen Zeitrahmen zwischen den Ausgaben fallen. Nach Ablauf des Zeitintervalls wird der letzte Wert emittiert.
const { Observable } = Rx;
const { debounceTime, take } = RxOperators;
Observable.create((observer) => {
let i = 1;
observer.next(i++);
// Emit it every 1000ms
setInterval(() => {
observer.next(i++)
}, 1000);
// Emit it every 1500ms
setInterval(() => {
observer.next(i++)
}, 1500);
}).pipe(
debounceTime(700), // Warten Sie 700ms auf Werte, bevor Sie sie verarbeiten
take(3)
); 
Operatoren — takeWhile

Emittiert Werte, solange takeWhile nicht false zurückgibt, danach wird vom Strom abgemeldet.
const { Observable } = Rx;
const { debounceTime, takeWhile } = RxOperators;
Observable.create((observer) => {
let i = 1;
observer.next(i++);
// Emit value every 1000ms
setInterval(() => {
observer.next(i++)
}, 1000);
}).pipe(
takeWhile(producer => producer < 5)
); 
Operatoren — combineLatest
Der kombinierte Operator combineLatest ähnelt promise.all. Er kombiniert mehrere Streams in einem. Nachdem jeder Stream mindestens einmal einen Wert emittiert hat, erhalten wir die letzten Werte von jedem in Form eines Arrays. Nach jedem Emit eines der kombinierten Streams werden neue Werte zurückgegeben.

const { combineLatest, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emit value every 1000ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
});
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emit value every 750ms
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
});
combineLatest(observer_1, observer_2).pipe(take(5));
Operatoren — zip
Zip erwartet einen Wert aus jedem Stream und bildet ein Array auf der Grundlage dieser Werte. Wenn aus einem Stream kein Wert empfangen wird, wird die Gruppe nicht gebildet.

const { zip, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emit value every 1000ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
});
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emit value every 750ms
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
});
const observer_3 = Observable.create((observer) => {
let i = 1;
// Emit value every 500ms
setInterval(() => {
observer.next('c: ' + i++);
}, 500);
});
zip(observer_1, observer_2, observer_3).pipe(take(5));
Operatoren — forkJoin
forkJoin kombiniert ebenfalls Streams, jedoch gibt er einen Wert nur aus, wenn alle Streams abgeschlossen sind (complete).

const { forkJoin, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emit value every 1000ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
}).pipe(take(3));
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emit value every 750ms
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
}).pipe(take(5));
const observer_3 = Observable.create((observer) => {
let i = 1;
// Emit value every 500ms
setInterval(() => {
observer.next('c: ' + i++);
}, 500);
}).pipe(take(4));
forkJoin(observer_1, observer_2, observer_3);
Operatoren — map
Der Transformationsoperator map wandelt den Wert des Emittierten in einen neuen um.

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

Operatoren – share, tap
Der Tap-Operator ermöglicht Side-Effekte, das heißt, irgendwelche Aktionen, die die Reihenfolge nicht beeinflussen.
Der Utility-Operator Share kann aus einem kalten Stream einen heißen machen.

Wir sind mit den Operatoren fertig. Lassen Sie uns zum Subject übergehen.
Gedanken laut
Und hier habe ich eine Tasse Tee getrunken. Diese Beispiele haben mich müde gemacht 😀
Familie der Subjects
Die Familie der Subjects ist ein anschauliches Beispiel für heiße Streams. Diese Klassen sind eine Art Hybrid, die gleichzeitig als Observable und Observer fungieren. Da das Subject ein heißer Stream ist, muss man sich von ihm abmelden. Die Hauptmethoden sind:
- next – Übermittlung neuer Daten an den Stream
- error – Fehler und Beendigung des Streams
- complete – Beendigung des Streams
- subscribe – für den Stream anmelden
- unsubscribe – vom Stream abmelden
- asObservable – in einen Observer umwandeln
- toPromise – in ein Promise umwandeln
Es gibt 4 bis 5 Arten von Subjects.
Gedanken laut
Ich habe 4 in dem Stream gesagt, aber es scheint, dass sie noch einen hinzugefügt haben. Wie man so sagt: Lebe und lerne.
Einfaches Subject new Subject()– die einfachste Art von Subjects. Wird ohne Parameter erstellt. Gibt nur Werte weiter, die nach der Anmeldung empfangen wurden.
BehaviorSubject new BehaviorSubject(defaultData<T>) – meiner Meinung nach die am weitesten verbreitete Art von Subjects. Nimmt einen Standardwert als Eingabe an. Speichert immer die Daten des letzten Emits, die bei der Anmeldung übermittelt werden. Diese Klasse hat auch eine nützliche Methode value, die den aktuellen Wert des Streams zurückgibt.
ReplaySubject new ReplaySubject(bufferSize?: number, windowTime?: number) – Kann optional als ersten Parameter die Größe des Puffers von Werten annehmen, die er speichern wird, und als zweiten die Zeit, in der wir Änderungen benötigen.
AsyncSubject new AsyncSubject() – Bei der Anmeldung passiert nichts, und der Wert wird nur bei complete zurückgegeben. Nur der letzte Wert des Streams wird zurückgegeben.
WebSocketSubject new WebSocketSubject(urlConfigOrSource: string | WebSocketSubjectConfig<T> | Observable<T>, destination?: Observer<T>) – Darüber schweigt die Dokumentation, und ich sehe ihn zum ersten Mal. Wer weiß, was er macht, schreibt bitte, wir ergänzen das.
Puh. So, das haben wir alles besprochen, was ich heute erzählen wollte. Ich hoffe, diese Informationen waren nützlich. Man kann sich eigenständig mit der Literaturübersicht im Tab nützliche Informationen vertraut machen.
Nützliche Informationen
- – Link zum Stream
- – Wie JS funktioniert: Überblick über die Engine, Ausführungsmechanismen, Aufrufstapel
- — Wie JS funktioniert: Ereignisschleife, Asynchronität und fünf Möglichkeiten zur Verbesserung des Codes mit async / await
- — Wie die Ereignisschleife in JavaScript funktioniert
- — Evolution des asynchronen JavaScript
- — Was ist RxJS und warum ist es nützlich, darüber Bescheid zu wissen
- — Praktische Anwendung von RxJS
- — RxJS Observables Tutorial — Erstellen & Abonnieren von Observables
- — RXJS: Heiße und kalte Observables
- — Klassen, Funktionen zur Erstellung von Observable. Operatoren.
- — RxJS Operatoren anhand von Beispielen
- — API-Liste
- — Arten von Subject und Zeitplänen in RxJS
Quelle: habr.com
