Cum a devenit Kafka realitate

Cum a devenit Kafka realitate

Salut, Habr!

Lucrez în echipa Tinkoff, care se ocupă cu dezvoltarea propriei centre de notificări. În cea mai mare parte, dezvolt pe Java folosind Spring Boot și rezolv diverse probleme tehnice care apar în proiect.

Cele mai multe dintre microserviciile noastre interacționează asincron între ele prin intermediul unui broker de mesaje. Anterior, foloseam IBM MQ ca broker, care nu mai făcea față încărcăturii, dar oferea în același timp garanții ridicate de livrare.

Ca alternativă, ne-a fost propus Apache Kafka, care are un potențial ridicat de scalare, dar, din păcate, necesită practic o abordare individualizată pentru configurarea în diferite scenarii. În plus, mecanismul de livrare at least once, care funcționează în Kafka în mod implicit, nu permite susținerea nivelului necesar de consistență din cutie. Mai departe, voi împărtăși experiența noastră de configurare a Kafka, în special voi explica cum să configurăm și să trăim cu exactly once delivery.

Livrare garantată și nu doar atât

Parametrii despre care vom discuta în continuare te vor ajuta să previi o serie de probleme cu setările de conectare implicite. Dar mai întâi, aș dori să acord atenție unui parametru care va facilita orice debugg.

Acest lucru te va ajuta client.id pentru Producer și Consumer. La prima vedere, ca valoare, poți folosi numele aplicației, și în majoritatea cazurilor, asta va funcționa. Totuși, situația în care aplicația folosește mai mulți Consumers și le aloci același client.id duce la următorul avertisment:

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

Dacă vrei să folosești JMX într-o aplicație cu Kafka, aceasta ar putea fi o problemă. Pentru acest caz, cel mai bine este să folosești ca valoare client.id o combinație între numele aplicației și, de exemplu, numele topic-ului. Rezultatul configurației noastre poate fi văzut în output-ul comenzii kafka-consumer-groups din utilitarele de la Confluent:

Cum a devenit Kafka realitate

Acum să analizăm scenariul livrării garantate a mesajelor. Un Kafka Producer are un parametru acks, care permite configurarea după câte acknowledge liderul cluster-ului trebuie să considere mesajul ca fiind scris cu succes. Acest parametru poate lua următoarele valori:

  • 0 — acknowledge-urile nu vor fi considerate.
  • 1 — parametrul implicit, necesită acknowledge doar de la 1 replică.
  • −1 — necesare acknowledgments de la toate replicile sincronizate (configurarea cluster-ului min.insync.replicas).

Din valorile enumerate, se observă că acks egal cu −1 oferă cele mai puternice garanții că mesajul nu se va pierde.

Cum știm cu toții, sistemele distribuite sunt nesigure. Pentru a ne proteja de defecțiuni temporare, Kafka Producer oferă parametrul retries, care permite stabilirea numărului de încercări de retransmitere în decursul delivery.timeout.ms. Deoarece parametrul retries are valoarea implicită Integer.MAX_VALUE (2147483647), numărul de retransmiteri ale mesajului poate fi reglat modificând doar delivery.timeout.ms.

Să ne îndreptăm spre livrarea exactly once

Setările enumerate permit Producer-ului nostru să livreze mesaje cu o garanție ridicată. Acum să discutăm despre cum putem garanta scrierea unei singure copii a mesajului în topicul Kafka? În cazul cel mai simplu, pentru aceasta, Producer-ul trebuie să seteze parametrul enable.idempotence la true. Idempotenta garantează scrierea unui singur mesaj într-o anumită partiție a unui topic. O precondiție pentru activarea idempotentei este setarea valorilor acks = all, retry > 0, max.in.flight.requests.per.connection ≤ 5. Dacă acești parametri nu sunt definiți de dezvoltator, valorile menționate mai sus vor fi setate automat.

Când idempotenta este configurată, este esențial să ne asigurăm că mesajele identice ajung întotdeauna în aceleași partiții. Acest lucru se poate realiza prin configurarea cheii și a parametrului partitioner.class pe Producer. Să începem cu cheia. Pentru fiecare trimitere, aceasta ar trebui să fie constantă. Este ușor de realizat folosind un identifiant de afaceri din mesajul original. Parametrul partitioner.class are o valoare implicită — DefaultPartitioner. Cu această strategie de partajare implicită, procedăm astfel:

  • Dacă partiția este specificată explicit la trimiterea mesajului, o folosim.
  • Dacă partiția nu este specificată, dar cheia este indicată — alegem partiția pe baza hash-ului cheii.
  • Dacă nici partiția, nici cheia nu sunt specificate — alegem partițiile în mod secvențial (round-robin).

În plus, utilizarea cheii și a trimiterii idempotente cu parametrul max.in.flight.requests.per.connection = 1 îți oferă un proces ordonat de gestionare a mesajelor pe Consumer. Este important de reținut că, dacă pe clusterul tău este activat controlul accesului, va fi nevoie de permisiuni pentru scrierea idempotentă în topic.

Dacă simți că îți lipsesc funcționalitățile de trimitere idempotentă pe cheie sau logica pe partea Producer necesită menținerea consistenței datelor între diferite partiții, atunci tranzacțiile te vor ajuta. În plus, prin intermediul unei tranzacții în lanț, este posibil să sincronizezi condiționat scrierea în Kafka, de exemplu, cu scrierea în baza de date. Pentru a activa trimiterea tranzacțională pe Producer, acesta trebuie să fie idempotent și să definești suplimentar transactional.id. Dacă pe clusterul tău Kafka este activat controlul accesului, va fi nevoie de permisiuni de scriere pentru scrierea tranzacțională, la fel ca pentru scrierea idempotentă, care pot fi furnizate printr-un model folosind valoarea stocată în transactional.id.

Formal, ca identificator al tranzacției, poate fi utilizată orice urmă, cum ar fi numele aplicației. Dar dacă rulezi mai multe instanțe ale aceleași aplicații cu același transactional.id, prima instanță lansată va fi oprită cu o eroare, deoarece Kafka o va considera un proces zombie.

org.apache.kafka.common.errors.ProducerFencedException: Producer a încercat să execute o operațiune cu o epocă veche. Ori există un producer mai nou cu același transactionalId, ori tranzacția producer-ului a expirat de către broker.

Pentru a rezolva această problemă, adăugăm la numele aplicației un sufix bazat pe numele gazdei, pe care îl obținem din variabilele de mediu.

Producer-ul este configurat, dar tranzacțiile în Kafka controlează doar domeniul de vizibilitate al mesajului. Indiferent de starea tranzacției, mesajul ajunge imediat în topic, dar are atribute sistemice suplimentare.

Pentru ca aceste mesaje să nu fie citite prea devreme de către Consumer, acesta trebuie să stabilească parametrul isolation.level la valoarea read_committed. Un astfel de Consumer va putea citi mesajele netranzacționale ca înainte, iar pe cele tranzacționale doar după comitere.
Dacă ai setat toate configurările enumerate anterior, atunci ai configurat exactly once delivery. Felicitări!

Dar mai există un detaliu. Transactional.id pe care l-am configurat mai sus este, de fapt, un prefix al tranzacției. La managerul de tranzacții, i se adaugă un număr secvențial. Identificatorul obținut este eliberat la transactional.id.expiration.ms, care este configurat pe un cluster Kafka și are o valoare implicită de „7 zile”. Dacă în acest timp aplicația nu a primit nicio mesaj, atunci la următoarea încercare de trimitere a unei tranzacții, veți primi InvalidPidMappingException. După aceasta, coordonatorul tranzacțiilor va atribui un nou număr secvențial pentru următoarea tranzacție. În acest caz, mesajul poate fi pierdut dacă InvalidPidMappingException nu este gestionată corect.

În loc de concluzii

Așa cum puteți observa, nu este suficient doar să trimiteți mesaje în Kafka. Trebuie să alegeți o combinație de parametri și să fiți pregătiți să faceți rapid modificări. În acest articol, am încercat să explic în detaliu configurarea livrării exactly once și am descris câteva probleme legate de configurarea client.id și transactional.id cu care ne-am confruntat. Mai jos sunt prezentate pe scurt setările pentru Producer și Consumer.

Producer:

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

Consumer:

  1. isolation.level = read_committed

Pentru a minimiza erorile în aplicațiile viitoare, am creat un wrapper peste configurația spring, unde sunt deja setate valori pentru anumite din parametrii enumerați.

Iată câteva materiale pentru auto-studiu:

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster