Salut tuturor. Sunt Oleinik Sergey. Recent am avut un stream despre programarea reactivă, unde am vorbit despre asincronismul din JavaScript. Astăzi aș dori să rezum acest material.

Dar înainte de a începe materialul principal, trebuie să facem o introducere. Așadar, să începem cu definițiile: ce este un stivă și o coadă?
Stivă — este o colecție ale cărei elemente sunt obținute pe principiul „ultimul intrat, primul ieșit” LIFO
Coadă — este o colecție ale cărei elemente sunt obținute pe principiul „primul intrat, primul ieșit” FIFO
Bine, să continuăm.

JavaScript este un limbaj de programare pe un singur fir de execuție. Aceasta înseamnă că are un singur fir de execuție și o singură stivă, în care funcțiile sunt așezate în așteptare pentru execuție. Prin urmare, într-un singur moment, JavaScript poate executa doar o singură operație, iar celelalte operații așteaptă să fie procesate în stivă până când sunt apelate.
Stiva de apeluri — este o structură de date care, simplificat, înregistrează informații despre locul din program unde ne aflăm. Dacă intrăm într-o funcție, adăugăm un record despre aceasta în partea de sus a stivei. Când ne întoarcem din funcție, scoatem din stivă cel mai de sus element și ne regăsim acolo de unde am apelat acea funcție. Aceasta este tot ce poate face o stivă. Dar acum apare o întrebare foarte interesantă. Cum funcționează atunci asincronismul în JavaScript?

De fapt, pe lângă stivă, în browsere există o coadă specială pentru lucrul cu așa-numitul WebAPI. Funcțiile din această coadă vor fi executate în ordinea în care sunt adăugate doar după ce stiva este complet eliberată. Numai după aceea sunt mutate din coadă în stivă pentru execuție. Dacă în stivă există măcar un element, acestea nu pot să intre în stivă. Din acest motiv, apelurile funcțiilor cu timeout sunt adesea inexacte ca timp, deoarece funcția nu poate intra din coadă în stivă, cât timp aceasta este plină.
Să luăm următorul exemplu și să ne ocupăm de „execuția” lui pas cu pas. De asemenea, să vedem ce se întâmplă în sistem în același timp.
console.log('Hi);
setTimeout(function cb1() {
console.log('cb1');
}, 5000);
console.log('Bye');
1) Deocamdată nu se întâmplă nimic. Consola browserului este goală, stiva de apeluri este liberă.

2) Apoi comanda console.log('Hi') este adăugată în stiva de apeluri.

3) Și aceasta se execută

4) Apoi, console.log('Hi') este eliminat din stiva de apeluri.

5) Acum trecem la comanda setTimeout(function cb1() {… }). Aceasta este adăugată în stiva de apeluri.

6) Comanda setTimeout(function cb1() {… }) se execută. Browserul creează un timer, care face parte din Web API. Acesta va începe un contor înapoi.

7) Comanda setTimeout(function cb1() {… }) a terminat și este eliminată din stiva de apeluri.

8) Comanda console.log('Bye') este adăugată în stiva de apeluri.

9) Comanda console.log('Bye') se execută.

10) Comanda console.log('Bye') este eliminată din stiva de apeluri.

11) După ce au trecut cel puțin 5000 ms, timerul terminatează și plasează callback-ul cb1 în coada de callback-uri.

12) Ciclul de evenimente ia funcția cb1 din coada de callback-uri și o plasează în stiva de apeluri.

13) Funcția cb1 se execută și adaugă console.log('cb1') în stiva de apeluri.

14) Comanda console.log('cb1') se execută.

15) Comanda console.log('cb1') este eliminată din stiva de apeluri.

16) Funcția cb1 este eliminată din stiva de apeluri.
Să ne uităm la un exemplu în dinamică:

Iată, am discutat despre cum asynchronicitatea este implementată în JavaScript. Acum să vorbim pe scurt despre evoluția codului asincron.
Evoluția codului asincron.
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);
})
})
})
})
})
});Programarea asincronă, așa cum o știm în JavaScript, poate fi realizată doar prin funcții. Acestea pot fi transmise ca orice altă variabilă altor funcții. Astfel au apărut callback-urile. Și este grozav, distractiv și plin de viață, până când devine trist, melancolic și deprimant. De ce? E simplu:
- Pe măsură ce complexitatea codului crește, proiectul se transformă rapid în blocuri de cod greu de înțeles, mult învelite - «callback hell».
- Gestionarea erorilor poate fi ușor omisă.
- Nu poți returna expresii cu return.
Odată cu apariția Promise, situația s-a îmbunătățit puțin.
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);
});- Au apărut lanțuri de promisiuni, ceea ce a îmbunătățit lizibilitatea codului.
- A apărut o metodă separată pentru interceptarea erorilor.
- A apărut posibilitatea executării paralele cu ajutorul Promise.all
- Putem rezolva asincronicitatea imbricată cu ajutorul async/await
Dar promisiunea are limitările sale. De exemplu, o promisiune nu poate fi anulată fără ceva efort, iar, ceea ce e cel mai important – lucrează cu o singură valoare.
Ei bine, iată că ne-am apropiat ușor de programarea reactivă. Ești obosit? Dar, în bună măsură, poți să mergi să îți prepari o cană de ceai, să te gândești și să te întorci să continui citirea. Eu voi continua.

Programarea reactivă– o paradigmă de programare orientată pe fluxurile de date și propagarea schimbărilor. Să analizăm mai în detaliu ce înseamnă fluxul de date.
// Получаем ссылку на элемент
const input = ducument.querySelector('input');
const eventsArray = [];
// Пушим каждое событие в массив eventsArray
input.addEventListener('keyup',
event => eventsArray.push(event)
);Să presupunem că avem un câmp de input. Creăm un array și la fiecare eveniment de keyup al input-ului vom salva evenimentul în array-ul nostru. De asemenea, ar fi de notat că array-ul nostru este sortat cronologic, adică indexul evenimentelor mai recente este mai mare decât indexul celor mai vechi. Acest array reprezintă un model simplificat al unui flux de date, dar este încă departe de a fi un flux. Pentru ca acest array să poată fi numit un flux, el trebuie să fie capabil să notifice abonații atunci când primesc date noi. Astfel, am ajuns la definiția fluxului.
Flux de date
const { interval1 } = Rx;
const { take } = RxOperators;
interval(1000).pipe(
take(4)
)
Fluxul– este un array de date, sortat cronologic, care poate comunica când datele s-au schimbat. Acum imaginează-ți cât de convenabil devine să scrii cod în care pentru o singură acțiune trebuie să apelezi mai multe evenimente în diferite locuri din cod. Pur și simplu ne abonăm la flux și acesta ne va informa singur când vor avea loc schimbări. Și asta poate face biblioteca RxJs.

RxJS – este o bibliotecă pentru lucrul cu programe asincrone și bazate pe evenimente care folosesc secvențe observabile. Biblioteca oferă tipul principal Observable, mai multe tipuri auxiliare (Observer, Schedulers, Subjects) și operatori pentru a lucra cu evenimente ca și cu colecții (map, filter, reduce, every și similare din JavaScript Array).
Să ne familiarizăm cu conceptele de bază ale acestei biblioteci.
Observable, Observer, Producer
Observable — primul tip de bază pe care îl vom examina. Această clasă conține cea mai mare parte a implementării RxJs. Este legată de un flux observabil, la care se poate abona prin metoda subscribe.
În Observable se implementează un mecanism auxiliar pentru a crea actualizări, numit Observer. Sursa valorilor pentru Observer este numită Producer. Aceasta poate fi un array, un iterator, un socket web, un anumit eveniment etc. Așadar, se poate spune că observable este un mediator între Producer și Observer.
Observable procesează trei tipuri de evenimente pentru Observer:
- next – date noi
- error – o eroare, dacă secvența s-a încheiat din cauza unei situații excepționale. Acest eveniment presupune de asemenea terminarea secvenței.
- complete — un semnal de finalizare a secvenței. Asta înseamnă că nu vor mai fi date noi.
Să vedem o demonstrație:

La început, vom procesa valorile 1, 2, 3, iar după 1 secundă, vom obține 4 și vom încheia fluxul nostru.
Gânduri cu voce tare
Și atunci mi-am dat seama că era mai interesant să povestesc decât să scriu despre asta. 😀
Subscription
Când ne abonați la flux, creăm o nouă clasă de subscription, care ne oferă posibilitatea de a anula abonarea prin metoda unsubscribe. De asemenea, putem grupa abonarile prin metoda add. Este logic că putem dezgrupa fluxurile prin remove. Metodele add și remove primesc ca intrare o altă subscriere. Aș dori să subliniez că atunci când ne dezabonăm, ne dezabonăm de la toate subscrierile copil, ca și cum am fi apelat metoda unsubscribe pentru ele. Să continuăm.
Tipuri de fluxuri
| HOT | COLD |
|---|---|
| Producer este creat din afara observable | Producer este creat în interiorul observable |
| Datele sunt transmise în momentul creării observable | Datele sunt comunicate în momentul abonării |
| Este necesară o logică suplimentară pentru dezabonare | Fluxul se încheie de la sine |
| Utilizează o relație de tip unul-la-mulți | Utilizează o relație de tip unul-la-unul |
| Toate abonările au o valoare comună | Abonările sunt independente |
| Datele pot fi pierdute dacă nu există abonare | Reemite toate valorile fluxului pentru o nouă abonare |
Dacă ar fi să fac o analogie, aș imagina un flux cald ca un film într-un cinematograf. În momentul în care ai venit, aceea e clipa în care începi vizionarea. Un flux rece l-aș compara cu un apel la asistența tehnică. Oricine sună ascultă înregistrarea unui răspuns automat de la început până la sfârșit, dar poți renunța la apel folosind unsubscribe.
Aș dori să menționez că există așa-numitele fluxuri calde (această definiție am întâlnit-o foarte rar și doar în comunități externe) — este un flux care se transformă dintr-un flux rece într-un flux cald. Se ridică întrebarea — unde ar trebui să le folosim)) Voi da un exemplu din practică.
Lucrez cu Angular. Acesta folosește activ rxjs. Pentru a obține date de la server, aștept un flux rece și folosesc acest flux în șablon cu ajutorul asyncPipe. Dacă folosesc acest pipe de mai multe ori, atunci, referindu-mă la definiția fluxului rece, fiecare pipe va solicita date de la server, ceea ce, pe scurt, este bizar. Dar dacă transform un flux rece într-unul cald, atunci cererea se va face o singură dată.
De fapt, înțelegerea formelor de fluxuri este destul de complicată pentru începători, dar este importantă.
Operatori
return this.http.get(`${environment.apiUrl}/${this.apiUrl}/trade_companies`)
.pipe(
tap(({ data }: TradeCompanyList) => this.companies$$.next(cloneDeep(data))),
map(({ data }: TradeCompanyList) => data)
);Oportunitățile de a lucra cu fluxurile ne sunt oferite de operatori. Aceștia ajută la controlul evenimentelor care se desfășoară în Observable. Vom analiza câțiva dintre cei mai populari, iar informații mai detaliate despre operatori puteți găsi în linkurile din informațiile utile.
Operatori — of
Să începem cu operatorul auxiliar of. Acesta creează un Observable pe baza unei valori simple.

Operatori — filter

Operatorul de filtrare filter, așa cum se poate înțelege din nume, filtrează semnalul fluxului. Dacă operatorul returnează adevărat, atunci trece mai departe.
Operatori — take

take — Primește valoarea numărului de emisii, după care încheie fluxul.
Operatori — debounceTime

debounceTime — elimină valorile emise care apar în intervalul specificat de timp între ieșirile date — după expirarea intervalului de timp emite ultima valoare.
const { Observable } = Rx;
const { debounceTime, take } = RxOperators;
Observable.create((observer) => {
let i = 1;
observer.next(i++);
// Emit value every 1000ms
setInterval(() => {
observer.next(i++)
}, 1000);
// Emit value every 1500ms
setInterval(() => {
observer.next(i++)
}, 1500);
}).pipe(
debounceTime(700), // Wait for 700ms before processing values
take(3)
); 
Operators — takeWhile

Emite valori până când takeWhile returnează false, după care se dezabonează de la flux.
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)
); 
Operators — combineLatest
Operatorul combinat combineLatest este similar cu promise.all. Acesta combină mai multe fluxuri într-unul. După ce fiecare flux emite măcar o dată, obținem ultimele valori de la fiecare sub formă de masiv. Ulterior, după orice emitere din fluxurile combinate, vor fi furnizate noi valori.

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));
Operators — zip
Zip așteaptă valori din fiecare flux și formează un masiv pe baza acestor valori. Dacă o valoare nu vine dintr-un flux, grupul nu va fi format.

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 750
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
});
const observer_3 = Observable.create((observer) => {
let i = 1;
// Emit value every 500
setInterval(() => {
observer.next('c: ' + i++);
}, 500);
});
zip(observer_1, observer_2, observer_3).pipe(take(5));
Operators — forkJoin
forkJoin combină de asemenea fluxurile, dar emite o valoare doar atunci când toate fluxurile sunt completate.

