
În trecut, Am discutat despre fundamentele teoretice ale arhitecturii reactive. A venit momentul să vorbim despre fluxurile de date, căile de implementare a sistemelor reactive Erlang/Elixir și despre modelele de schimb de mesaje în acestea:
- Cerere-răspuns
- Cerere-Răspuns Chunked
- Răspuns cu Cerere
- Publicare-abonare
- Publicare-abonare inversată
- Distribuția sarcinilor
SOA, MSA și schimbul de mesaje
SOA, MSA – arhitecturi sistemice care definesc regulile de construire a sistemelor, în timp ce messaging oferă primitive pentru implementarea acestora.
Nu vreau să promovez o arhitectură sau alta pentru construirea sistemelor. Susțin aplicarea celor mai eficiente și utile practici pentru proiectul și afacerea specifică. Indiferent de paradigma aleasă, este mai bine să creăm blocuri sistemice cu referire la principiul Unix: componente cu legături minime, responsabile pentru entități separate. Metodele API efectuează acțiuni cât mai simple cu entitățile.
Messaging ‒ după cum sugerează numele ‒ este un broker de mesaje. Principalul său scop este de a primi și transmite mesaje. El se ocupă de interfețele de expediere a informațiilor, de formarea canalelor logice de transmitere a informațiilor în cadrul sistemului, de rutare și echilibrare, precum și de gestionarea defectelor la nivel sistemic.
Messaging-ul dezvoltat nu încearcă să concureze cu rabbitmq sau să-l înlocuiască. Principalele sale caracteristici sunt:
- Distribuția.
Punctele de schimb pot fi create pe toate nodurile clusterei, cât mai aproape de codul care le utilizează. - Simplitate.
Orientare către minimizarea codului standard și ușurința de utilizare. - Performanță superioară.
Nu încercăm să replicăm funcționalitatea rabbitmq, ci ne concentrăm doar pe stratul arhitectural și de transport, pe care îl integrăm cât mai simplu în OTP, minimizând costurile. - Flexibilitatea.
Fiecare serviciu poate combina numeroase modele de schimb. - Resiliența, înrădăcinată în design.
- Scalabilitate.
Messaging-ul crește împreună cu aplicația. Pe măsură ce sarcina crește, punctele de schimb pot fi mutate pe mașini separate.
Observație. Din perspectiva organizării codului, pentru sisteme complexe în Erlang/Elixir, proiectele meta sunt foarte potrivite. Tot codul proiectului se află într-un singur depozit - un proiect umbrelă. În acest context, microserviciile sunt cât mai izolate posibil și efectuează operații simple, fiecare fiind responsabilă pentru o entitate separată. Cu un astfel de abordare, este ușor de menținut API-ul întregului sistem, modificările sunt simple de realizat, iar scrierea testelor unitare și de integrare este convenabilă.
Componentele sistemului interacționează direct sau prin intermediul unui broker. Din perspectiva messaging-ului, fiecare serviciu are mai multe faze de viață:
- Inițializarea serviciului.
În această etapă se face configurarea și pornirea procesului care execută serviciul și a dependențelor. - Crearea unui punct de schimb.
Serviciul poate folosi un punct de schimb static, definit în configurația nodului, sau poate crea puncte de schimb dinamic. - Înregistrarea serviciului.
Pentru ca serviciul să poată gestiona solicitările, trebuie înregistrat la punctul de schimb. - Funcționarea normală.
Serviciul îndeplinește o muncă utilă. - Finalizarea muncii.
Există 2 tipuri de finalizare a activității: normală și de urgență. În cazul normal, serviciul se deconectează de la punctul de schimb și se oprește. În cazuri de urgență, messaging-ul execută unul dintre scenariile de gestionare a defecțiunilor.
Pare destul de complex, dar în cod nu este totul atât de înfricoșător. Exemple de cod cu comentarii vor fi prezentate în analiza șabloanelor puțin mai târziu.
Exchanges
Punctul de schimb este un proces de messaging care implementează logica interacțiunii cu componentele în cadrul șablonului de schimb de mesaje. În toate exemplele prezentate mai jos, componentele interacționează prin puncte de schimb, combinația cărora formează messaging-ul.
Modelele de schimb de mesaje (MEPs)
În mod global, modelele de schimb pot fi împărțite în bidirecționale și unidirecționale. Primele implică un răspuns la mesajul primit, cele din urmă nu. Un exemplu clasic de model bidirecțional în arhitectura client-server este modelul Request-response. Să examinăm modelul și modificările sale.
Request–response sau RPC
RPC este utilizat atunci când avem nevoie să obținem un răspuns de la un alt proces. Acest proces poate fi pornit pe același nod sau poate fi situat pe un alt continent. Mai jos este prezentată schema de interacțiune între client și server prin messaging.

Deoarece messaging-ul este complet asincron, pentru client schimbul este împărțit în 2 faze:
Trimiterea unei solicitări
messaging:request(Exchange, ResponseMatchingTag, RequestDefinition, HandlerProcess).Exchange ‒ numele unic al punctului de schimb
ResponseMatchingTag ‒ eticheta locală pentru procesarea răspunsului. De exemplu, în cazul trimiterii mai multor cereri identice, care aparțin unor utilizatori diferiți.
RequestDefinition ‒ corpul cererii
HandlerProcess ‒ PID-ul handler-ului. Acest proces va primi răspunsul de la server.Procesarea răspunsului
handle_info(#'$msg'{exchange = EXCHANGE, tag = ResponseMatchingTag, message = ResponsePayload}, State)ResponsePayload ‒ răspunsul serverului.
Pentru server, procesul constă de asemenea în 2 faze:
- Inițializarea punctului de schimb
- Procesarea cererilor primite
Vom ilustra acest șablon cu cod. Să presupunem că trebuie să implementăm un serviciu simplu care oferă o singură metodă pentru ore exacte.
Codul serverului
Vom extrage definirea API-ului serviciului în api.hrl:
%% =====================================================
%% entități
%% =====================================================
-record(time, {
unixtime :: non_neg_integer(),
datetime :: binary()
}).
-record(time_error, {
code :: non_neg_integer(),
error :: term()
}).
%% =====================================================
%% metode
%% =====================================================
-record(time_req, {
opts :: term()
}).
-record(time_resp, {
result :: #time{} | #time_error{}
}).Vom defini controlerul serviciului în time_controller.erl
%% Exemplul arată doar codul relevant. Introducându-l în șablonul gen_server se poate obține un serviciu funcțional.
%% inițializarea gen_server
init(Args) ->
%% conectare la punctul de schimb
messaging:monitor_exchange(req_resp, ?EXCHANGE, default, self())
{ok, #{}}.
%% procesarea evenimentului de pierdere a conexiunii cu punctul de schimb. Acest eveniment este primit și dacă punctul de schimb nu a început încă.
handle_info(#exchange_die{exchange = ?EXCHANGE}, State) ->
erlang:send(self(), monitor_exchange),
{noreply, State};
%% procesarea API-ului
handle_info(#time_req{opts = _Opts}, State) ->
messaging:response_once(Client, #time_resp{
result = #time{ unixtime = time_utils:unixtime(now()), datetime = time_utils:iso8601_fmt(now())}
});
{noreply, State};
%% finalizarea activității gen_server
terminate(_Reason, _State) ->
messaging:demonitor_exchange(req_resp, ?EXCHANGE, default, self()),
ok.Codul clientului
Pentru a trimite o cerere serviciului, în orice loc al clientului se poate apela API-ul de cerere messaging:
case messaging:request(?EXCHANGE, tag, #time_req{opts = #{}}, self()) of
ok -> ok;
_ -> %% logică de repetare sau eroare
endÎntr-un sistem distribuit, configurația componentelor poate fi foarte variată și în momentul cererii messaging poate să nu fi pornit, sau controlerul serviciului nu va fi pregătit să proceseze cererea. De aceea, este necesar să verificăm răspunsul messaging și să gestionăm cazul de eșec.
După trimiterea cu succes, clientul va primi un răspuns sau o eroare din partea serviciului.
Să procesăm ambele cazuri în handle_info:
handle_info(#'$msg'{exchange = ?EXCHANGE, tag = tag, message = #time_resp{result = #time{unixtime = Utime}}}, State) ->
?debugVal(Utime),
{noreply, State};
handle_info(#'$msg'{exchange = ?EXCHANGE, tag = tag, message = #time_resp{result = #time_error{code = ErrorCode}}}, State) ->
?debugVal({error, ErrorCode}),
{noreply, State};Cerere-Răspuns Chunked
Este mai bine să nu permiteți transmiterea de mesaje mari. Acest lucru afectează răspunsul și funcționarea stabilă a întregului sistem. Dacă răspunsul la o solicitare consumă multă memorie, atunci fragmentarea pe părți este obligatorie.

Voi oferi câteva exemple de astfel de cazuri:
- Componentele schimbă date binare, cum ar fi fișierele. Fragmentarea răspunsului în părți mici ajută la gestionarea eficientă a fișierelor de orice dimensiune și la evitarea depășirilor de memorie.
- Listări. De exemplu, trebuie să selectăm toate înregistrările dintr-un tabel imens din bază și să le transmitem unei alte componente.
Numesc astfel de răspunsuri tren. În orice caz, 1024 de mesaje de 1 MB sunt mai bune decât un singur mesaj de 1 GB.
În clusterul Erlang, obținem un avantaj suplimentar - reducerea sarcinii pe punctul de schimb și pe rețea, deoarece răspunsurile sunt direcționate imediat către destinatar, ocolind punctul de schimb.
Răspuns cu Cerere
Aceasta este o modificare destul de rară a modelului RPC pentru construirea sistemelor de dialog.

Publicare-abonare (copac de distribuție a datelor)
Sistemele orientate pe evenimente livrează datele consumatorilor pe măsură ce devin disponibile. Astfel, sistemele sunt mai înclinate spre modelul push decât spre pull sau poll. Această caracteristică permite evitarea risipei de resurse prin a solicita constant și a aștepta date.
În imagine este prezentat procesul de distribuție a mesajelor către consumatorii abonați la un anumit subiect.

Exemple clasice de utilizare a acestui șablon sunt distribuția stării: lumea jocurilor video, datele de pe piețe, informații utile în fluxuri de date.
Să analizăm codul abonatului:
init(_Args) ->
%% ne abonăm la punctul de schimb, cheie = key
messaging:subscribe(?SUBSCRIPTION, key, tag, self()),
{ok, #{}}.
handle_info(#exchange_die{exchange = ?SUBSCRIPTION}, State) ->
%% dacă punctul de schimb nu este disponibil, încercăm să ne reconectăm
messaging:subscribe(?SUBSCRIPTION, key, tag, self()),
{noreply, State};
%% preluăm mesajele primite
handle_info(#'$msg'{exchange = ?SUBSCRIPTION, message = Msg}, State) ->
?debugVal(Msg),
{noreply, State};
%% la oprirea consumatorului - ne deconectăm de la punctul de schimb
terminate(_Reason, _State) ->
messaging:unsubscribe(?SUBSCRIPTION, key, tag, self()),
ok.Sursa poate apela funcția de publicare a mesajelor în orice loc convenabil:
messaging:publish_message(Exchange, Key, Message).Exchange ‒ numele punctului de schimb,
Cheie ‒ cheia de rutare,
Mesaj ‒ încărcătura utilă.
Publicare-abonare inversată

După implementarea pub-sub, se poate obține un model confortabil pentru logare. Setul de surse și consumatori poate fi complet diferit. În imagine este prezentat un caz cu un singur consumator și multe surse.
Model de distribuție a sarcinilor
În aproape fiecare proiect apar sarcini de procesare întârziate, cum ar fi generarea de rapoarte, livrarea de notificări, obținerea de date din sisteme externe. Capacitatea sistemului care execută aceste sarcini se scalează cu ușurință prin adăugarea de procesatori. Tot ce ne rămâne de făcut este să formăm un cluster de procesatori și să distribuim sarcinile uniform între ei.
Să examinăm situațiile apărute prin exemplul a 3 procesatori. Chiar și în etapa de distribuție a sarcinilor apare întrebarea echității distribuției și a supraaglomerării procesatorilor. Răspunderea pentru echitate va fi asigurată de distribuția round-robin, iar pentru a preveni situațiile de supraaglomerare a procesatorilor, vom introduce o limită, prefetch_limit. În regimuri tranzitorii, prefetch_limit nu va permite unui procesator să primească toate sarcinile.
Messaging gestionează cozile și prioritatea de procesare. Procesatorii primesc sarcinile pe măsură ce sosesc. Executarea unei sarcini poate fi finalizată cu succes sau poate fi refuzată:
messaging:ack(Tack)‒ este apelat în cazul procesării reușite a mesajului;messaging:nack(Tack)‒ este apelat în toate situațiile neprevăzute. După returnarea sarcinii, messaging o va transmite unui alt procesator.

Să presupunem că, în timpul procesării a trei sarcini, a avut loc o eroare complexă: procesatorul 1 a picat după primirea sarcinii, fără a reuși să comunice ceva punctului de schimb. În acest caz, punctul de schimb, după expirarea timeout-ului ack, va transmite sarcina unui alt procesator. Procesatorul 3, dintr-un anumit motiv, a refuzat sarcina și a trimis nack; în cele din urmă, sarcina a trecut și la alt procesator care a executat-o cu succes.
Rezumat preliminar
Am trecut în revistă principalele elemente ale sistemelor distribuite și am obținut o înțelegere de bază a aplicării acestora în Erlang/Elixir.
Combinând modelele de bază, se pot construi paradigme complexe pentru a rezolva sarcinile apărute.
În partea finală a ciclului, vom explora întrebările generale privind organizarea serviciilor, rutarea și balansarea, precum și vom discuta partea practică a scalabilității și rezilienței sistemelor.
Sfârșitul părții a doua.
Fotografie
Ilustrațiile au fost create cu ajutorul websequencediagrams.com
Sursa: habr.com
