
Tere, Habr.
Hiljuti jagasin ma selle kohta, milliseid parameetreid me tihti kasutame Kafka producentide ja tarbijate jaoks, et tagada garanteeritud kohaletoimetamine. Selles artiklis rÀÀgin, kuidas korraldasime sĂŒndmuse taastöötluse, mis saadi Kafka'lt vĂ€lise sĂŒsteemi ajutise kĂ€ttesaamatuse tĂ”ttu.
Kaasaegsed rakendused toimivad vĂ€ga keerulises keskkonnas. ĂriĂŒksus, mis on pakitud kaasaegsesse tehnoloogilisse virna, töötades Docker'i pildis, mida haldab orkestraator nagu Kubernetes vĂ”i OpenShift, ja kommunikeerides teiste rakendustega vĂ”i ettevĂ”tte lahendustega fĂŒĂŒsiliste ja virtuaalsete ruuterite ahela kaudu. Sellises keskkonnas vĂ”ib alati midagi katki minna, seega on sĂŒndmuste taastöötlus, kui ĂŒks vĂ€lisest sĂŒsteemist ei ole saadaval â oluline osa meie Ă€riprotsessidest.
Kuidas oli enne Kafka
Varasemalt kasutasime projektis IBM MQ asĂŒnkroonsete sĂ”numite saatmiseks. Kui teenuse töötamise kĂ€igus tekkis mĂ”ni viga, vĂ”is saadud sĂ”num paigutada dead-letter-queue (DLQ), et seda hiljem kĂ€sitsi analĂŒĂŒsida. DLQ loodi koos sissetuleva jĂ€rjekorraga, sĂ”numi ĂŒmberpaigutamine toimus IBM MQ sees.
Kui viga oli ajutine ja me suutsime selle mÀÀrata (nĂ€iteks ResourceAccessException HTTP-kĂ”ne ajal vĂ”i MongoTimeoutException MongoDb pĂ€ringu puhul), siis rakendati kordusvĂ€ljakutsete strateegiat. ĂkskĂ”ik millest rakenduse loogika branch'ist, algne sĂ”num omistati kas sĂŒsteemijĂ€rjekorda edasilĂŒkkamiseks vĂ”i eraldi rakendusse, mis kunagi ammu loodi sĂ”numite kordus saatmiseks. SĂ”numi pĂ€isesse kirjutatakse kordusvĂ€ljakutsete number, mis on seotud viivituse intervalliga vĂ”i rakenduse tasemel strateegia lĂ”puga. Kui me jĂ”uame strateegia lĂ”puni, kuid vĂ€line sĂŒsteem on endiselt kĂ€ttesaamatu, paigutatakse sĂ”num DLQ-sse kĂ€sitsi analĂŒĂŒsimiseks.
Lahenduse otsimine
, vĂ”ib leida jĂ€rgmist . LĂŒhidalt öeldes pakutakse iga viivituse intervalli jaoks luua teema ja rakenduse kĂŒljes tarbijaid, mis loevad sĂ”numid vajaliku viivitusega.

MalgrĂ© le grand nombre d'avis positifs, cela ne me semble pas tout Ă fait rĂ©ussi. Tout d'abord, parce que le dĂ©veloppeur, en plus de mettre en Ćuvre les exigences commerciales, devra passer beaucoup de temps Ă mettre en Ćuvre le mĂ©canisme dĂ©crit.
Lisaks, kui Kafka klastri juures on lubatud juurdepÀÀsu haldamine, tuleb kulutada aega teemade loomisele ja nendele vajalike juurdepÀÀsude tagamisele. Peale selle on vajalik valida Ôige retention.ms parameeter iga retry-teema jaoks, et sÔnumid jÔuaksid uuesti edastada ja ei kaoks. Rakenduse loomine ja juurdepÀÀsude taotlemine tuleb korrata iga olemasoleva vÔi uue teenuse jaoks.
NĂŒĂŒd vaatame, milliseid sĂ”numi uuesti töötlemise mehhanisme pakub meile spring ĂŒldiselt ja spring-kafka konkreetsemalt. Spring-kafka-l on ĂŒlekantav sĂ”ltuvus spring-retry-st, mis pakub abstraktsioone erinevate BackOffPolicy-de haldamiseks. See on ĂŒsna paindlik tööriist, kuid selle oluline puudus on sĂ”numite hoidmine uuesti saatmiseks rakenduse mĂ€lus. See tĂ€hendab, et rakenduse taaskĂ€ivitamine uuenduse vĂ”i tööea vigade tĂ”ttu toob kaasa kĂ”ikide uuesti töötlemist ootavate sĂ”numite kaotuse. Kuna see on meie sĂŒsteemi jaoks kriitiline punkt, ei uurinud me seda edaspidi.
Ise spring-kafka pakub mitmeid rakendusi ContainerAwareErrorHandler-le, nÀiteks , mille abil saab sÔnumit hiljem töödelda, mitte nihutades offset'i vea tekkimise korral. Alates versioonist spring-kafka 2.3 on vÔimalus mÀÀrata BackOffPolicy.
See lĂ€henemine vĂ”imaldab korduvkasutatavatel sĂ”numitel rakenduse taaskĂ€ivitamisest ĂŒle elada, kuid DLQ mehhanism on endiselt puudu. Just selle valiku tegime 2019. aasta alguses, optimistlikult arvates, et DLQ-d ei ole vaja (meil vedas ja see tĂ”epoolest ei olnud vajalik mitme kuu jooksul rakenduse sellise kordusprotsessimise sĂŒsteemiga töötamise ajal). Ajutised vead viisid SeekToCurrentErrorHandleri aktiveerimiseni. ĂlejÀÀnud vead kanti logisse, mis viis offset'i nihkeni ning töötlemine jĂ€tkus jĂ€rgmise sĂ”numiga.
LÔplik otsus
SeekToCurrentErrorHandleril pÔhinev rakendus viis meid oma sÔnumite kordussaata mehhanismi vÀljatöötamiseni.
Esiteks soovisime kasutada juba olemasolevat kogemust ja laiendada seda rakenduse loogika pĂ”hjal. Rakenduse puhul, millel on lineaarne loogika, oleks optimaalne lĂ”petada uute sĂ”numite lugemine lĂŒhikese aja jooksul, mis on mÀÀratud uuestis kutsumise strateegia raames. Muude rakenduste puhul soovime, et oleks olemas ĂŒhtne punkt, mis tagab uuestis kutsumise strateegia rakendamise. Lisaks peaks see ĂŒhtne punkt omama DLQ funktsionaalsust mĂ”lema lĂ€henemise puhul.
Uuestis kutsumise strateegia peaks olema salvestatud rakenduses, mis vastutab jÀrgmise intervalli saamise eest ajutise vea ilmnemisel.
Consumer'i peatamine lineaarse loogikaga rakendustes
Spring-kafka kasutamisel vÔib Consumer'i peatamise kood vÀlja nÀha umbes nii:
public void pauseListenerContainer(MessageListenerContainer listenerContainer,
Instant retryAt) {
if (nonNull(retryAt) && listenerContainer.isRunning()) {
listenerContainer.stop();
taskScheduler.schedule(() -> listenerContainer.start(), retryAt);
return;
}
// DLQ jaoks
}NÀites retryAt on aeg, millal tuleb uuesti kÀivitada MessageListenerContainer, kui see veel töötab. Uuesti kÀivitamine toimub eraldi teemas, mis on kÀivitatud TaskScheduler'is, mille rakendust pakub ka spring.
MÀÀrame retryAt vÀÀrtuse jÀrgmisel viisil:
- Otsitakse korduvate kutsete loendi vÀÀrtust.
- Loendi vÀÀrtuse pÔhjal otsitakse praegune viivituse intervall korduvate kutsete strateegias. Strateegia kuulutatakse vÀlja rakenduses endas, selle salvestamiseks valisime JSON formaadi.
- JSON-massis leitud intervall sisaldab sekundeid, mille möödumisel tuleb töötlemine uuesti kÀivitada. See sekundite arv liidetakse praegusele ajale, luues vÀÀrtuse retryAt jaoks.
- Kui intervalli ei leita, siis on retryAt vÀÀrtus null ja sÔnum saadetakse DLQ-sse kÀsitsi lahendamiseks.
Sellise lÀhenemise puhul jÀÀb alles vaid sÀilitada iga töötlemisel oleva teadete korduvate kÔnede arv, nÀiteks rakenduse mÀlus. Katsete arvu salvestamine mÀllu ei ole selle lÀhenemise jaoks kriitiline, kuna lineaarse loogikaga rakendus ei saa tervikuna töötlemist teostada. Erinevalt spring-retry'ist, rakenduse taaskÀivitamine ei too kaasa kÔigi teadete kadumist korduvaks töötlemiseks, vaid lihtsalt strateegia taaskÀivitamist.
See lĂ€henemine aitab vĂ€hendada koormust vĂ€lisele sĂŒsteemile, mis vĂ”ib olla kergesti kergesti ĂŒle koormatud ja seetĂ”ttu mittesaadav. TeisisĂ”nu, lisaks korduvatele töötlemistele oleme saavutanud mustri rakendamise. .
Meie juhul on vea kĂŒnnis vaid 1, ja et minimeerida sĂŒsteemi seiskamist ajutise nĂ€tivĂ”rgu katkestuse tĂ”ttu, kasutame vĂ€ga granulaarset korduvate kĂ”nede strateegiat vĂ€ikeste viibimisega intervallidega. See ei pruugi sobida kĂ”ikidele kontserni rakendustele, seega tuleb vea kĂŒnnise ja intervalli suuruse vaheline suhe valida, tuginedes sĂŒsteemi eripĂ€radele.
Eraldi rakendus, mis kÀsitleb rakendustelt saadud sÔnumeid, millel on mÀÀramatu loogika.
Siin on koodinÀide, mis saadab sÔnumit sellesse rakendusse (Retryer), mis saadab jÀlle teemale DESTINATION, kui jÔuab aega RETRY_AT:
public void retry(ConsumerRecord record, String retryToTopic,
Instant retryAt, String counter, String groupId, Exception e) {
Headers headers = ofNullable(record.headers()).orElse(new RecordHeaders());
List arrayOfHeaders =
new ArrayList(Arrays.asList(headers.toArray()));
updateHeader(arrayOfHeaders, GROUP_ID, groupId::getBytes);
updateHeader(arrayOfHeaders, DESTINATION, retryToTopic::getBytes);
updateHeader(arrayOfHeaders, ORIGINAL_PARTITION,
() -> Integer.toString(record.partition()).getBytes());
if (nonNull(retryAt)) {
updateHeader(arrayOfHeaders, COUNTER, counter::getBytes);
updateHeader(arrayOfHeaders, SEND_TO, "retry"::getBytes);
updateHeader(arrayOfHeaders, RETRY_AT, retryAt.toString()::getBytes);
} else {
updateHeader(arrayOfHeaders, REASON,
ExceptionUtils.getStackTrace(e)::getBytes);
updateHeader(arrayOfHeaders, SEND_TO, "backout"::getBytes);
}
ProducerRecord messageToSend =
new ProducerRecord(retryTopic, null, null, record.key(), record.value(), arrayOfHeaders);
kafkaTemplate.send(messageToSend);
} NÀiteks on nÀha, et palju teavet edastatakse pÀistes. RETRY_AT vÀÀrtus asub samamoodi nagu tarbija peatamise kordamise mehhanismi puhul. Peale DESTINATION ja RETRY_AT edastame:
- GROUP_ID, mille alusel grupeerime sĂ”numeid kĂ€sitsi analĂŒĂŒsimiseks ja otsingu lihtsustamiseks.
- ORIGINAL_PARTITION, et pĂŒĂŒda sĂ€ilitada sama tarbija uuesti töötlemiseks. See parameeter vĂ”ib olla null, sel juhul saadakse uus partitsioon originaalsĂ”numi record.key() vĂ”tme jĂ€rgi.
- Ăksuste uuendatud COUNTER vÀÀrtus, et jĂ€rgida korduvate kutsete strateegiat.
- SEND_TO â konstant, mis nĂ€itab, kas sĂ”num tuleks saata uuesti töötlemiseks pĂ€rast RETRY_AT saavutamist vĂ”i paigutada DLQ-sse.
- REASON â pĂ”hjus, miks sĂ”numi töötlemine katkestati.
Retryer salvestab sĂ”numid edasise saatmise ja kĂ€sitsi analĂŒĂŒsimise jaoks PostgreSQL-is. Ajastatult kĂ€ivitatakse ĂŒlesanne, mis leiab sĂ”numid, mille RETRY_AT on möödunud, ja saadab need tagasi DESTINATIONi TOPICALi ORIGINAL_PARTITION partitsioonile record.key() vĂ”tmega.
PĂ€rast sĂ”numi saatmist kustutatakse need PostgreSQL-ist. SĂ”numite kĂ€sitsi töötlemine toimub lihtsas kasutajaliideses, mis suheldes Retryeriga ĂŒle REST API. Peamised funktsioonid hĂ”lmavad sĂ”numite uuesti saatmist vĂ”i DLQ-st kustutamist, veateabe vaatamist ja sĂ”numite otsimist, nĂ€iteks vea nime jĂ€rgi.
Kuna meie klasterites on lubatud juurdepÀÀsuhaldus, on vaja tÀiendavalt taotleda juurdepÀÀse teemale, mida kuulab Retryer, ja anda Retryerile vÔimalus kirjutada DESTINATION teemale. See on ebamugav, kuid vÔrreldes ajavahemikute pÔhise lÀhenemisega on meil tÀielik DLQ ja kasutajaliides selle haldamiseks.
On juhtumeid, kui sisendteemat loevad mitmed erinevad tarbijagruppide rakendused, millel on erinev loogika. SĂ”numi uuesti töötlemine Retryeri kaudu ĂŒhe sellise rakenduse jaoks toob teises esile duplikaadi. Selle vastu kaitsmiseks loome eraldi teema uuesti töötlemiseks. Sisend- ja retry-teemat vĂ”ib lugeda sama tarbija ilma piiranguteta.

Vaikimisi ei paku see lĂ€henemine circuit breakerâi vĂ”imalust, kuid selle saab rakendusele lisada vĂ”i uus , mĂ€hkides vĂ€liste teenuste kutsumise kohad vastavatesse abstraktsioonidesse. Lisaks ilmub vĂ”imalus valida strateegia mustrile, mis vĂ”ib samuti olla kasulik. NĂ€iteks spring-cloud-netflixis vĂ”ib see olla teepool vĂ”i seeder.
KokkuvÔte
Tulemuseks on meil eraldi rakendus, mis vĂ”imaldab sĂ”numi töötlemist korrata, kui mĂ”ni vĂ€line sĂŒsteem on ajutiselt kĂ€ttesaamatu.
Ere rakenduse peamine eelis on see, et seda saavad kasutada vĂ€lised sĂŒsteemid, mis töötavad samas Kafka klastris, ilma mĂ€rkimisvÀÀrsete muudatusteta oma kĂŒljel! Selline rakendus peab ainult saama juurdepÀÀsu retry-teemadele, tĂ€itma mĂ”ned Kafka pealkirjad ja saatma sĂ”numi Retryerisse. Ăhtegi tĂ€iendavat infrastruktuuri ei pea ĂŒles tĂ”stma. Ja et vĂ€hendada sĂ”numite saadetavaid koguseid rakendusest Retryerisse ja tagasi, oleme eraldanud rakendused, millel on lineaarne loogika, ning teinud neisse uuesti töötlemise kautta Consumeri peatumise.
Allikas: habr.com
