Si Kafka u bë realitet

Si Kafka u bë realitet

Përshëndetje, Habr!

Punoj në ekipin Tinkoff, i cili merret me zhvillimin e qendrës sonë të njoftimit. Kryesisht zhvilloj në Java duke përdorur Spring Boot dhe zgjidh probleme teknike të ndryshme që shfaqen në projekt.

Shumica e mikroshërbimeve tona ndërveprojnë asinkronisht përmes një brokeri mesazhi. Më parë, përdorëm IBM MQ si broker, i cili nuk arriti të përballojë ngarkesën, por kishte garanci të ulta për dërgesat.

Si zëvendësim, na u propozua Apache Kafka, e cila ka potencial të lartë për shkallëzim, por, fatkeqësisht, kërkon një qasje pothuajse të personalizuar për konfigurimin për skenarë të ndryshëm. Për më tepër, mekanizmi i dërgesës së paktën një herë, i cili funksionon në Kafka si parazgjedhje, nuk lejonte mbajtjen e nivelit të nevojshëm të konsistencës nga kutia. Më poshtë do të ndaj përvojën tonë të konfigurimit të Kafka, duke treguar veçanërisht se si të konfigurojmë dhe të jetojmë me dërgesën saktësisht një herë.

Dërgesa e garantuar dhe jo vetëm

Parametrat, për të cilët do të flasim më poshtë, do të ndihmojnë në parandalimin e një sërë problemeve me konfigurimin e parazgjedhur. Por fillimisht dëshirojmë të kushtojmë vëmendje një parametri që do ta lehtësojë potencialin e proçesit të debagut.

Këtu do të ndihmojë client.id për Producuesin dhe Konsumatorin. Në shikim të parë, si vlerë mund të përdorim emrin e aplikacionit, dhe në shumicën e rasteve do të funksionojë. Megjithatë, situata kur në aplikacion përdoren disa Konsumatorë dhe u jepni atyre të njëjtin client.id, çon në paralajmërimin e mëposhtëm:

org.apache.kafka.common.utils.AppInfoParser — Gabim te regjistrimi i AppInfo mbean javax.management.InstanceAlreadyExistsException: kafka.consumer:type=app-info,id=kafka.test-0

Nëse dëshironi të përdorni JMX në një aplikacion me Kafka, kjo mund të jetë një problem. Për këtë rast, më e mira është që si vlerë client.id të përdoret një kombinim i emrit të aplikacionit dhe, për shembull, emrit të temës. Rezultatin e konfigurimit tonë mund ta shikoni në daljen e komandës kafka-consumer-groups nga utilitarët e Confluent:

Si Kafka u bë realitet

Tani le të shqyrtojmë skenarin e dërgesës së garantuar të mesazhit. Producuesi i Kafka ka parametrin acks, i cili lejon të konfiguroni se pas sa pranimeve lideri i klasit duhet të marrë këtë mesazh si të shkruar me sukses. Ky parametr mund të marrë vlerat e mëposhtme:

  • 0 — pranimet nuk do tĂ« merren parasysh.
  • 1 — parametri i parazgjedhur, Ă«shtĂ« e nevojshme pranimi vetĂ«m nga 1 replikĂ«.
  • −1 — kĂ«rkohet njohja nga tĂ« gjitha replikat e sinkronizuara (konfigurimi i klasterit min.insync.replicas).

Nga vlerat e pĂ«rmendura, Ă«shtĂ« e qartĂ« se acks qĂ« Ă«shtĂ« −1 ofron garancitĂ« mĂ« tĂ« forta qĂ« mesazhi nuk do tĂ« humbasĂ«.

Siç e dimë të gjithë, sistemet e shpërndara janë të pasigurta. Për t'u mbrojtur nga defektet përkohësore, Kafka Producer ofron parametrin retries, i cili lejon caktimin e numrit të përpjekjeve për dërgim brenda delivery.timeout.ms. Duke qenë se parametri retries ka vlerën e paracaktuar Integer.MAX_VALUE (2147483647), numri i përsëritjeve të dërgesës së mesazhit mund të rregullohet duke ndryshuar vetëm delivery.timeout.ms.