const { forkJoin, Observable } = Rx;
const { take } = RxOperators;
const observer_1 = Observable.create((observer) => {
let i = 1;
// Emitem o valoare la fiecare 1000ms
setInterval(() => {
observer.next('a: ' + i++);
}, 1000);
}).pipe(take(3));
const observer_2 = Observable.create((observer) => {
let i = 1;
// Emitem o valoare la fiecare 750
setInterval(() => {
observer.next('b: ' + i++);
}, 750);
}).pipe(take(5));
const observer_3 = Observable.create((observer) => {
let i = 1;
// Emitem o valoare la fiecare 500
setInterval(() => {
observer.next('c: ' + i++);
}, 500);
}).pipe(take(4));
forkJoin(observer_1, observer_2, observer_3);
Operatorii — map
Operatorul de transformare map converteste valoarea emisiunii într-o nouă valoare.

const { Observable } = Rx;
const { take, map } = RxOperators;
Observable.create((observer) => {
let i = 1;
// Emitem o valoare la fiecare 1000ms
setInterval(() => {
observer.next(i++);
}, 1000);
}).pipe(
map(x => x * 10),
take(3)
);

Operatorii – share, tap
Operatorul tap permite efecte secundare, adică orice acțiuni care nu afectează secvența.
Operatorul utilitar share poate transforma un flux rece într-unul cald.

Am terminat cu operatorii. Să trecem la Subject.
Gânduri cu voce tare
Și aici am plecat să beau ceai. M-au obosit aceste exemple 😀
Familia subject-urilor
Familia subject-urilor reprezintă un exemplu izbitor de fluxuri calde. Aceste clase reprezintă un hibrid, care funcționează simultan ca observable și observer. Deoarece subject este un flux cald, trebuie să te dezabonezi de la el. Vorbind despre metodele principale, acestea sunt:
- next – transmiterea de noi date în flux
- error – eroare și încheierea fluxului
- complete – încheierea fluxului
- subscribe – a te abona la flux
- unsubscribe – a te dezabona de la flux
- asObservable – transformăm într-un observer
- toPromise – transformă într-un promisiune
Se disting 4-5 tipuri de subject-uri.
Gânduri cu voce tare
În stream am spus 4, dar se pare că au adăugat încă unul. Așa cum se spune, trăiește și învață.
Subject simplu new Subject()– cea mai simplă formă de subject-uri. Se creează fără parametri. Transmite valori primite doar după abonare.
BehaviorSubject new BehaviorSubject(defaultData) – din punctul meu de vedere, cel mai comun tip de subject-uri. Acceptă o valoare implicită. Păstrează întotdeauna datele ultimei emisii, pe care le transmite la abonare. Această clasă are de asemenea metode utile, cum ar fi value, care returnează valoarea curentă a fluxului.
ReplaySubject new ReplaySubject(bufferSize?: number, windowTime?: number) – poate accepta opțional ca prim argument dimensiunea buffer-ului valorilor pe care le va reține și ca al doilea argument timpul în care avem nevoie de modificări.
AsyncSubject new AsyncSubject() — la înscriere nimic nu se întâmplă, iar valoarea va fi returnată doar la complete. Va fi returnată doar ultima valoare din flux.
WebSocketSubject new WebSocketSubject(urlConfigOrSource: string | WebSocketSubjectConfig | Observable, destination?: Observer) — Despre el documentația nu spune nimic și eu îl văd pentru prima dată. Cine știe ce face, scrieți-ne, să completăm.
Uf, iată că am acoperit tot ce voiam să discut astăzi. Sper că informațiile de mai sus au fost utile. Poți consulta singur lista de literatură în secțiunea de informații utile.
Informații utile
- — link către stream
- — Cum funcționează JS: o privire de ansamblu asupra motorului, mecanismelor de execuție, stivei de apeluri
- — Cum funcționează JS: ciclul de evenimente, asincronicitatea și cinci metode de îmbunătățire a codului cu async / await
- — Cum funcționează Event Loop în JavaScript
- — Evoluția JavaScript-ului asincron
- — Ce este RxJS și de ce este util să știm despre el
- — Aplicarea practică a RxJS
- — RxJS Observables Tutorial — Crearea și abonarea la Observables
- — RXJS: Observables fierbinți și reci
- — Clase, funcții pentru crearea Observable. Operatori.
- — Operatori RxJS prin exemple
- — Lista API
- — Tipuri de Subject și Programări în RxJS
Sursa: habr.com
