Kuidas Kafka tõelisuseks sai

Kuidas Kafka tõelisuseks sai

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-0

Kui 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:

Kuidas Kafka tõelisuseks sai

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 DefaultPartitioner. 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:

  1. acks = all
  2. retries > 0
  3. enable.idempotence = true
  4. max.in.flight.requests.per.connection ≤ 5 (1 — korraliku saatmise jaoks)
  5. transactional.id = ${application-name}-${hostname}

Consumer:

  1. 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

Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid 🔥 Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster