Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit

Para se të fillojë një rrjedhë të re në kurs «Inxhiner i Të Dhënave» përgatitëm përkthimin e një materiali interesant.

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit

Përmbledhje

Do të flasim për një model mjaft të njohur, me të cilin aplikacionet përdorin disa depo të dhënash, ku secila depo përdoret për qëllime të saj, për shembull, për të ruajtur formën kanonike të të dhënave (MySQL etj.), për të ofruar mundësi të zgjeruara kërkimi (ElasticSearch etj.), duke e ruajtur në cache (Memcached etj.) dhe të tjera. Në përgjithësi, kur përdoren disa depo të dhënash, njëra prej tyre punon si depo kryesore, ndërsa të tjerat si depo të nxjerra. Problemi i vetëm është se si të sinkronizohen këto depo të dhënash.

Ne shqyrtuan disa modele të ndryshme që ishin përpjekur të zgjidhnin problemin e sinkronizimit të disa depoëve, si regjistrimi i dyfishtë, transaksionet e shpërndara etj. Sidoqoftë, këto qasje kanë kufizime të rëndësishme në përdorimin e tyre në jetë reale, besueshmëri dhe mirëmbajtje teknike. Përveç sinkronizimit të të dhënave, disa aplikacione gjithashtu kërkojnë të pasurojnë të dhënat duke thirrur shërbime të jashtme.

Për të zgjidhur këto probleme, u zhvillua Delta. Në fund të fundit, Delta përfaqëson një platformë të integruar, të menaxhuar nga ngjarjet për sinkronizimin dhe pasurimin e të dhënave.

Zgjidhjet ekzistuese

Shënim i dyfishtë

Për të sinkronizuar dy depo të dhënash, mund të përdoret shënimi i dyfishtë, i cili regjistron në një depo dhe menjëherë pas kësaj në depo tjetër. Regjistrimi i parë mund të përsëritet, ndërsa i dyti mund të ndërpritet nëse i pari dështoi pas përfundimit të numrit të përpjekjeve. Megjithatë, dy depo mund të ndalojnë së sinkronizuari nëse regjistrimi në depo të dytë dështon. Kjo problem zakonisht zgjidhet me krijimin e një procedure rikuperimi, e cila mund të transferojë periodikisht të dhënat nga depoja e parë në të dytën ose ta bëjë këtë vetëm në rast se në të dhëna zbulohet ndonjë diferencë.

Problemet:

Ekzekutimi i procedurës së rikuperimit është një punë specifike që nuk mund të ri përdoret. Për më tepër, të dhënat midis depozitave mbeten jashtë sinkronizimit deri në përfundimin e procedurës së rikuperimit. Zgjidhja bëhet më e komplikuar nëse përdoren më shumë se dy depozita të dhënash. Dhe, përfundimisht, procedura e rikuperimit mund të shtojë ngarkesë në burimin origjinar të të dhënave.

Tabela e regjistrimeve të ndryshimeve

Kur ndodhin ndryshime në grupin e tabelave (p.sh., inserimi, përditësimi dhe fshirja e një regjistrimi), regjistrimet e ndryshimeve shtohen në tabelën e regjistrimeve, si një pjesë e të njëjtës transaksion. Një rrjedhë tjetër ose proces vazhdimisht e kërkon ngjarjet nga tabela e regjistrimeve dhe i shkruan ato në një ose më shumë depozita të dhënash, duke fshirë ngjarjet nga tabela e regjistrimeve pas konfirmimit të shkrimit nga të gjitha depozitat.

Problemet:

Ky modeli duhet të realizohet si një bibliotekë dhe në mënyrë ideale pa ndryshuar kodin e aplikacionit që e përdor atë. Në një mjedis poliglot, realizimi i një biblioteke të tillë duhet të ekzistojë në çdo gjuhë të nevojshme, por sigurimi i qëndrueshmërisë së funksioneve dhe sjelljes mes gjuhëve është shumë i vështirë.

Problemi tjetër është se për të marrë ndryshimet e skemës, në sistemet që nuk mbështesin ndryshimet transaksionale të skemës [1][2], si për shembull MySQL. Prandaj, modeli për ekzekutimin e ndryshimit (për shembull, ndryshimi i skemës) dhe regjistrimi transaksional i tij në tabelën e logjeve të ndryshimeve nuk do të funksionojë gjithmonë.

Transaksionet e shpërndara

Transaksionet e shpërndara mund të përdoren për të ndarë një transaksion midis disa depozitave të dhënash të ndryshme, në mënyrë që operacioni të jetë ose i konfirmuar në të gjitha depozitë të përdorura, ose të mos konfirmohet në asnjërën prej tyre.

Problemet:

Transaksionet e shpërndara janë një problem shumë i madh për depozitat e ndryshme të të dhënave. Nga natyra, ato mund të mbështeten vetëm në minimalen e zakonshme të sistemeve që marrin pjesë. Për shembull, transaksionet XA bllokojnë ekzekutimin nëse ndodh një dështim gjatë fazës së përgatitjes në procesin e aplikacionit. Për më tepër, XA nuk siguron zbulimin e bllokimeve dhe nuk mbështet skemat optimiste të menaxhimit të paralelizmit. Përveç kësaj, disa sisteme si ElasticSearch nuk mbështesin XA ose ndonjë model tjetër heterogjen të transaksioneve. Si rezultat, sigurimi i atomizmit të shkruar në teknologjitë e ndryshme të ruajtjes së të dhënave mbetet një detyrë mjaft e vështirë për aplikacionet [3].

Delta

Delta u zhvillua për të adresuar kufizimet e zgjidhjeve ekzistuese të sinkronizimit të të dhënave, gjithashtu ajo lejon pasurimin e të dhënave në kohë reale. Qëllimi ynë ishte të abstragjonim të gjitha këto aspekte të komplikuara nga zhvilluesit e aplikacioneve, në mënyrë që ata të mund të përqendrohen plotësisht në realizimin e funksionalitetit të biznesit. Më pas, do të përshkruajmë "Movie Search", rastin e vërtetë të përdorimit të Delta nga Netflix.

Në Netflix, përdoret gjerësisht arkitektura mikroshërbimeve dhe çdo mikroshërbim zakonisht shërben për një lloj të dhënash. Informacioni themelor mbi filmat është i organizuar në një mikroshërbim të quajtur Shërbimi i Filmit, si dhe të dhënat e lidhura, si informacioni mbi producentët, aktorët, furnizuesit dhe të tjera, menaxhohen nga disa mikroshërbime të tjera (në veçanti Shërbimi i Marrëveshjeve, Shërbimi i Talenteve dhe Shërbimi i Furnizuesve).
Përdoruesit biznesi në Netflix Studios shpesh kanë nevojë të kërkojnë filma sipas kritereve të ndryshme, prandaj është shumë e rëndësishme për ta të kenë mundësinë të kryejnë kërkime për të gjithë të dhënat e lidhura me filmat.

