
Salam, Habr.
Son zamanlarda Kafka Producer və Consumer üçün komandamızın daha çox istifadə etdiyi parametrlər barədə. Bu məqalədə, xarici sistemin müvəqqəti əlçatan olmaması nəticəsində Kafka-dan alınan hadisənin təkrar emalını necə təşkil etdiyimizi anlatmaq istəyirəm.
Müasir tətbiqlər çox mürəkkəb bir mühitdə fəaliyyət göstərir. Müasir texnoloji yığınla sarılan biznes məntiği, Kubernetes və ya OpenShift kimi orkestrator tərəfindən idarə olunan Docker imici içindədir və digər tətbiqlərlə və ya enterprise həlləri ilə fiziki və virtual yönləndiricilər silsiləsi vasitəsilə ünsiyyət qurur. Belə bir mühitdə hər zaman bir şeyin pozulması mümkündür, buna görə də xarici sistemlərdən birinin əlçatan olmaması halında hadisələrin təkrar emalı bizim biznes proseslərimizin mühüm bir hissəsidir.
Kafka-dan əvvəl
Layihədə əvvəlcə asinxron mesaj göndərmə üçün IBM MQ istifadə edirdik. Xidmətin iş prosesi zamanı hər hansı bir xəta baş verdikdə, alınan mesaj dead-letter-queue (DLQ) daxilində yarıdılıb manual təhlil üçün göndərilə bilərdi. DLQ, giriş qovşağının yanında yaradılırdı və mesajın köçürülməsi IBM MQ daxilində baş verirdi.
Əgər xəta müvəqqəti xarakter daşıyırdısa və biz bunu müəyyən edə bilirdiksə (məsələn, HTTP çağırışı zamanı ResourceAccessException və ya MongoDb-yə edilən sorğu zamanı MongoTimeoutException), o zaman təkrar çağırmalar strategiyası devreye girirdi. Tətbiqin məntiqinin şaxələnməsindən asılı olmayaraq, ilkin mesaj ya gecikdirilmiş göndəriş üçün sistem sıraına, ya da bir zamanlar mesajları təkrar göndərmək üçün hazırlanmış ayrı bir tətbiqə köçürülürdü. Bu zaman mesajın başlığına göndəriş sayı yazılırdı ki, bu da gecikmə intervalı və ya tətbiq səviyyəsində strategiyanın sona çatmasına bağlıdır. Strategiyanın sonuna çatmışıqsa, lakin xarici sistem hələ də əlçatan deyilsə, mesaj DLQ-ya manual təhlil üçün göndəriləcək.
Həll axtarışı
, aşağıdakıları tapa bilərsiniz . Qısa desək, hər gecikmə intervallarına bir mövzu açmaq və tətbiq tərəfində müvafiq gecikmə ilə mesajları oxuyan Consumer-ləri həyata keçirmək təklif olunur.

Baxışların çoxsaylı müsbət olmasına baxmayaraq, bu, mənə tamamilə uğurlu görünmür. Birincisi, ona görə ki, inkişaf etdiricinin biznes tələblərini yerinə yetirməkdən əlavə, təsvir olunan mexanizmi həyata keçirmək üçün çox vaxt sərf etməsi lazım olacaq.
Bundan əlavə, əgər Kafka klasterində giriş idarəetməsi aktivdirsə, topiklərin yaradılması və onlara lazım olan girişlərin təmin edilməsi üçün bir qədər vaxt sərf olunmalıdır. Bununla yanaşı, hər bir yenidən cəhd olunan topik üçün retention.ms parametrini düzgün seçmək lazımdır ki, mesajlar təkrar göndərilsin və itirilməsin. Girişləri həyata keçirmək və sorğunu hər bir mövcud və ya yeni xidmət üçün təkrarlamaq lazım olacaq.
Gəlin indi spring-in ümumilikdə və spring-kafka-nın xüsusilə mesajların təkrar emalı üçün təqdim etdiyi mexanizmlərə baxaq. Spring-kafka, müxtəlif BackOffPolicy-ləri idarə etmək üçün abstractions təqdim edən spring-retry-ə asılılıq yaradır. Bu, olduqca çevik bir alətdir, lakin onun əhəmiyyətli bir çatışmazlığı, təkrar göndərilməsi üçün mesajların tətbiqin yaddaşında saxlanılmasıdır. Bu, tətbiqin yeniləməsi və ya istismar zamanı bir səhv səbəbindən yenidən başlaması nəticəsində bütün gözləmədə olan mesajların itirilməsinə səbəb olur. Bu məsələ bizim sistemimiz üçün kritik olduğundan, biz onu daha da araşdırmadıq.
Köməkçi spring-kafka bir neçə ContainerAwareErrorHandler implementasiyası təqdim edir, məsələn , bununla, bir səhv yaranarsa offset-mə dəyişmədən mesajı sonra emal etmək mümkündür. Spring-kafka 2.3 versiyasından başlayaraq BackOffPolicy təyin etmək imkanı təmin edilib.
Bu yanaşma, təkrar emal olunan mesajların tətbiqin yenidən başlamasını yaşamasına imkan tanıyır, lakin DLQ mexanizmi hələ də mövcud deyil. Biz 2019-cu ilin əvvəlində bu variantı seçdik, optimistcəsinə DLQ-ya ehtiyac olmayacağını düşünərək (şansımız gətirdi və gerçəkdən də bir neçə aylıq tətbiq istismarı ərzində DLQ-ya ehtiyac olmadı). Müvəqqəti səhvlər SeekToCurrentErrorHandler-in işə düşməsinə səbəb oldu. Digər səhvlər isə loga yazılır, offset-in dəyişməsinə yol açır və emal növbəti mesajla davam edir.
Nəticə qərarı
SeekToCurrentErrorHandler-a əsaslanan tətbiqetmənin həyata keçirilməsi, bizə öz mesajların yenidən göndərilməsi mexanizmini inkişaf etdirmək üçün ilham verdi.
Əvvəlcə biz mövcud təcrübəni istifadə etmək və onu tətbiqin məntiqinə görə genişləndirmək istəyirdik. Xətti məntiqə malik olan tətbiq üçün yeni mesajların oxunmasının dayandırılması üçün kiçik bir müddət təyin etmək optimal olardı, bu müddət isə təkrar çağırma strategiyası çərçivəsində müəyyən edilir. Digər tətbiqlər üçün isə weblayihəyə təmin etməsi üçün tək bir nöqtənin olması arzuolunandır ki, bu, iki yanaşma üçün təkrar çağırma strategiyasını həyata keçirir. Bununla yanaşı, bu tək nöqtə həm də hər iki yanaşma üçün DLQ funksionallığına sahib olmalıdır.
Təkrar çağırma strategiyası, müvəqqəti xətanın baş verməsi zamanı növbəti intervalları almaqdan məsul olan tətbiqin içində saxlanmalıdır.
Xətti məntiqə malik tətbiq üçün Consumer-in dayanması
spring-kafka ilə işlədikdə, Consumer-in dayanma kodu təxminən belə görünə bilər:
public void pauseListenerContainer(MessageListenerContainer listenerContainer,
Instant retryAt) {
if (nonNull(retryAt) && listenerContainer.isRunning()) {
listenerContainer.stop();
taskScheduler.schedule(() -> listenerContainer.start(), retryAt);
return;
}
// DLQ-ya göndərMüxtəlif retryAt - MessageListenerContainer-in yenidən işə salınması üçün vaxtdır, əgər o, hələ də işləyirsə. Yenidən başlaması TaskScheduler-də gerçekleştirilen ayrı bir ipdə baş verəcək, onun implementasiyasını da spring təklif edir.
retryAt dəyərini aşağıdakı kimi tapırıq:
- Yenidən çağırmaların sayğacının dəyəri axtarılır.
- Sayğacın dəyərinə uyğun olaraq, təkrar çağırma strategiyasında cari gecikmə intervalları axtarılır. Strategiyamız tətbiqin içində elan edilir, onun saxlanması üçün JSON formatını seçdik.
- JSON massivində tapılan intervallar, emalın təkrarı üçün lazım olan saniyələrin sayını ehtiva edir. Bu sayı cari vaxta əlavə olunur, nəticədə retryAt üçün dəyər meydana gəlir.
- Əgər interval tapılmasa, onda retryAt dəyəri null olur və mesaj DLQ-ya göndərilir ki, bu da əl ilə analiz ediləcək.
Bu yanaşmada, indiki zaman üçün emal olunan hər mesaj üçün hər dəfə çağırmaların sayını saxlamaq lazımdır, məsələn, tətbiqin yaddaşında. Cəhd sayını yaddaşda saxlamaq bu yanaşma üçün kritik deyil, çünki xətti məntiqə malik bir tətbiq ümumilikdə emalı həyata keçirə bilmir. Spring-retry ilə fərqli olaraq, tətbiqin yenidən başladılması, bütün mesajların təkrar emal üçün itirilməsinə səbəb olmur, yalnız strategiyanın yenidən başlamasına səbəb olur.
Bu yanaşma, çox böyük yük səbəbindən mövcud olmaya bilən xarici sistemin yükünü azaltmağa kömək edir. Başqa sözlə, təkrar emala əlavə olaraq, biz .
Bizim halımızda, səhv eşiklərin sayı cəmi 1-dir və müvəqqəti şəbəkə kesintisi səbəbiylə sistemin dayanıqlığını minimalizə etmək üçün, biz kiçik gecikmə intervalları ilə çox detallaşdırılmış bir təkrar çağırma strategiyası istifadə edirik. Bu, şirkət qrupunun bütün tətbiqləri üçün uyğun olmaya bilər, buna görə də səhv eşikləri ilə interval ölçüsü arasındakı nisbət sistemin xüsusiyyətlərinə əsaslanaraq tənzimlənməlidir.
Dəyişkən məntiqə sahib tətbiqlərdən mesajları emal etmək üçün ayrıca bir tətbiq
Bu tətbiqə (Retryer) mesaj göndərmək üçün bir kod nümunəsi:
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ümunədən göründüyü kimi, başlıqlarda çoxlu məlumat ötürülür. RETRY_AT dəyəri, istehsalçı dayandırılması yolu ilə təkrar mexanizmi üçün olduğu kimi qalır. DESTINATION və RETRY_AT-dan əlavə, biz ötürürük:
- GROUP_ID, mesajları əl ilə analiz üçün qruplaşdırmaq və axtarışı asanlaşdırmaq üçün.
- ORIGINAL_PARTITION, eyni istehsalçını təkrar emal üçün saxlamağa çalışmaq üçün. Bu parametr null ola bilər, belə olduqda yeni parti açıq mesajın record.key() açarına görə alınacaq.
- Təkrar çağırma strategiyasını izləmək üçün güncellenmiş COUNTER dəyəri.
- SEND_TO — RETRY_AT-a çatdıqda mesajı təkrar emal üçün göndərmək, ya da DLQ-ya yerləşdirmək üçün olan konstant.
- REASON — mesajın emalının niyə dayandırıldığını bildirən səbəb.
Retryer, PostgreSQL-də təkrar göndərmə və əl ilə təhlil üçün mesajları saxlayır. Zamanlayıcıya əsasən, RETRY_AT-ı keçmiş mesajları axtarıb onları DESTINATION topikinin ORIGINAL_PARTITION partiyasına record.key() açarı ilə geri göndərən bir tapşırıq işə salınır.
Mesaj yollandıqdan sonra PostgreSQL-dən silinir. Mesajların əl ilə təhlili, Retryer ilə REST API vasitəsilə əlaqələndirilmiş sadə UI-da həyata keçirilir. Onun əsas xüsusiyyətləri DLQ-dan mesajların yenidən göndərilməsi və ya silinməsi, xəta məlumatlarını baxmaq və məsələn, xəta adları üzrə mesajları axtarmaqdır.
Bizim klasterlərdə giriş idarəetməsi aktiv olduğu üçün, Retryer’in dinlədiyi topik üçün əlavə icazələr istəmək və Retryer’in DESTINATION topikində yazma imkanı vermək lazımdır. Bu narahatdır, amma intervalla topik yanaşmasından fərqli olaraq, tam DLQ və onunla idarə etmə üçün bir UI əldə etmiş oluruq.
Bəzi hallarda, daxil olan topik bir neçə fərqli istehsalçı qrupları tərəfindən oxunur, onların proqramları fərqli məntiqi həyata keçirir. Retryer vasitəsi ilə bir proqram üçün mesajın təkrar emalı digər proqramda dublikat yaradacaq. Bunu qorumaq üçün, biz bir ayrı təkrar emal topiki yaradırıq. Daxil olan və təkrar-emal topiklərini eyni istehsalçı hər hansı məhdudiyyət olmadan oxuya bilər.

Default olaraq, bu yanaşma circuit breaker imkanı təqdim etmir, lakin bunu tətbiqə əlavə etmək mümkündür ya da yeni , xarici xidmətlərin çağırış yerlərini müvafiq abstraksiyalara sararaq. Bundan əlavə, seçə biləcəyiniz strategiya imkanı yaranır nümunənin, bu da yararlı ola bilər. Məsələn, spring-cloud-netflix-də bu, thread pool və ya semaphore ola bilər.
Çıktı
Nəticədə, vaxtilə xarici sistemlərdən birinin müvəqqəti əlçatan olmaması zamanı mesajın emalını təkrarlamağa imkan verən ayrıca bir tətbiq əldə etdik.
Tətbiqin əsas üstünlüklərindən biri odur ki, eyni Kafka-klasterində işləyən xarici sistemlər ondan əhəmiyyətli dəyişikliklər etmədən istifadə edə bilərlər! Bu tətbiq yalnız retry-topikə giriş əldə etməli, bir neçə Kafka başlığını doldurmalı və mesajı Retryer-ə göndərməlidir. Heç bir əlavə infrastruktur yaratmaq lazım deyil. Mesajların tətbiqdən Retryer-ə və geri ötürülməsinin sayını azaltmaq üçün, xətti məntiqli tətbiqləri ayırdıq və onlarda istehlakçının dayanması vasitəsilə təkrar emalı etdik.
Mənbə: habr.com