Le të kalojmë te dërgimi me saktësi të vetme

Caktimet e pĂ«rmendura lejojnĂ« qĂ« Produsi ynĂ« tĂ« dĂ«rgojĂ« mesazhe me njĂ« garantim tĂ« lartĂ«. Tani le tĂ« flasim se si tĂ« garantojmĂ« regjistrimin e njĂ« kopjeje tĂ« vetme tĂ« mesazhit nĂ« Kafka-topic? NĂ« rastin mĂ« tĂ« thjeshtĂ«, pĂ«r kĂ«tĂ« Produsi duhet tĂ« vendosĂ« parametrin enable.idempotence nĂ« vlerĂ«n true. Idempotenca garanton regjistrimin e vetĂ«m njĂ« mesazhi nĂ« njĂ« parti tĂ« caktuar tĂ« njĂ« topiku. Parakusht pĂ«r aktivizimin e idempotencĂ«s janĂ« vlerat acks = all, retry > 0, max.in.flight.requests.per.connection ≀ 5. NĂ«se kĂ«to parametra nuk janĂ« caktuar nga zhvilluesi, do tĂ« vendosen automatikisht vlerat e pĂ«rmendura mĂ« sipĂ«r.

Kur idempotenca Ă«shtĂ« e konfiguruar, duhet tĂ« sigurohemi qĂ« mesazhet e njĂ«jta tĂ« shkojnĂ« gjithmonĂ« nĂ« tĂ« njĂ«jtat parti. Kjo mund tĂ« bĂ«het duke rregulluar çelĂ«sin dhe parametrin partitioner.class nĂ« Produs. Le tĂ« fillojmĂ« me çelĂ«sin. PĂ«r çdo dĂ«rgesĂ« ai duhet tĂ« jetĂ« i njĂ«jtĂ«. Kjo Ă«shtĂ« e lehtĂ« pĂ«r t'u arritur duke pĂ«rdorur njĂ« identifikues biznesi nga mesazhi origjinal. Parametri partitioner.class ka vlerĂ«n e paracaktuar — DefaultPartitioner. Me kĂ«tĂ« strategji parashikimi tĂ« paracaktuar veprojmĂ« kĂ«shtu:

  • NĂ«se partia Ă«shtĂ« e specifikuar qartĂ« gjatĂ« dĂ«rgimit tĂ« mesazhit, atĂ«herĂ« e pĂ«rdorim atĂ«.
  • NĂ«se partia nuk Ă«shtĂ« e specifikuar, por çelĂ«si Ă«shtĂ« i caktuar — zgjedhim partinĂ« sipas hashes nga çelĂ«si.
  • NĂ«se as partia as çelĂ«si nuk janĂ« caktuar — zgjedhim partitĂ« me radhĂ« (round-robin).

PĂ«r mĂ« tepĂ«r, pĂ«rdorimi i çelĂ«sit dhe dĂ«rgimi idempotent me parametrin max.in.flight.requests.per.connection = 1 ju jep procesim tĂ« renditur tĂ« mesazheve nĂ« Consumer. ËshtĂ« mirĂ« tĂ« mbani mend se, nĂ«se nĂ« klasterin tuaj Ă«shtĂ« e aktivizuar menaxhimi i qasjes, do t'ju nevojiten tĂ« drejtat pĂ«r regjistrimin idempotent nĂ« temĂ«.

Nëse ndonjëherë ju mungojnë mundësitë e dërgimit idempotent sipas çelësit ose logjika në anën e Producer kërkon ruajtjen e konsistencës së të dhënave midis particioneve të ndryshme, atëherë transaksionet do t'ju vijnë në ndihmë. Për më tepër, me anë të transaksionit të ndërlikuar mund të sinkronizoni kushtimisht regjistrimin në Kafka, për shembull, me regjistrimin në DB. Për të aktivizuar dërgimin transaksional në Producer, është necesare që ai të ketë idempotencë dhe të vendosni transactional.id. Nëse klasteri juaj Kafka ka menaxhim qasje, atëherë për regjistrimin transaksional, ashtu si për atë idempotent, do t'ju nevojiten të drejtat për regjistrim, të cilat mund të jepen nëpërmjet maske të përdorur me vlerën që ruhet në transactional.id.

Formalisht si identifikues transaksioni mund të përdoret çdo varg, për shembull emri i aplikacionit. Por nëse jeni duke bërë ekzekutim të disa instancave të njëjtit aplikacion me të njëjtin transactional.id, atëherë instanca e parë e nisur do të ndalet me gabim, pasi Kafka do ta konsiderojë atë si proces zombi.

org.apache.kafka.common.errors.ProducerFencedException: Prodhuesi përpiqet të kryejë një operacion me një epokë të vjeter. Ose ka një prodhues më të ri me të njëjtin transactionalId, ose transaksioni i prodhuesit ka skaduar nga brokeri.

Për të zgjidhur këtë problem, ne shtojmë në emrin e aplikacionit një sufix në formën e emrit të hostit, të cilin e marrim nga variablat e mjedisit.

Prodhuesi është i konfiguruar, por transaksionet në Kafka menaxhojnë vetëm fushën e dukshmërisë së mesazhit. Pavarësisht nga statusi i transaksionit, mesazhi menjëherë hyn në temë, por disponon disa atribute shtesë sistematike.

Që këto mesazhe të mos lexohen përpara kohe nga Consumer, ai duhet të vendosë parametrin isolation.level në vlerën read_committed. Këtë Consumer do të mund të lexojë mesazhet e pa-transaksionuara si më parë, dhe ato transaksionale vetëm pas angazhimit.
Nëse keni vendosur të gjitha konfigurimet e përmendura më parë, atëherë keni konfiguruar dorëzim exactly once. Urime!

Por ka një nuancë tjetër. Transactional.id, i cili u konfiguruar më lart, në të vërtetë është një prefiks i transaksionit. Në menaxherin e transaksioneve, atij i shtohet një numër radhor. Identifikuesi i marrë jepet në transactional.id.expiration.ms, i cili konfigurohet në klasterin Kafka dhe ka një vlerë të paracaktuar "7 ditë". Nëse gjatë kësaj kohe aplikacioni nuk ka pranuar asnjë mesazh, atëherë gjatë përpjekjes për të dërguar transaksionin e ardhshëm do të merrni InvalidPidMappingException. Pas kësaj, koordinatori i transaksioneve do të japë një numër rendor të ri për transaksionin e ardhshëm. Në këtë rast, mesazhi mund të humbasë, nëse InvalidPidMappingException nuk trajtohet siç duhet.

Në vend të rezultateve

Siç mund të vëreni, nuk mjafton thjesht të dërgoni mesazhe në Kafka. Duhet të zgjidhni kombinimin e parametrave dhe të jeni të gatshëm për të bërë ndryshime të shpejta. Në këtë artikull, kam përpjekur të tregoj në detaje konfigurimin e dërgimit 'exactly once' dhe kam përshkruar disa probleme me konfigurimet e client.id dhe transactional.id, me të cilat u përballëm. Më poshtë janë përmbledhje të konfigurimeve të Prodhueseve dhe Konsumatorëve.

Producenti:

  1. acks = all
  2. retries > 0
  3. enable.idempotence = true
  4. max.in.flight.requests.per.connection ≀ 5 (1 — pĂ«r dĂ«rgesĂ« tĂ« renditur)
  5. transactional.id = ${application-name}-${hostname}

Konsumatori:

  1. isolation.level = read_committed

Për të minimizuar gabimet në aplikacionet e ardhshme, ne krijuam një mbështjellës mbi konfigurimin e spring-ut, ku tashmë janë caktuar vlerat për disa nga parametrat e përmendur.

Dhe këtu janë disa materiale për studim të pavarur:

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster