SĂŒndmuste uuesti töötlemine, mis on saadud Kafka'st

SĂŒndmuste uuesti töötlemine, mis on saadud Kafka'st

Tere, Habr.

Hiljuti jagasin ma oma kogemust selle kohta, milliseid parameetreid me tiimis kĂ”ige sagedamini kasutame Kafka Producentide ja Tarbijate jaoks, et saavutada garanteeritud kohaletoimetamine. Selles artiklis tahan rÀÀkida, kuidas me korraldasime sĂŒndmuse uuesti töötlemise, mis saadi Kafka'st, vĂ€lise sĂŒsteemi ajutise kĂ€ttesaamatuse tĂ”ttu.

Kaasaegsed rakendused töötavad vĂ€ga keerulises keskkonnas. Ärilogika, mis on mĂ€hitud kaasaegsesse tehnoloogilisse kihti, töötab Docker'i pildis, mida haldab koordineerija nagu Kubernetes vĂ”i OpenShift, ning suhtleb teiste rakenduste vĂ”i ettevĂ”tte lahendustega lĂ€bi fĂŒĂŒsiliste ja virtuaalsete marsruuterite ahela. Sellises keskkonnas vĂ”ib alati midagi katki minna, seetĂ”ttu on sĂŒndmuste uuesti töötlemine juhul, kui mĂ”ni vĂ€line sĂŒsteem on kĂ€ttesaamatu, meie Ă€ri protsesside oluline osa.

Kuidas oli enne Kafka

Varem kasutasime projektis IBM MQ asĂŒnkroonseks sĂ”numite edastamiseks. Teenuse töö kĂ€igus tekkinud vea korral vĂ”is saadud sĂ”num paigutada sisemise surnud sĂ”numite jĂ€rjekorda (DLQ) edasiseks kĂ€sitlemiseks. DLQ loodi siseneva jĂ€rjekorra kĂ”rvale, sĂ”numi edasiviimine toimus IBM MQ sees.

Kui viga oli ajutine ja me suudame selle tuvastada (nĂ€iteks ResourceAccessException HTTP-kĂ”ne kĂ€igus vĂ”i MongoTimeoutException MongoDb pĂ€ringus), siis rakendati korduste strateegiat. SĂ”ltumata rakenduse loogika harudest, paigutati algne sĂ”num kas sĂŒsteemijĂ€rjekorda edasilĂŒkatud saatmiseks vĂ”i eraldi rakendusse, mis kunagi ammu loodi sĂ”numite uuesti saatmiseks. Samal ajal kirjutatakse sĂ”numi pĂ€isesse saadetiste korduse number, mis on seotud viivituse intervalliga vĂ”i rakenduse taseme strateegia lĂ”puga. Kui oleme jĂ”udnud strateegia lĂ”ppu, kuid vĂ€line sĂŒsteem on endiselt kĂ€ttesaamatu, siis paigutatakse sĂ”num DLQ-sse kĂ€sitlemiseks.

Lahenduse otsing

Otsides internetist, vĂ”ib leida jĂ€rgmist otsuse. LĂŒhidalt öeldes pakutakse vĂ€lja iga teema jaoks luua topikud igaks viivituse intervalliks ja rakendada rakenduse poolel Tarbijad, kes loevad sĂ”numeid vajaliku viivitusega.

SĂŒndmuste uuesti töötlemine, mis on saadud Kafka'st

Hoolimata paljusid positiivsetest ĂŒlevaadetest, nĂ€ib see mulle mitte just kĂ”ige parem. Esiteks seetĂ”ttu, et arendajal tuleb peale Ă€rinĂ”uete rakendamise kulutada palju aega ka kirjeldatud mehhanismi rakendamisele.

Lisaks, kui Kafka-klastris on aktiveeritud juurdepÀÀsude haldamine, peab kulutama aega teemade loomisele ja nendele vajalike juurdepÀÀsude tagamisele. Sellele lisandub vajadus leida iga retry-teema jaoks Ôige retention.ms parameeter, et sÔnumid jÔuaksid uuesti saata ja ei kaoks. Juhtimise ja juurdepÀÀsutaotluse rakendamine tuleb korrata iga olemasoleva vÔi uue teenuse jaoks.

Vaatame nĂŒĂŒd, milliseid mehhanisme sĂ”numite uuesti töötlemiseks pakub meile Spring tervikuna ja spring-kafka eelkĂ”ige. Spring-kafka-l on ĂŒlekantav sĂ”ltuvus spring-retry'st, mis pakub abstraheerimisvĂ”imalusi erinevate BackOffPolicy haldamiseks. See on ĂŒsna paindlik tööriist, kuid selle oluline puudus on see, et sĂ”numid uuesti saatmiseks salvestatakse rakenduse mĂ€llu. See tĂ€hendab, et rakenduse taaskĂ€ivitamine seoses uuendamise vĂ”i tööea jooksul tekkinud veaga toob kaasa kĂ”ikide uuesti töötlemist ootavate sĂ”numite kadumise. Kuna see punkt on meie sĂŒsteemi jaoks kriitiline, ei kaalunud me seda edasi.

Spring-kafka ise pakub mitmeid ContainerAwareErrorHandler'i rakendusi, nÀiteks SeekToCurrentErrorHandler, millega saab, mitte nihutades offsetit vea tekkimisel, sÔnumi hiljem töödelda. Alates versioonist spring-kafka 2.3 on vÔimalik mÀÀrata BackOffPolicy.

See lĂ€henemine vĂ”imaldab uuesti töödeldavatel sĂ”numitel ellu jÀÀda rakenduse taaskĂ€ivitamise, kuid DLQ mehhanism on endiselt puudulik. Just selle variandi valisime 2019. aasta alguses, optimistlikult pidades, et DLQ-d ei tule vaja (meil vedas ja tĂ”epoolest ei olnud seda mitu kuud rakenduse sellise uuesti töötlemise sĂŒsteemiga). Ajutised vead pĂ”hjustasid SeekToCurrentErrorHandler'i aktiveerimist. ÜlejÀÀnud vead kirjutatakse logisse, pĂ”hjustavad offseti nihkumist ja töötlemine jĂ€tkub jĂ€rgmise sĂ”numiga.

LÔplik lahendus

SeekToCurrentErrorHandler'ile pÔhinev rakendus viitis meid arendada oma mehhanismi sÔnumite edastamiseks uuesti.

Esiteks soovisime kasutada juba olemasolevat teadmist ja laiendada seda rakenduse loogika pĂ”hjal. Lineaarse loogikaga rakenduse jaoks oleks optimaalne peatada uute sĂ”numite lugemine lĂŒhikeseks ajaks, mis on mÀÀratud uuesti kĂ”nede strateegia raames. ÜlejÀÀnud rakenduste puhul soovisime omada ĂŒhte punkti, mis tagaks uuesti kĂ”nede strateegia elluviimise. Lisaks peaks see ĂŒhtne punkt omama DLQ funktsionaalsust mĂ”lema lĂ€henemise jaoks.

Uuesti kÔnede strateegia peaks olema salvestatud rakenduses, mis vastutab jÀrgmise intervalli saamise eest ajutiste vigade korral.

Consumer'i peatamine lineaarse loogikaga rakenduse jaoks

Spring-kafka kasutamise korral vÔib Consumer'i peatamise kood vÀlja nÀha nii:

public void pauseListenerContainer(MessageListenerContainer listenerContainer, 
                                   Instant retryAt) {
        if (nonNull(retryAt) && listenerContainer.isRunning()) {
            listenerContainer.stop();
            taskScheduler.schedule(() -> listenerContainer.start(), retryAt);
            return;
        }
        // DLQ-le
    }

NÀites on retryAt aeg, millal tuleks uuesti kÀivitada MessageListenerContainer, kui see veel töötab. Uuesti kÀivitamine toimub eraldi lÔimes, mis on kÀivitatud TaskScheduler'is, mille rakenduse pakub samuti spring.

Löime retryAt vÀÀrtuse jÀrgmise meetodi abil:

  1. Otsitakse uuesti kÔnede arvu vÀÀrtust.
  2. Vastavalt uuesti kÔnede arvu vÀÀrtusele otsitakse praegune viivituse intervall uuesti kÔnede strateegias. Strateegia kuulutatakse vÀlja rakenduses, mille salvestamiseks valisime JSON formaadi.
  3. JSON-massis leitud intervall sisaldab sekundeid, mille jooksul tuleb töötlemine uuesti korrata. See sekundite arv liidetakse praegusele ajale, moodustades retryAt vÀÀrtuse.
  4. Kui intervalli ei leita, siis retryAt vÀÀrtus on null ja sÔnum saadetakse DLQ-sse kÀsitsi töötlemiseks.

Selle lÀhenemise puhul on vajalik sÀilitada igas sÔnumis, mis on praegu töötlemisel, korduvate katsete arv, nÀiteks rakenduse mÀlu. Katsete arvu sÀilitamine mÀlus ei ole kriitiline, kuna lineaarse loogikaga rakendus ei pruugi kogu töötlemist korraldada. Erinevalt spring-retry'st ei too rakenduse taaskÀivitamine kaasa kÔigi sÔnumite kadumist uuesti töötlemiseks, vaid lihtsalt strateegia taaskÀivitamise.

See lĂ€henemine aitab vĂ€hendada koormust vĂ€lisele sĂŒsteemile, mis vĂ”ib olla vĂ€ga suure koormuse tĂ”ttu kĂ€ttesaamatu. TeisisĂ”nu, lisaks uuesti töötlemisele oleme saavutanud mustri elluviimise. circuit breaker.

Meie puhul on veapiirang vaid 1, ja et minimeerida sĂŒsteemi seiskumist ajutiste vĂ”rgu katkestuste tĂ”ttu, kasutame vĂ€ga granulaarses korduste strateegias lĂŒhikesi viivitusintervalle. See vĂ”ib mitte sobida kĂ”igile ettevĂ”tte grupi rakendustele, seetĂ”ttu tuleb veapiirangu ja intervalli suuruse suhet valida, tuginedes sĂŒsteemi omadustele.

Eraldi rakendus sÔnumite töötlemiseks rakendustelt, millel on mÀÀramatud loogikad.

Siin on nÀide koodist, mis saadab sÔnumi sellesse rakendusse (Retryer), mis tÀidab uuesti saatmise DESTINATION teema juurde, kui RETRY_AT aeg on kÀes:


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Ă€idatud nĂ€itest on nĂ€ha, et palju teavet edastatakse pealkirjades. RETRY_AT vÀÀrtus on samuti sama, mis Consumer’i peatamise kordusmehhanismi puhul. Lisaks DESTINATIONile ja RETRY_AT-ile edastame:

  • GRUP_ID, mille kaudu rĂŒhmitame sĂ”numid kĂ€sitsi analĂŒĂŒsimiseks ja otsingu lihtsustamiseks.
  • ALGE_PARTITSIOON, et pĂŒsida sama Consumeri juures ĂŒmber töötlemiseks. See parameeter vĂ”ib olla null, sel juhul omandatakse uus partitsioon originaalsĂ”numi record.key() vĂ”tme jĂ€rgi.
  • Uuendatud COUNTER vÀÀrtus, et jĂ€rgida korduvate kĂ”nede strateegiat.
  • SEND_TO — konstant, mis nĂ€itab, kas sĂ”num saata korduvaks töötlemiseks RETRY_AT saavutatuna vĂ”i paigutada DLQ-sse.
  • PÕHJUS — pĂ”hjus, miks sĂ”numi töötlemine katkestati.

Retryer salvestab sĂ”numeid korduvaks saatmiseks ja kĂ€sitsi analĂŒĂŒsimiseks PostgreSQL-is. Ajastatult kĂ€ivitatakse ĂŒlesanne, mis leiab sĂ”numid, mille RETRY_AT on saabunud, ja saadab need tagasi originaalsesse partitsiooni DESTINATION teemas record.key() vĂ”tmega.

PĂ€rast sĂ”numi saatmist eemaldatakse need PostgreSQL-ist. SĂ”numite kĂ€sitsi analĂŒĂŒs toimub lihtsas kasutajaliideses, mis suhtleb Retryeriga REST API kaudu. Selle peamised omadused on sĂ”numite uuesti saatmine vĂ”i DLQ-st kustutamine, vigade teabe vaatamine ja sĂ”numite otsimine, nĂ€iteks veateate jĂ€rgi.

Kuna meie klastrites on sisse lĂŒlitatud juurdepÀÀsu haldamine, tuleb tĂ€iendavalt taotleda juurdepÀÀsu teemale, mida kuulab Retryer, ja anda Retryerile vĂ”imalus kirjutada DESTINATION teemas. See on ebamugav, kuid erinevalt ajavahemikus pĂ”hinevast lĂ€henemisest on meil tĂ€ielik DLQ ja UI selle haldamiseks.

On juhtumeid, kus sisenemise teemat loevad erinevad consumer-grupid, mille rakendused rakendavad erinevat loogikat. SĂ”numi korduv töötlemine Retryeri kaudu ĂŒhe sellise rakenduse jaoks pĂ”hjustab teise rakenduse duplikaadi. Selle kaitsmiseks loome eraldi teema korduvaks töötlemiseks. Sisenemis- ja retry-teemat vĂ”ib lugeda sama Consumeriga piiramatu arv kordi.

SĂŒndmuste uuesti töötlemine, mis on saadud Kafka'st

Vaikimisi ei paku see lĂ€henemine circuit breaker’i vĂ”imalust, kuid seda saab rakendusele lisada spring-cloud-netflix vĂ”i uue spring cloud circuit breaker, ĂŒmbritsedes vĂ€liste teenuste kutsumise kohti vastava abstraktsiooniga. Lisaks ilmneb vĂ”imalus valida strateegia bulkhead mustri jaoks, mis vĂ”ib samuti kasulik olla. NĂ€iteks spring-cloud-netflixis vĂ”ib see olla niidi bassein vĂ”i semafor.

KokkuvÔte

Kuna tulemuseks on meil eraldi rakendus, mis vĂ”imaldab sĂ”numi töötlemist korrata, kui mĂ”ni vĂ€line sĂŒsteem on ajutiselt kĂ€ttesaamatu.

Rakenduse ĂŒks peamisi eeliseid on see, et seda saavad kasutada vĂ€list sĂŒsteemi, mis töötab samas Kafka-klistris, ilma oluliste muudatusteta oma poolel! Selline rakendus peab vaid saama juurdepÀÀsu retry-teemale, tĂ€itma mĂ”ned Kafka-pealkirjad ja saatma sĂ”numi Retryerisse. Ei ole vaja luua tĂ€iendavat infrastruktuuri. Ja et vĂ€hendada sĂ”numite edasisaatmist rakendusest Retryerisse ja tagasi, oleme eraldanud rakendused, millel on lineaarsed loogikad, ja viime nende kaudu kordustöötlemise lĂ€bi Consumeri peatamise.

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster