
Në të kaluarën Kemi analizuar teoritë themelore të arkitekturës reaktive. Ka ardhur koha të flasim për gratë e të dhënave, rrugët e realizimit të sistemeve reaktive Erlang/Elixir dhe modelet e shkëmbimit të mesazheve në to:
- Request-response
- Request-Chunked Response
- Response with Request
- Publish-subscribe
- Inverted Publish-subscribe
- Task distribution
SOA, MSA dhe shkëmbimi i mesazheve
SOA, MSA janë arkitekturat sistemore që përcaktojnë rregullat e ndërtimit të sistemeve, ndërsa messaging ofron primitive për realizimin e tyre.
Nuk dëshiroj të propagandoj asnjë arkitekturë në ndërtimin e sistemeve. Unë jam për aplikimin e praktikave më të efektshme dhe të dobishme për projektin dhe biznesin specifik. Pa marrë parasysh paradigmat që zgjedhim, është më mirë të krijojmë blloqet sistemore duke u mbështetur në Unix-way: komponimet me lidhje minime, që përgjigjen për entitetet e veçanta. Metodat API realizojnë veprime maksimalisht të thjeshta me entitetet.
Messaging ‒ siç kuptohet nga emri ‒ është brokeri i mesazheve. Qëllimi i tij kryesor është të pranojë dhe të dërgojë mesazhe. Ai është përgjegjës për ndërfaqet e dërgimit të informacionit, formimin e kanaleve logjikë të transmetimit të informacionit brenda sistemit, rrugëzimin dhe balancimin, si dhe trajtimin e dështimeve në nivelin sistemor.
Messaging që po zhvillohet nuk përpiqet të konkurrojë me rabbitmq ose ta zëvendësojë atë. Veçoritë e tij kryesore janë:
- Distribuese.
Pikat e shkëmbimit mund të krijohen në të gjitha nyjat e klasterit, sa më afër me kodin që i përdor ato. - Thjeshtësia.
Orientimi drejt minimizimit të kodit model dhe lehtësisë në përdorim. - Performanca më e mirë.
Ne nuk përpiqemi të përsërisim funksionalitetin e rabbitmq, por përqendrohemi vetëm në nivelin arkitekturor dhe të transportit, që të integrohet sa më thjesht në OTP, duke minimizuar kostot. - Fleksibiliteti.
Çdo shërbim mund të përfshijë shumë modele shkëmbimi. - Qëndrueshmëri në rast dështimi, e integruar në dizajn.
- Shkallëzimi.
Messaging rritet së bashku me aplikacionin. Me rritjen e ngarkesës, mund të transferohen pikat e shkëmbimit në makina të veçanta.
Vërejtje. Nga pikëpamja e organizatës së kodit, për sistemet komplekse në Erlang/Elixir metaprojektet janë shumë të përshtatshme. E gjithë kodi i projektit ndodhet në një depo ‒ projektin nën çadër. Duke e bërë këtë, mikroshërbimet janë maksimalisht të izoluara dhe kryejnë operacione të thjeshta, përgjegjëse për një entitet të veçantë. Me këtë qasje, është e lehtë të mbash API-në erë, të bësh ndryshime dhe të shkruash teste unike dhe integruese.
Komponentët e sistemit ndërlidhen drejtpërdrejt ose përmes një brokeri. Nga pozita e mesazheve, çdo shërbim ka disa faza të jetës:
- Inicimi i shërbimit.
Në këtë fazë ndodhin konfigurimi dhe nisja e procesit ekzekutues të shërbimit dhe varësive të tij. - Krijimi i pikës së shkëmbimit.
Shërbimi mund të përdorë një pikë shkëmbimi statike, të caktuar në konfigurimin e nodit, ose të krijojë pikë shkëmbimi dinamikisht. - Regjistrimi i shërbimit.
Që shërbimi të mund të shërbejë kërkesat, ai duhet të regjistrohet në pikën e shkëmbimit. - Funksionimi normal.
Shërbimi kryen punën e nevojshme. - Përfundimi i punës.
Ka dy lloje përfundimesh: standarde dhe emergjente. Në rastin standard, shërbimi çkombinohet nga pika e shkëmbimit dhe ndalet. Në raste emergjente, mesazhet realizojnë një nga skenarët e përpunimit të dështimeve.
Duket mjaft e komplikuar, por në kod nuk është aq e frikshme. Një shembuj kodi me komente do të jepen në analizën e modeleve më vonë.
Këmbime
Pika e shkëmbimit ‒ një proces mesazhi, që realizon logjikën e ndërveprimit me komponentët brenda strukturës së shkëmbimit të mesazheve. Në të gjitha shembujt e paraqitur më poshtë, komponentët ndërveprojnë përmes pikave të shkëmbimit, kombinimi i të cilave krijon mesazhin.
Modelet e shkëmbimit të mesazheve (MEPs)
Globalisht, modelet e shkëmbimit mund të ndahen në dyanshëm dhe njëanshëm. Të parat përfshijnë një përgjigje ndaj mesazhit të pranuar, ndërsa të dyjat jo. Një shembull klasik i një modeli dyanshëm në arkitekturën klient-server është modeli i Kërkesë-përgjigje. Le të shqyrtojmë modelin dhe modifikimet e tij.
Kërkesë-përgjigje ose RPC
RPC përdoret kur na nevojitet të marrim një përgjigje nga një proces tjetër. Ky proces mund të jetë i nisur në të njëjtin nod apo të jetë në një kontinent tjetër. Më poshtë paraqitet skema e ndërveprimit të klientit dhe serverë përmes mesazheve.

Duke qenë se mesazhi është plotësisht asinkron, për klientin shkëmbimi ndahet në 2 faza:
Dërgimi i kërkesës
mesazhimi: kërkesë(Exchange, ResponseMatchingTag, RequestDefinition, HandlerProcess).Shkëmbim ‒ emri unik i pikës së shkëmbimit
ResponseMatchingTag ‒ etiketë lokale për përpunimin e përgjigjes. Për shembull, në rast se dërgojmë disa kërkesa identike, të cilat iu takojnë përdoruesve të ndryshëm.
RequestDefinition ‒ trupi i kërkesës
HandlerProcess ‒ PID-i i përpunuesit. Ky proces do të marrë përgjigjen nga serveri.Përpunimi i përgjigjes
handle_info(#'$msg'{exchange = EXCHANGE, tag = ResponseMatchingTag,message = ResponsePayload}, State)ResponsePayload ‒ përgjigja e serverit.
Për serverin, procesi përbëhet gjithashtu nga 2 faza:
- Inizializimi i pikës së shkëmbimit
- Përpunimi i kërkesave të marra
Të ilustrojmë këtë model me kod. Le të supozojmë se duhet të realizojmë një shërbim të thjeshtë, i cili ofron vetëm një metodë për kohën e saktë.
Kodi i serverit
Do ta nxjerrim përcaktimin e API-së së shërbimit në api.hrl:
%% =====================================================
%% entitetet
%% =====================================================
-record(time, {
unixtime :: non_neg_integer(),
datetime :: binary()
}).
-record(time_error, {
code :: non_neg_integer(),
error :: term()
}).
%% =====================================================
%% metodat
%% =====================================================
-record(time_req, {
opts :: term()
}).
-record(time_resp, {
result :: #time{} | #time_error{}
}).Do ta përcaktojmë kontrolluesin e shërbimit në time_controller.erl
%% Në këtë shembull tregohet vetëm kodi i rëndësishëm. Duke e vendosur atë në modelin gen_server, mund të marrim një shërbim funksional.
%% inizializimi i gen_server
init(Args) ->
%% lidhja me pikën e shkëmbimit
messaging:monitor_exchange(req_resp, ?EXCHANGE, default, self())
{ok, #{}}.
%% përpunimi i ngjarjeve të humbjes së lidhjes me pikën e shkëmbimit. Ky është gjithashtu ngjarja që vjen nëse pika e shkëmbimit nuk ka filluar ende.
handle_info(#exchange_die{exchange = ?EXCHANGE}, State) ->
erlang:send(self(), monitor_exchange),
{noreply, State};
%% përpunimi i API-së
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};
%% përfundimi i punës së gen_server
terminate(_Reason, _State) ->
messaging:demonitor_exchange(req_resp, ?EXCHANGE, default, self()),
ok.Kodi i klientit
Për të dërguar një kërkesë në shërbim, në çdo vend të klientit mund të thërrasim API-në e kërkesës së messaging:
rast messaging:request(?EXCHANGE, tag, #time_req{opts = #{}}, self()) nga
ok -> ok;
_ -> %% logjika e përsëritjes ose dështimit
fundNë një sistem të shpërndarë, konfigurimi i komponentëve mund të jetë i ndryshëm dhe në momentin e kërkesës, messaging mund të mos ketë filluar akoma, ose kontrolluesi i shërbimit mund të mos jetë gati për të shërbyer kërkesën. Prandaj, është e nevojshme të verifikojmë përgjigjen e messaging dhe të përpunojmë rastin e dështimit.
Pas dërgimit të suksesshëm, klienti do të marrë një përgjigje ose një gabim nga shërbimi.
Do t’i përpunojmë të dy rastet 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};Request-Chunked Response
Është më mirë të shmangni dërgimin e mesazheve të mëdha. Kjo ndikon në përgjegjshmërinë dhe stabilitetin e gjithë sistemit. Nëse përgjigjja ndaj një kërkese merr shumë memorie, atëherë ndarja në pjesë është e domosdoshme.

Ja disa shembuj të tillë:
- Komponentët shkëmbejnë të dhëna binare, për shembull skedarë. Ndarja e përgjigjes në pjesë të vogla ndihmon në punën efikase me skedarë të çdo madhësie dhe në parandalimin e mbushjes së memories.
- Listimet. Për shembull, na nevojitet të zgjedhim të gjitha regjistrimet nga një tabelë të madhe në bazë dhe t'i dërgojmë një komponenti tjetër.
I quaj përgjigjet e tilla ‘lokomotiva’. Në çdo rast, 1024 mesazhe prej 1 MB janë më mirë se një mesazh i vetëm me madhësi 1 GB.
Në një klaster Erlang, ne fitojmë një avantazh shtesë - uljen e ngarkesës në pikën e shkëmbimit dhe rrjet, pasi përgjigjet dërgohen menjëherë tek marrësi, duke anashkaluar pikën e shkëmbimit.
Response with Request
Kjo është një modifikim mjaft i rrallë i modelit RPC për ndërtimin e sistemeve dialoguese.

Publiko-abonimi (pema e shpërndarjes së të dhënave)
Sistemet e orientuara nga ngjarjet dorëzojnë të dhënat te konsumatoret sipas gatishmërisë së tyre. Kështu, sistemet janë më të prirura drejt modelit push sesa pull ose poll. Kjo veçori parandalon shpenzimin e burimeve duke kërkuar dhe pritur vazhdimisht të dhëna.
Në ilustrim janë paraqitur proceset e shpërndarjes së mesazheve te konsumatoret që janë abonuar në një temë të caktuar.

Shembujt klasikë të përdorimit të këtij modeli janë shpërndarja e gjendjes: botës së lojrave në lojrat kompjuterike, të dhënave të tregut në bursë, informacionit të dobishëm në datafeed.
Le të shqyrtojmë kodin e abonentit:
init(_Args) ->
%% abonohuni në shkëmbimin, çelësi = key
messaging:subscribe(?SUBSCRIPTION, key, tag, self()),
{ok, #{}}.
handle_info(#exchange_die{exchange = ?SUBSCRIPTION}, State) ->
%% nëse pika e shkëmbimit nuk është e available, përpiqemi të ripërqafojmë
messaging:subscribe(?SUBSCRIPTION, key, tag, self()),
{noreply, State};
%% procesoni mesazhet që vijnë
handle_info(#'$msg'{exchange = ?SUBSCRIPTION, message = Msg}, State) ->
?debugVal(Msg),
{noreply, State};
%% kur ndalej konsumatori - shkëputemi nga pika e shkëmbimit
terminate(_Reason, _State) ->
messaging:unsubscribe(?SUBSCRIPTION, key, tag, self()),
ok.Burimi mund të thërrasë funksionin e publikimit të mesazheve në çdo vend të përshtatshëm:
messaging:publish_message(Exchange, Key, Message).Shkëmbim ‒ emri i pikës së shkëmbimit,
Çelësi ‒ çelësi i маршрутizimit
Mesazhi ‒ ngarkesa
Inverted Publish-subscribe

Duke zhvilluar pub-sub, mund të arrijmë një model të përshtatshëm për regjistrim. Grupi i burimeve dhe konsumatorëve mund të jetë krejtësisht i ndryshëm. Në figurë tregohet rasti me një konsumator dhe shumë burime.
Modeli i shpërndarjes së detyrave
Në almost çdo projekt ndodhin detyra të përpunimit të shtyer, siç janë krijimi i raporteve, dërgimi i njoftimeve, marrja e të dhënave nga sisteme të jashtme. Kapaciteti i sistemit që kryen këto detyra lehtë mund të shkallëzohet duke shtuar përpunues. E vetmja gjë që na mbetet është të formojmë një grup të përpunuesve dhe të shpërndajmë detyrat njëlloj mes tyre.
Le të shqyrtojmë situatat që lindin me shembuj 3 përpunuesish. Edhe në fazën e shpërndarjes së detyrave lind pyetja e drejtësisë së shpërndarjes dhe mbingarkesës së përpunuesve. Drejtësia do të garantohet nga shpërndarja round-robin, dhe për të parandaluar mbingarkesën e përpunuesve, do të vendosim një kufizim prefetch_limit. Në mënyrat e kalimit prefetch_limit nuk do të lejojë një përpunues që të marrë të gjitha detyrat.
Messaging menaxhon radhët dhe prioritetin e përpunimit. Përpunuesit marrin detyrat me kalimin e tyre. Ekzekutimi i një detyre mund të përfundojë me sukses ose me dështim:
messaging:ack(Tack)‒ thirret në rastin e përpunimit të suksesshëm të mesazhitmessaging:nack(Tack)‒ thirret në të gjitha situatat e jashtzakonshme. Pas kthimit të detyrës, messaging do ta kalojë atë tek një përpunues tjetër.

Le të supozojmë se gjatë përpunimit të tre detyrave ndodhi një dështim i ndërlikuar: përpunuesi 1 pasi mori detyrën ra, pa arritur të njoftojë asgjë në pikën e shkëmbimit. Në këtë rast, pika e shkëmbimit pas skadimit të kohës për njohjen e miratimit do ta kalojë detyrën tek një përpunues tjetër. Përpunuesi 3 nga ndonjë arsye e refuzoi detyrën dhe dërgoi nack, si rezultat detyra gjithashtu kaloi tek një përpunues tjetër që e kreu me sukses.
Përfundim paraprak
Kemi shqyrtuar bloket thelbësore të sistemeve të shpërndara dhe kemi marrë një kuptim themelor të aplikimit të tyre në Erlang/Elixir.
Duke kombinuar modelet themelore, është e mundur të krijojmë paradigma komplekse për të zgjidhur detyrat që lindin.
Në pjesën përfundimtare të ciklit, ne do të shqyrtojmë çështje të përgjithshme rreth organizimit të shërbimeve, ruterimit dhe balancimit, si dhe do të flasim për anën praktike të shkallëzueshmërisë dhe qëndrueshmërisë së sistemeve.
Fundi i pjesës së dytë.
Foto
Ilustrimet janë përgatitur me ndihmën e websequencediagrams.com
Burimi: habr.com
