
Tere, Habr!
Ma töötan Tinkoffi meeskonnas, mis arendab oma teavituskeskust. Peamiselt arendan Java-s, kasutades Spring Booti, ja lahendan erinevaid tehnilisi probleeme, mis projektis esinevad.
Enamik meie mikroteenustest suhtleb omavahel asĂŒnkroonselt sĂ”numite vahendaja kaudu. Varem kasutasime selleks IBM MQ-d, mis enam ei suutnud koormust taluda, kuid millel olid kĂ”rged saatmisgarantiid.
Asenduseks pakuti meile Apache Kafka, mis on suurepÀraste skaleeritavuse vÔimalustega, kuid vajab kahjuks peaaegu individuaalset lÀhenemist konfigureerimisele erinevates stsenaariumides. Lisaks ei vÔimaldanud Kafka vaikimisi töötav vÀhemalt korra saatmise mehhanism tagada vajalikku kooskÔlastatust. JÀrgmises jagan meie kogemusi Kafka konfigureerimisel, eriti kuidas seadistada ja elada tÀpselt korra saatmisega.
Tagatud saatmine ja mitte ainult
JĂ€rgnevad seaded aitavad vĂ€ltida mitmeid vaikimisi ĂŒhenduskonfiguratsiooniga seotud probleeme. Kuid esmalt soovin tĂ€helepanu pöörata ĂŒhele seadmele, mis hĂ”lbustab vĂ”imalikku tĂ”rkeotsingut.
Selleks aitab client.id tootja ja tarbija jaoks. Esmapilgul vÔib vÀÀrtusena kasutada rakenduse nime, ja enamikul juhtudel töötab see. Kuid olukord, kus rakenduses on mitu tarbijat ja mÀÀrate neile sama client.id, viib jÀrgmise hoiatuse tekkimiseni:
org.apache.kafka.common.utils.AppInfoParser â Viga AppInfo mbean registreerimisel javax.management.InstanceAlreadyExistsException: kafka.consumer:type=app-info,id=kafka.test-0Kui soovite kasutada JMX-i Kafka rakenduses, siis vĂ”ib see olla probleem. Sel juhul on kĂ”ige parem kasutada client.id vÀÀrtusena rakenduse nime ja nĂ€iteks teema kombinatsiooni. Meie konfiguratsiooni tulemust saab vaadata kĂ€su kafka-consumer-groups Confluenti utiliidist:

NĂŒĂŒd vaatame ĂŒle sĂ”numi tagamisel toimimise stsenaariumi. Kafka tootjal on seadistuse acks, mis vĂ”imaldab mÀÀrata, pĂ€rast kui palju kinnitusi tuleb klastrijuhi poolt loetav sĂ”num edastatuks lugeda. See seade vĂ”ib vĂ”tta jĂ€rgmisi vÀÀrtusi:
- 0 â kinnitusi ei arvestata.
- 1 â vÀÀrtus mĂ”istes, vaja kinnitust ainult ĂŒhe replikalt.
- â1 â vajalikud kinnitused kĂ”igilt sĂŒnkroniseeritud replikalt (klastri seadistus) min.insync.replicas).
Nendest vÀÀrtustest on selge, et acks, mis on vĂ”rdne â1, annab kĂ”ige tugevamad garantiid, et sĂ”num ei kao.
Kuidas me kĂ”ik teame, jaotatud sĂŒsteemid pole usaldusvÀÀrsed. Aja jooksul tekkinud tĂ”rgete eest kaitsmiseks pakub Kafka tootja seadet retries, mis vĂ”imaldab mÀÀrata uuspostituste arvu delivery.timeout.ms. Kuna retries seadistus on vaikimisi Integer.MAX_VALUE (2147483647), saab sĂ”numi korduspostitamist reguleerida lihtsalt delivery.timeout.ms vÀÀrtust muutes.
Liigume tÀpselt korra saatmise poole
Toodud seadistused vĂ”imaldavad meie tootjal sĂ”numeid kĂ”rge garantii tasemega edastada. RÀÀgime nĂŒĂŒd, kuidas tagada, et Kafka teemas oleks salvestatud ainult ĂŒks sĂ”numikohane koopia? Lihtsaimatel juhtudel peab selle saavutamiseks tootjas olema seatud seade enable.idempotence vÀÀrtuseks true. Idempotentsus tagab, et kindlasse teema partitsiooni salvestatakse ainult ĂŒks sĂ”num. Idempotentsuse lubamise eeltingimuseks on vÀÀrtused acks = all, retry > 0, max.in.flight.requests.per.connection †5. Kui neid seadistusi ei mÀÀra arendaja, seadistatakse automaatselt ĂŒlaltoodud vÀÀrtused.
Kui idempotentsus on seatud, tuleb tagada, et samad sÔnumid lÀheksid iga kord samadesse partitsioonidesse. Seda saab saavutada, seadistades vÔtme ja parameetri partitioner.class tootjas. Alustame vÔtmega. Iga saatmine peab olema sama. Seda on lihtne saavutada, kasutades mingit Àriidentifikaatorit algsest sÔnumist. Parameeter partitioner.class on vaikimisi vÀÀrtusega . Selle pÔhilise partitsioneerimisstrateegia puhul toimime jÀrgmiselt:
- Kui partitsioon on sÔnumi saatmise ajal selgelt mÀÀratud, kasutame seda.
- Kui partitsioon ei ole mÀÀratud, kuid vĂ”tme mÀÀrati â valime partitsiooni vĂ”tme hash'i jĂ€rgi.
- Kui ei ole mÀÀratud ei partitsiooni ega vĂ”tme â valime partitsioonid jĂ€rjestikku (round-robin).
Lisaks, vÔtme ja idempotentse saatmise kasutamine parameetriga max.in.flight.requests.per.connection = 1 annab teile struktureeritud sÔnumite töötlemise Consumeris. Tuleb eraldi meeles pidada, et kui teie klastris on seadistatud juurdepÀÀsu haldus, siis vajate Ôigusi idempotentse kirje tegemiseks teemasse.
Kui teile ei piisa idempotentse saatmise vĂ”imalustest vĂ”tme jĂ€rgi vĂ”i kui tootja poolne loogika nĂ”uab andmete konsistentsuse sĂ€ilitamist erinevate partitsioonide vahel, siis tulevad appi tehingud. Lisaks vĂ”imaldab aheltehingud tingimuslikult sĂŒnkroonida kirje Kafka-s, nĂ€iteks andmebaasi kirje tegemisega. Tehingute saatmise lubamiseks peab tootja olema idempotentne ja lisaks seadistama transactional.id. Kui teie Kafka klaster on seadistatud juurdepÀÀsu halduseks, siis tehinguliste kirjete jaoks, nagu ka idempotentsete jaoks, on vajalikud kirjutamise Ă”igused, mis vĂ”ivad olla antud maskiga, kasutades vÀÀrtust, mis on salvestatud transactional.id-s.
Formaalsetele tehingu identifikaatorina vÔib kasutada mis tahes stringi, nÀiteks rakenduse nime. Kuid kui kÀitate mitu instantsi sama rakenduse sama transactional.id-ga, siis esimene kÀivitatud instants peatatakse vea tÔttu, kuna Kafka peab seda zombiprotsessiks.
org.apache.kafka.common.errors.ProducerFencedException: Tootja ĂŒritas toimingut, millel oli vana epohh. Kas on olemas uuem tootja sama transactionalId-ga vĂ”i on tootja tehing aegunud vahendaja poolt.Selle probleemi lahendamiseks lisame rakenduse nimele suffiksi, mis koosneb hostinimest, mille saame keskkonnamuutujatest.
Tootja on seadistatud, kuid tehingud Kafka-s haldavad ainult sĂ”numi nĂ€htavuse ulatust. Olenemata tehingu seisundist satub sĂ”num kohe teemasse, kuid omab lisaks sĂŒsteemi atribuute.
Et sellised sÔnumid ei loetaks Consumer'iga liiga vara, peab ta seadistama parameetri isolation.level vÀÀrtuseks read_committed. Selline Consumer vÔib lugeda mitte-tehingulisi sÔnumeid nagu varem, kuid tehingulisi ainult pÀrast kinnitust.
Kui olete seadistanud kÔik eelnevalt loetletud seaded, siis olete seadistanud exactly once delivery. Palju Ônne!
Kuid on veel ĂŒks nĂŒanss. Transactional.id, mida me ĂŒlal seadistame, on tegelikult tehingu prefix. Tehingu haldur lisab sellele jĂ€rjestusnumbri. Saadud identifikaator antakse transactional.id.expiration.ms, mis on seadistatud Kafka klastris ja millel on vaikimisi vÀÀrtus â7 pĂ€evaâ. Kui selle aja jooksul rakendus ei saa ĂŒhtegi sĂ”numit, siis jĂ€rgmise tehingulise saatmise katse korral saate InvalidPidMappingException. PĂ€rast seda vĂ€ljastab tehingu koordinaator jĂ€rgmise tehingu jaoks uue jĂ€rjestusnumbri. Sellega vĂ”ib sĂ”num olla kaduma lĂ€inud, kui InvalidPidMappingException'i pole Ă”igesti kĂ€sitletud.
KokkuvÔtte asemel
Nagu vĂ”ib nĂ€ha, ei piisa lihtsalt sĂ”numite saatmisest Kafka-sse. Tuleb valida parameetrite kombinatsioon ja olla valmis kiirete muudatuste tegemiseks. Selles artiklis pĂŒĂŒdsin ĂŒksikasjalikult nĂ€idata exactly once delivery seadistust ja kirjeldasin mitmeid probleeme client.id ja transactional.id konfigureerimisel, millega me silmitsi seisisime. Allpool on lĂŒhidalt toodud tootja ja consumer seadistused.
Tootja:
- acks = all
- retries > 0
- enable.idempotence = true
- max.in.flight.requests.per.connection †5 (1 â korraliku saatmise jaoks)
- transactional.id = ${application-name}-${hostname}
Consumer:
- isolation.level = read_committed
Vigade minimeerimiseks tulevastes rakendustes oleme teinud oma mÀhise, mis pÔhineb spring-konfiguratsiooni peal, kus on juba seatud vÀÀrtused mÔnedele loetletud parameetritele.
Siin on paar materjali iseseisvaks Ôppimiseks:
Allikas: habr.com