Para se të shfaqej Delta, ekipi i kërkimit të filmave duhet të merrte të dhëna nga disa mikroshërbime, përpara se të indeksonte të dhënat e filmave. Përveç kësaj, ekipi duhej të zhvillonte një sistem që përditësonte rregullisht indeksin e kërkimit duke kërkuar ndryshime nga mikroshërbime të tjera, edhe nëse nuk kishte asnjë ndryshim. Ky sistem u bë shumë shpejt komplikuar dhe ishte e vështirë për t'u mbajtur.

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit
Figurë 1. Sistemi i polling para Delta
Pas fillimit të përdorimit të Delta, sistemi u thjeshtua në një sistem të menaxhuar nga ngjarje, siç tregohet në figurën e mëposhtme. Ngjarjet CDC (Change-Data-Capture) dërgohen në temat Keystone Kafka përmes Delta-Connector. Aplikacioni Delta, i ndërtuar me përdorimin e Delta Stream Processing Framework (i bazuar në Flink), merr ngjarjet CDC nga tema, i pasuron ato duke thirrur mikroshërbime të tjera, dhe, në fund, i dërgon të dhënat e pasuruara në indeksin e kërkimit në Elasticsearch. I gjithë procesi ndodh thuajse në kohë reale, domethënë, sa herë që ndryshimet regjistrohen në depo, indekset e kërkimit përditësohen.

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit
Figura 2. Pipeline i të dhënave kur përdorim Delta
Në seksionet e mëposhtme ne do të përshkruajmë funksionimin e Delta-Connector, i cili lidhet me depo dhe publikon ngjarjet CDC në nivelin e transportit, i cili përbën infrastrukturën e transmetimit të të dhënave në kohë reale, që drejton ngjarjet CDC në temat Kafka. Dhe në fund, ne do të flasim për strukturën e përpunimit të rrjedhave Delta, të cilën zhvilluesit e aplikacioneve mund ta përdorin për logjikën e përpunimit dhe pasurimit të të dhënave.

CDC (Change-Data-Capture)

Ne kemi zhvilluar shërbimin CDC me emrin Delta-Connector, i cili mund të regjistrojë ndryshimet e kaqzuara nga depoja e të dhënave në kohë reale dhe t'i shkruajë ato në një rrjedhë. Ndryshimet në kohë reale merren nga regjistri i transaksioneve dhe grumbujt e depozitave. Grumbujt përdoren pasi regjistrat e transaksioneve zakonisht nuk ruajnë gjithë historinë e ndryshimeve. Ndryshimet zakonisht serializohen si ngjarje Delta, në mënyrë që marrësi të mos shqetësohet për origjinën e ndryshimit.

Delta-Connector mbështet disa funksione shtesë, të tilla si:

  • MundĂ«sia pĂ«r tĂ« shkruar nĂ« dalje tĂ« personalizuara pĂ«rtej Kafka.
  • MundĂ«sia pĂ«r tĂ« aktivizuar grumbuj manualĂ« nĂ« çdo kohĂ« pĂ«r tĂ« gjitha tabelat, pĂ«r njĂ« tabelĂ« tĂ« caktuar ose pĂ«r çelĂ«sa tĂ« caktuar tĂ« parĂ«.
  • Grumbujt mund tĂ« merren nĂ« pjesĂ«, kĂ«shtu qĂ« nuk ka nevojĂ« tĂ« filloni gjithçka nga fillimi nĂ« rast dĂ«shtimi.
  • Nuk ka nevojĂ« tĂ« vendosni bllokime nĂ« tabela, e cila Ă«shtĂ« shumĂ« e rĂ«ndĂ«sishme pĂ«r tĂ« siguruar qĂ« trafiku i shkruar nĂ« bazĂ«n e tĂ« dhĂ«nave kurrĂ« tĂ« mos bllokohet nga shĂ«rbimi ynĂ«.
  • DisponueshmĂ«ri e lartĂ« pĂ«r shkak tĂ« instanceve rezervĂ« nĂ« AWS Availability Zones.

Tani ne mbështesim MySQL dhe Postgres, përfshirë edhe kur bëhet fjalë për implementimin në AWS RDS dhe Aurora. Gjithashtu mbështesim Cassandra (multi-master). Më shumë detaje rreth Delta-Connector mund të merrni në këtë blogun e tij.

Kafka dhe niveli i transportit

Niveli i transportit të ngjarjeve Delta është ndërtuar mbi shërbimin e shkëmbimit të mesazheve të platformës Keystone.

Në mënyrë historike, publikimi i mesazheve në Netflix është optimizuar për të rritur disponueshmërinë, jo qëndrueshmërinë (shih artikulli i mëparshëm). Kompromisi ishte potenciali për mos përputhje të të dhënave të brokerit në skenarë të ndryshëm kufitarë. Për shembull, elektori i papastër është përgjegjës që marrësi potencialisht të dyfishojë ose të humbasë ngjarje.

Me Delta ne donim të kishim më shumë garanci të rëndësishme për qëndrueshmërinë, për të siguruar dorëzimin e ngjarjeve CDC në magazinat përkatëse. Për këtë, ne propozuam një klaster Kafka të dizajnuar posaçërisht si një objekt i klasës së parë. Mund të shihni disa konfigurime të brokerit në tabelën më poshtë:

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit

Në klasterët Keystone Kafka, elektori i papastër zakonisht janë të aktivizuara për të siguruar disponueshmërinë e botuesit. Kjo mund të çojë në humbjen e mesazheve në rast se një replika e nesinhronizuar zgjidhet si lider. Për një grup të ri me besueshmëri të lartë Kafka, parametri elektori i papastër është i çaktivizuar për të parandaluar humbjen e mesazheve.

Po ashtu, ne e kemi rritur faktorin e replikimit nga 2 në 3 dhe replicas minimale të sinkronizuara nga 1 në 2. Botuesit që shkruajnë në këtë grup kërkojnë acks nga të gjithë të tjerët, duke garantuar që 2 nga 3 replikat do të kenë mesazhet më të fundit të dërguara nga botuesi.

Kur një instancë e brokerit përfundohet, një instancë e re zëvendëson të vjetrën. Megjithatë, brokeri i ri do të duhet të arrijë replikat që nuk janë sinkronizuar, gjë që mund të zgjasë disa orë. Për të shkurtuar kohën e rikuperimit të këtij skenari, filluam të përdorim ruajtjen e blloqeve të të dhënave (Amazon Elastic Block Store) në vend të disqeve lokale të brokerëve. Kur një instancë e re zëvendëson një instancë brokeri të përfunduar, ajo bashkangjitet me volumet EBS që ishin te instanca e përfunduar dhe fillon të arrijë mesazhet e reja. Ky proces e shkurton kohën e shlyerjes së vonesave nga disa orë në disa minuta, pasi instanca e re nuk ka më nevojë të replikojë nga një gjendje e zbrazët. Në përgjithësi, ciklet e jetës së ruajtjes dhe të brokerit ndihmojnë në uljen e ndikimit të efektit të ndërrimit të brokerit.

Për të rritur edhe më tej garantimin e dorëzimit të të dhënave, përdorim sistemin e gjurmimit të mesazheve për të zbuluar çdo humbje mesazhesh në kushte ekstreme (për shembull, përshkallëzimin e orëve në liderin e seksionit).

Paraqitja e Procesimit të Rrymës

Niveli i përpunimit në Delta është ndërtuar mbi platformën Netflix SPaaS, e cila siguron integrimin e Apache Flink me ekosistemin Netflix. Platforma ofron një ndërfaqe përdoruese që menaxhon implementimin e detyrave Flink dhe orkestrimin e klubeve Flink mbi platformën tonë të menaxhimit të kontejnerëve Titus. Ndërfaqja gjithashtu menaxhon konfigurimet e detyrave dhe lejon përdoruesit të bëjnë ndryshime në konfigurim dinamikisht pa pasur nevojë të rikompilojnë detyrat Flink.

Delta ofron një kornizë përpunimi rrjedhës (stream processing framework) të të dhënave në bazë të Flink dhe SPaaS, e cila përdor të bazuar në annotime DSL (Gjuha e Veçantë e Domeneve), për të abstaraktuar detajet teknike. Për shembull, për të përcaktuar hapin me të cilin do të pasurohen ngjarjet, duke bërë thirrje në shërbime të jashtme, përdoruesit duhet të shkruajnë DSL-në e mëposhtme, dhe korniza do të krijojë një model mbi të cilin do të ekzekutohet Flink.

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit
Figurë 3. Shembulli i pasurimit në DSL në Delta

Korniza për procesimin e të dhënave jo vetëm që redukton njohuritë e nevojshme, por gjithashtu ofron funksione të zakonshme për përpunimin e rrjedhave, siç janë deduplikimi, skematizimi dhe fleksibiliteti dhe qëndrueshmëria për të zgjidhur probleme të zakonshme në operim.

Delta Stream Processing Framework përbëhet nga dy module kyçe, moduli DSL & API dhe moduli Runtime. Moduli DSL & API ofron DSL dhe API UDF (Funksion i Përcaktuar nga Përdoruesi) për të lejuar përdoruesit të shkruajnë logjikën e tyre për procesimin (p.sh., filtrimin ose transformimet). Moduli Runtime ofron implementimin e parserit DSL, i cili ndalon përfaqësimin e brendshëm të hapave të përpunimit në modelet DAG. Komponenti i Ekzekutimit interpreton modelet DAG për të inicializuar operatorët realë Flink dhe për të filluar aplikacionin Flink. Arkitektura e kornizës ilustrohet në figurën e mëposhtme.

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit
Figura 4. Arkitektura e Delta Stream Processing Framework

Ky qasje ka disa përparësi:

  • PĂ«rdoruesit mund tĂ« pĂ«rqendrohen nĂ« logjikĂ«n e tyre tĂ« biznesit pa pasur nevojĂ« tĂ« thellohen nĂ« specifikat e Flink ose strukturĂ«n SPaaS.
  • Optimizimi mund tĂ« kryhet nĂ« njĂ« mĂ«nyrĂ« qĂ« Ă«shtĂ« transparente pĂ«r pĂ«rdoruesit, dhe gabimet mund tĂ« rregullohen pa pasur nevojĂ« tĂ« bĂ«hen ndonjĂ« ndryshim nĂ« kodin e pĂ«rdoruesit (UDF).
  • Funksionimi i aplikacioneve Delta Ă«shtĂ« lehtĂ«suar pĂ«r pĂ«rdoruesit, pasi platforma ofron fleksibilitet dhe qĂ«ndrueshmĂ«ri qĂ« vijnĂ« me paketĂ« dhe grumbullon njĂ« numĂ«r tĂ« madh statistikash tĂ« detajuara qĂ« mund tĂ« pĂ«rdoren pĂ«r njoftime.

Përdorimi në prodhim

Delta ka qenë në prodhim për më shumë se një vit dhe luan një rol kyç në shumë aplikacione të Netflix Studio. Ajo ka ndihmuar ekipet të realizojnë përdorime të tilla si indeksimi i kërkimeve, ruajtja e të dhënave dhe proceset e punës të menaxhuara nga ngjarjet. Më poshtë është një pasqyrë e arkitekturës së lartë të platformës Delta.

Delta: Platforma e sinkronizimit të të dhënave dhe pasurimit
Figura 5. Arkitektura e lartë e Delta.

Faleminderit

Dëshirojmë të falënderojmë individët e mëposhtëm, të cilët morën pjesë në krijimin dhe zhvillimin e Delta në Netflix: Allen Wang, Charles Zhao, Jaebin Yoon, Josh Snyder, Kasturi Chatterjee, Mark Cho, Olof Johansson, Piyush Goyal, Prashanth Ramdas, Raghuram Onti Srinivasan, Sandeep Gupta, Steven Wu, Tharanga Gamaethige, Yun Wang dhe Zhenzhong Xu.

Burimet

  1. dev.mysql.com/doc/refman/5.7/en/implicit-commit.html
  2. dev.mysql.com/doc/refman/5.7/en/cannot-roll-back.html
  3. Martin Kleppmann, Alastair R. Beresford, Boerge Svingen: Procesimi i ngjarjeve online. Commun. ACM 62(5): 43–49 (2019). DOI: doi.org/10.1145/3312527

Regjistrohu për webinarin falas: «Data Build Tool për magazinën Amazon Redshift».

Burimi: habr.com

Bli njĂ« hosting tĂ« besueshĂ«m pĂ«r faqet me mbrojtje DDoS, VPS VDS serverĂ« đŸ”„ Bli njĂ« hosting tĂ« besueshĂ«m pĂ«r faqet me mbrojtje DDoS, VPS VDS serverĂ« | ProHoster