Kiçik bir kitabın tərcüməsinin davamı:
«Mesaj Brokerlərini Anlamaq»,
müəllif: Jakub Korab, nəşriyyat: O’Reilly Media, Inc., nəşr tarixi: İyun 2017, ISBN: 9781492049296.
Əvvəlki tərcümə olunmuş hissə:
FƏSİL 3
Kafka
Kafka, LinkedIn-də ənənəvi mesaj brokerlərinin bəzi məhdudiyyətlərini aşmaq və müxtəlif "nöqtə-nöqtə" qarşılıqlı əlaqələr üçün bir neçə mesaj brokeri tənzimləməkdən qaçmaq məqsədilə hazırlanmışdır. Bu kitabda "Şaquli və üfüqi genişləndirmə" bölməsində izah edilən şəraitdə. LinkedIn-də istifadə ssenariləri əsasən tək yönlü olaraq çox böyük məlumat axınlarının, məsələn, səhifə klikləri və giriş jurnallarının yığılması üzərində qurulmuşdur, eyni zamanda bu məlumatların bir neçə sistemdən istifadə edilməsinə imkan tanımışdır, istehsalçıların və digər istehlakçıların performansına təsir etmədən. Əslində, Kafkanın var olma səbəbi, Universal Data Pipeline-ın təsvir etdiyi mesaj mübadiləsi arxitekturasını əldə etməkdir.
Bu son məqsəd nəzərə alındıqda, təbii olaraq başqa tələblər də meydana çıxdı. Kafka:
- Son dərəcə sürətli olmalıdır
- Mesajlarla işləyərkən çoxlu keçid təmin etməlidir
- «Nəşrçi-Abunəçi» və «Nöqtə-Nöqtə» modellərini dəstəkləməlidir
- İstehlakçıların artırılması ilə yavaşlamamalıdır. Məsələn, ActiveMQ-dakı istehlakçı sayı artdıqca, həm növbələr, həm də mövzuların performansı pisləşir
- Üfüqi şəkildə genişləndirilə bilən olmalıdır; bir broker, mesajları yalnız disk sürətinin maksimum sürəti ilə saxlayırsa, performansı artırmaq üçün bir broker instansiyasının sərhədini aşmağı düşünmək mənalıdır
- Mesajların saxlanmasına və yenidən alınmasına giriş hüquqlarını ayrılmalıdır
Bütün bunlara nail olmaq üçün Kafkada müştərilər və mesaj brokerləri arasındakı rolları və öhdəlikləri yenidən müəyyən edən bir arxitektura qəbul edilmişdir. JMS modeli, mesajların yayılmasında brokerin məsul olduğu, müştərilərin yalnız mesaj göndərmək və qəbul etmək məsələsi ilə məşğul olması baxımından çox yönümlüdür. Kafkasa, müştəri mərkəzli yanaşmaya malikdir və müştəri, istehlakçılar arasında müvafiq mesajların ədalətli bölüşdürülməsi kimi ənənəvi bir brokerin bir çox funksiyalarını öz üzərinə götürərək, son dərəcə sürətli və genişlənə bilən bir broker alır. Ənənəvi mesaj mübadiləsi sistemləri ilə işləyən insanlar üçün Kafkada işləmək, baxışlarda fundamental dəyişikliklər tələb edir.
Bu mühəndislik sahəsi, adi broker ilə müqayisədə mesaj mübadiləsinin tutumunu bir neçə dəfə artırmağı bacaran bir infrastrukturun yaradılmasına səbəb oldu. Gördüyümüz kimi, bu yanaşma müəyyən yük tipləri və quraşdırılmış proqram üçün Kafka-nın uyğun olmadığı kompromislərlə birlikdə gündəmə gəlir.
Birləşdirilmiş ünvan modeli
Yuxarıda təsvir olunan tələbləri yerinə yetirmək üçün Kafka, "nəşr-abunə" və "nöqtə-nöqtə" mesaj mübadiləsini bir növ ünvan çərçivəsində birləşdirir — mövzu. Bu, mesaj mübadiləsi sistemləri ilə işləyən insanları çaşdırır, burada "mövzu" sözü yayımlama mexanizminə müraciət edir, bu mexanizmdən (mövzudan) oxuma reliyabil deyil (dayanıqlı deyil). Kafka-nın mövzularını bu kitabın girişində verilmiş tərifə uyğun olaraq hibrid növ ünvan olaraq həyata keçirmək olar.
Bu fəsilin qalan hissəsində, əgər biz açıq şəkildə başqa cür göstərməsək, "mövzu" termini Kafka mövzusuna aiddir.
Mövzuların necə davrandığını və hansı təminatları verdiyini tam başa düşmək üçün əvvəlcə onların Kafka-da necə tətbiq edildiyini nəzərdən keçirməliyik.
Kafka-da hər bir mövzunun öz jurnalısı var.
Kafka-ya mesaj göndərən istehsalçılar bu jurnala əlavə edirlər, alıcılar isə jurnaldan göstəricilər vasitəsilə oxuyurlar ki, bu göstəricilər davamlı olaraq irəliləyir. Vaxt-vaxtında Kafka, jurnaldakı ən köhnə hissələri silir, buna görə də bu hissələrdəki mesajlar oxunub-oxunmamasına baxmayaraq. Kafka-nın dizaynının mərkəzi hissəsi, brokerin mesajların oxunub-oxunmadığını düşünməməsi — bu müştərinin məsuliyyətidir.
Jurnal və göstərici terminləri bu yaxşı tanınmış terminlər burada anlayışı asanlaşdırmaq məqsədilə istifadə olunur.
Bu model ActiveMQ-dən tam fərqlidir, burada bütün növbələrdən olan mesajlar bir jurnalda saxlanır, broker isə mesajları oxunduqdan sonra silinmiş kimi işarələyir.
Gəlin indi mövzunun jurnalını daha dərindən öyrənək.
Kafka jurnalı bir neçə partiyadan ibarətdir (). Kafka, hər bir partidada sərt sıralanmanı təmin edir. Bu, bir partiyada müəyyən bir sıra ilə yazılan mesajların eyni sırada oxunacağı deməkdir. Hər bir partiya dövri (rolling) jurnal faylı olaraq həyata keçirilir və alt dəst (subset) öz istehsalçıları tərəfindən göndərilən mesajların hamısından ibarətdir. Yaradılan mövzu əslində bir partiyaya malikdir. Partiyaların ideyası Kafka-nın üfüqi miqyaslanma üçün mərkəzi ideyasıdır.

Şəkil 3-1. Kafka Partiyaları
Bir istehsalçı Kafka mövzusuna mesaj göndərərkən, o, mesajı hansı partiyaya göndərəcəyinə qərar verir. Bunu daha ətraflı müzakirə edəcəyik.
Mesajları oxumaq
Müştəri, mesajları oxumaq istəyən, adlandırılmış bir göstəricini idarə edir. istehlakçı qrupları (consumer group), hansı ki, işarə edir. ofset (offset) partisyonlardakı mesajların. Ofset, partisiya başlanğıcında 0-dan başlayan artan nömrəsi olan bir mövqedir. API-yə istifadəçi tərəfindən müəyyən edilən group_id vasitəsilə istinad edilən bu istehlakçı qrupları bir loji istehlakçıya və ya sistemə uyğundur..
Mesajlaşma mübadiləsi edən əksər sistemlər, mesajların paralel işlənməsi üçün bir neçə nümunə və axın vasitəsilə məlumat oxuyur. Beləliklə, adətən eyni istehlakçı qrupunu birgə istifadə edən bir çox istehlakçı nümunəsi olacaq.
Oxuma problemi aşağıdakı kimi təqdim edilə bilər:
- Mövzu bir neçə partisyona malikdir.
- Eyni anda bir çox istehlakçı qrupları mövzudan istifadə edə bilər.
- İstehlakçı qrupu bir neçə ayrı nümunəyə malik ola bilər.
Bu, «çoxdan çox» kompleks bir problemdir. Kafka-nın istehlakçı qrupları, istehlakçı nümunələri və partisyalar arasındakı münasibətləri necə idarə etdiyini başa düşmək üçün tədricən çətinləşən oxuma ssenarilərini nəzərdən keçirək.
İstehlakçılar və istehlakçı qrupları
Giriş nöqtəsi kimi tək partisiya ilə bir mövzunu götürək ().

Şəkil 3-2. İstehlakçı partisiyadan oxuyur.
İstehlakçı nümunəsi bu mövzuya öz group_id ilə qoşulduqda, ona oxumaq üçün bir partisiya və bu partisiyada ofset təyin edilir. Bu ofsetin mövqeyi müştəridə ən son mövqedən (ən yeni mesaj) və ya ən erkən mövqedən (ən köhnə mesaj) göstərici kimi konfiqurasiya edilir. İstehlakçı mövzudan mesajları sorğu edir (polls), bu da onların yazıdan ardıcıl şəkildə oxunmasına səbəb olur.
Ofsetin mövqeyi mütəmadi olaraq Kafka-ya geri komit edilir və daxili mövzuda saxlanılır. _istehlakçı_ofsetləri. Oxunan mesajlar hələ də silinmir, adi brokerdən fərqli olaraq, müştəri artıq baxılan mesajları yenidən işləmək üçün ofseti geri çevirə bilər.
İkinci loji istehlakçı fərqli group_id istifadə edərək qoşulduqda, o, birincidən asılı olmayan ikinci bir göstəricini idarə edir (). Beləliklə, Kafka mövzusu, yalnız bir istehlakçının olduğu bir növbə kimi fəaliyyət göstərir, eyni zamanda bir neçə istehlakçının abunə olduğu adi bir nəşr-abunə (pub-sub) mövzusu kimi, bütün mesajların saxlanılması və bir neçə dəfə işlənməsi əlavə üstünlüyü ilə.

Şəkil 3-3. İki istehlakçı fərqli istehlakçı qruplarında bir partisiyadan oxuyurlar.
İstehlakçılar istehlakçı qrupunda
Bir istehlakçı partisyadan məlumat oxuduqda, o, tamamilə göstəriciyi idarə edir və əvvəlki bölmədə təsvir edildiyi kimi mesajları işləyir.
Əgər bir neçə istehlakçı eyni group_id ilə bir partisyaya abunə olunubsa, sonuncu qoşulan istehlakçıya göstəricinin icarəsi veriləcək və bundan sonra o, bütün mesajları alacaq ().

Şəkil 3-4. Eyni istehlakçı qrupunda iki istehlakçının bir partisyadan oxuması
Bu işləmə rejimi, istehlakçıların sayı partisyaların sayını aşdıqda, monopol istehlakçı növü olaraq qəbul edilə bilər. Bu, istehlakçıların «aktiv-passiv» (ya da «isti-isti») klasterləşməsini təmin etmək üçün faydalıdır, baxmayaraq ki, bir neçə istehlakçının paralel işləməsi («aktiv-aktiv» və ya «isti-isti») çox daha adidir, nəinki gözləmə rejimində olan istehlakçılar.
Yuxarıda təsvir olunan mesaj paylaması davranışı, adi JMS sırası ilə müqayisədə təəccüblü ola bilər. Bu modeldə, sıraya göndərilən mesajlar iki istehlakçı arasında bərabər paylanacaq.
Ən çox rast gəlinən hallarda, bir neçə istehlakçı yaratdığımız zaman, bunu ya mesajların paralel işlənməsi, ya da oxuma sürətinin artırılması üçün, ya da oxuma prosesinin davamlılığını artırmaq üçün edirik. Partisyadan məlumatı eyni anda yalnız bir istehlakçı oxuya biləcəyi üçün, bu Kafka-da necə əldə edilir?
Bunu etmənin yollarından biri, bütün mesajları oxuyub onları bir ip axınına ötürmək üçün bir istehlakçıdan istifadə etməkdir. Bu yanaşma işlənmə sürətini artırsa da, istehlakçıların məntiqinin mürəkkəbliyini artırır və sistemin oxuma davamlılığını artırmır. Əgər bir istehlakçı enerji kəsilməsi və ya oxşar bir hadisə səbəbindən dayandırılırsa, oxuma dayanır.
Bu problemin Kafka-da kanonik həlliBir neçə partisyadan istifadə etməkdir.
Partisyonlama
Partisiyalar oxumağı paralel etmək və bir brokerin bir istehlak gücünü aşmaq üçün əsas mexanizmdir. Bunun daha yaxşı başa düşülməsi üçün, gəlin iki partisiyadan ibarət bir mövzunun olduğunu və bu mövzuya bir istehlakçının abunə olduğunu düşünək ().

Şəkil 3-5. Bir istehlakçı bir neçə partisiyadan oxuyur
Bu ssenaridə istehlakçıya hər iki partisiyada öz group_id-ə uyğun göstəricilərin idarə edilməsi verilir və həmçinin hər iki partisiyadan mesajlar oxumağa başlayır.
Eyni group_id üçün bu mövzuya əlavə bir istehlakçı əlavə edildikdə, Kafka bir partiyadan digərinə birini yenidən ayırır (reallocate). Bundan sonra hər bir istehlakçı nümunəsi bir partiyadan oxuyacaq.).
Mesajları 20 iplik üzrə paralel işlətmək üçün sizə minimum 20 partiya lazım olacaq. Əgər partiyalar daha az olarsa, sizin işləməyə heç nəsi olmayan istehlakçılarınız qalacaq ki, bu da monopol istehlakçılar müzakirəsində qeyd edilib.

Şəkil 3-6. Eyni istehlakçı qrupundakı iki istehlakçı fərqli partiyalardan oxuyur
Bu sxem broker Kafka-nın JMS növbəsini dəstəkləmək üçün lazım olan mesajların paylanması üzrə işini əhəmiyyətli dərəcədə asanlaşdırır. Burada aşağıdakı məsələlərə diqqət yetirmək lazım deyil:
- Növbəti mesajın hansı istehlakçıya çatdırılmalı olduğuna dair, dövrü paylanma (round-robin), cari əvvəlcədən götürmə tamponlarının tutumu və ya əvvəlki mesajlara əsaslanaraq.
- Hansı mesajların hansı istehlakçılara göndərildiyi və uğursuzluq halında yenidən təqdim olunmalı olub-olmadığı.
Broker Kafka-nın etməli olduğu yeganə şey, istehlakçı onları soruşduqda mesajları ardıcıl şəkildə ötürməkdir.
Lakin, paralel oxuma tələbləri və uğursuz mesajların yenidən göndərilməsi heç yerdən getmir — bunların məsuliyyəti brokerdən müştəriyə keçir. Bu, onların kodunuzda nəzərə alınmalı olduğunu bildirir.
Mesaj göndərilməsi
Mesajı hansı partiyaya göndərmək qərarı bu mesajın istehsalçısına aiddir. Bunun üçün tətbiq olunan mexanizmi başa düşmək üçün əvvəlcə həqiqətən göndərdiyimiz şeyə nəzər salmalıyıq.
JMS-də metadata (başlıqlar və xüsusiyyətlər) strukturu olan və yük (payload) daşıyan bir mesaj formasından istifadə edərkən, Kafka-da mesaj «açar-dəyər» cütüdür. Mesajın yüklənməsi dəyər (value) olaraq göndərilir. Açar, digər tərəfdən, əsasən partiyalaşdırma üçün istifadə olunur və biznes məntiqinə xüsusi açardaşımalıdır ki, əlaqəli mesajlar eyni partiyada olsun.
2-ci Fəsildə, əlaqəli hadisələrin bir istehlakçı tərəfindən ardıcıl şəkildə işlənməli olduğu online bahis senarisini müzakirə etdik:
- İstifadəçi hesabı qurulub.
- Pul hesaba köçürülür.
- Hesabdan pul çıxaran bir bahis edilir.
Hər bir hadisə bir mövzuya göndərilən mesaj kimi təqdim olunursa, bu halda təbii açar istifadəçi hesabının identifikatoru olacaq.
Kafka Producer API istifadə edərək mesaj göndəriləndə, bu mesaj partiyalaşdırma funksiyasına daxil edilir. Bu funksiya mesajı və Kafka klasterinin cari vəziyyətini nəzərə alaraq, mesajın göndəriləcəyi partiyanın identifikatorunu qaytarır. Bu funksiya Java-da Partitioner interfeysi vasitəsilə həyata keçirilir.
Bu interfeys aşağıdakı kimidir:
interface Partitioner {
int partition(String topic,
Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster);
}Partitioner-in partiyanı təyin etmək üçün həyata keçirdiyi tətbiq, standart olaraq açar üçün hash üsulunu (general-purpose hashing algorithm over the key) və ya açar göstərilmədikdə dövrü tədqiqatı (round-robin) istifadə edir. Bu standart dəyər əksər hallarda yaxşı işləyir. Lakin gələcəkdə öz strategiyanızı yazmağınız lazım ola bilər.
Öz partiyalaşdırma strategiyanızı yazmaq
Gəlin, mesajın yükü ilə yanaşı metadata göndərmək istədiyimiz bir misala baxaq. Bizim nümunədə yük, oyun hesabında depozit yatırılması üçün göstərişdir. Göstəriş, göndərilirken mütləq dəyişdirilməyəcəyini təmin etmək istədiyimiz bir şeydir və yalnız etibarlı yuxarı sistemin bu göstərişi başlatmasına əmin olmaq istəyirik. Bu halda, göndərən və qəbul edən sistemlər, mesajın doğruluğunu təsdiqləmək üçün imza istifadəsi ilə razılaşırlar.
Adi JMS-də biz sadəcə mesajın 'imza' xüsusiyyətini müəyyən edirik və onu mesaja əlavə edirik. Lakin, Kafka bizə metadata göndərmək üçün mexanizm təmin etmir — yalnız açar və dəyər.
Dəyər, biz saxlamaq istədiyimiz bank transferinin yükü (bank transfer payload) olduğundan, açarda istifadə etməli olduğumuz verilənlər strukturunu müəyyən etməkdən başqa çarəmiz qalmır. Açarlara partiyalaşdırma üçün bizə lazım olan hesab identifikatorunu düşündüyümüzü qəbul etsək, bütün hesabla əlaqəli mesajların ardıcıl işlənməli olduğunu, aşağıdakı JSON strukturu ilə təqdim edəcəyik:
{
"signature": "541661622185851c248b41bf0cea7ad0",
"accountId": "10007865234"
}İmza dəyəri yükə görə dəyişəcəyindən, Partitioner interfeysinin standart hash strategiyası bağlı mesajları etibarlı şəkildə qruplaşdırmayacaq. Buna görə, bu açarı analiz edəcək və accountId dəyərini ayıracaq öz strategiyamızı yazmalıyıq.
Kafka, mesajların zədələnməsini aşkar etmək üçün yoxlama cəmləri daxil edir və tam təhlükəsizlik funksiyaları toplusuna malikdir. Hətta bu halda belə, bəzən yuxarıdakı kimi spesifik sənaye tələbləri meydana gəlir.
İstifadəçi parçalanma strategiyası, bütün əlaqəli mesajların tək bir parçaya düşməsini təmin etməlidir. Bu, sadə görünə bilsə də, əlaqəli mesajların sıralanmasının əhəmiyyəti və mövzuda parçaların sayının nə qədər sabit olduğu bu tələbi çətinləşdirə bilər.
Mövzuda parçaların sayı zamanla dəyişə bilər, çünki ilkin gözləntiləri aşan trafik üçün əlavə edilə bilər. Beləliklə, mesajların açarları, ilkin göndərildiyi parça ilə əlaqələndirilə bilər ki, bu da istehsalçı nümunələri arasında bölüşdürülməsi lazım olan vəziyyətin bir hissəsini ifadə edir.
Digər bir faktor, mesajların parçalar arasında bərabər şəkildə yayılmasını nəzərə almaqdır. Adətən, açarlar mesajlar arasında bərabər paylanmır və hash funksiyaları kiçik açar dəstləri üçün ədalətli paylanmanı təmin etmir.
Qeyd etmək vacibdir ki, mesajları hər necə ayırsanız, ayırıcıdan bəlkə də təkrar istifadə etməli olacaqsınız.
Fərqli coğrafi yerlərdəki Kafka klasterləri arasında məlumatların replikasiyası tələblərini nəzərdən keçirək. Bu məqsəd üçün Kafka, bir klasterdən mesajları oxumaq və başqasına göndərmək üçün istifadə edilən MirrorMaker adlanan komanda xətti aləti ilə birlikdə təqdim olunur.
MirrorMaker, klasterlər arasında replikasiya zamanı mesajlar arasında nisbi sıralamanı saxlamaq üçün replikasiyaya alınan mövzu açarlarını başa düşməlidir, çünki bu mövzunun parçalarının sayı iki klasterdə bərabər olmaya bilər.
Xüsusi parçalanma strategiyaları nisbətən nadirdir, çünki standart hashing və dövrvi paylama əksər vəziyyətlərdə uğurla işləyir. Lakin, sərt sıralama zəmanətinə ehtiyacınız varsa və ya yük məzmunundan metadata çıxarmaq lazımdırsa, o zaman parçalanmaya daha yaxından nəzər yetirməlisiniz.
Kafka-nın miqyaslandırma və performans üstünlükləri, bəzi ənənəvi brokerin öhdəliklərini müştəriyə köçürməklə təmin olunur. Bu halda, paralel işləyən bir neçə istehlakçı arasında potensial olaraq bağlı mesajların paylanmasına qərar verilir.
JMS brokerləri eyni tələblərlə də məşğul olmalıdır. Maraqlıdır ki, JMS Mesaj Qrupları vasitəsilə həyata keçirilən əlaqəli mesajların eyni istehlakçıya göndərilməsi mexanizmi (sticky load balancing (SLB) strategiyasının bir növü) göndərənin mesajları əlaqəli olaraq işarələməsini tələb edir. JMS-də broker, bir çox istehlakçı arasında əlaqəli mesajlar qrupunun bir istehlakçıya göndərilməsini və istehlakçı sistemdən ayrıldıqda qrupun mülkiyyətinin ötürülməsini təmin edir.
İstehsalçı ilə bağlı razılaşmalar
Partisiyalaşdırma, mesaj göndərərkən nəzərə alınması gərəkən yeganə amil deyil. Gəlin Java API-də Producer sinfinin send() metodlarına bir nəzər salaq:
Future send(ProducerRecord record);
Future send(ProducerRecord record, Callback callback);Həmçinin, qeyd etmək lazımdır ki, hər iki metod Future qaytarır ki, bu da deməkdir ki, göndərmə əməliyyatı dərhal icra edilmir. Nəticədə, mesaj (ProducerRecord) hər aktiv partition üçün göndərmə tamponuna yazılır və Kafka müştəri kitabxanasında fon prosesi vasitəsilə brokerə ötürülür. Bu, işi inanılmaz dərəcədə sürətli edir, ancaq bu o deməkdir ki, təcrübəsiz şəkildə yazılmış bir tətbiq, prosesi dayandırıldığı halda mesajları itirə bilər.
Həmçinin, həmişə performans xərcləri ilə daha etibarlı göndərmə əməliyyatı həyata keçirmək yolları mövcuddur. Bu tamponun ölçüsünü 0 olaraq təyin etmək olar və göndərmə tətbiqinin prosesi brokerə mesajın göndərilməsini tamamlamasını gözləməyə məcbur olacaq:
RecordMetadata metadata = producer.send(record).get();Mesajları oxumaq haqqında bir daha
Mesaj oxumağın əlavə çətinlikləri vardır ki, bunları müzakirə etmək lazımdır. JMS API-dən fərqli olaraq, mesajlardan gələn zaman mesaj dinləyicisini (message listener) işə salmağı mümkün edən, Consumer Kafka yalnız sorğulama (polling) edir. Bu məqsəd üçün istifadə olunan poll()metodu ilə daha ətraflı tanış olaq:
ConsumerRecords poll(long timeout);Metodun qaytarılan dəyəri bir neçə obyektin ehtiva edildiyi konteyner qurğusudur. ConsumerRecord potensial bir neçə partition-dan alınan açar-dəyər cütü ilə müvafiq metadata saxlayan bir obyekt-holddur. ConsumerRecord Yuxarıda müzakirə edildiyi kimi, müvafiq olaraq uğurlu və ya uğursuz işlənildikdən sonra mesajlarla nə baş verdiyini bizim davamlı olaraq xatırlamamız lazımdır, məsələn, müştəri mesajı işlədə bilmirsə və ya o, dayandırılırsa. JMS-də bu təsdiqləmə rejimi (acknowledgement mode) vasitəsilə həll edilir. Broker ya uğurla işlənmiş mesajı siləcək, ya da işlənməmiş və ya uğursuz olanı (şərt olaraq, tranzaksiyalar istifadə olunduqda) yenidən göndərəcək.
Yuxarıda müzakirə edildiyi kimi, biz uğurlu və ya uğursuz işlənildikdən sonra mesajlarla nə baş verdiyini yadda saxlamağa davam etməliyik, məsələn, müştəri mesajı işlədə bilmir və ya prosesi dayandırırsa. JMS-də bu, təsdiqləmə rejimi (acknowledgement mode) vasitəsi ilə həll olunur. Broker ya uğurla işlənmiş mesajı siləcək, ya da yenidən işlənməmiş və ya uğursuz olan mesajı göndərəcək (şərtlə ki, tranzaksiyalar istifadə edilsin).
Kafka tamamilə fərqli çalışır. Mesajlar brokerdə oxunduqdan sonra silinmir və uğursuzluq baş verdikdə məsuliyyət oxuma kodunun özündədir.
Daha əvvəl qeyd etdiyimiz kimi, istehlakçı qrupları jurnalın sıralanması ilə bağlıdır. Bu sıralama ilə əlaqəli jurnal yeri, cavab olaraq verilecek növbəti mesajı göstərir. poll(). Oxuma zamanı sıralamanın artdığı an mühüm əhəmiyyət kəsb edir.
Daha əvvəl müzakirə olunan oxuma modelinə qayıdaraq, mesajın emalı üç mərhələdən ibarətdir:
- Mesajı oxumaq üçün çıxarmaq.
- Mesajı emal etmək.
- Mesajı təsdiqləmək.
Kafka istehlakçısı konfiqurasiya seçimi ilə gəlir enable.auto.commit. Bu, adətən "avto" sözünü ehtiva edən konfiqurasiyalarda olduğu kimi, istifadə olunan bir prefabrik ayarıdır.
Kafka 0.10-dan əvvəl, bu parametrlə işləyən müştəri, növbəti çağırış zamanı son oxunan mesajın sıralamasını göndərirdi poll() emaldan sonra. Bu, artıq çıxarılmış (fetched) olan mesajların, əgər müştəri onları artıq emal edibsə, lakin gözlənilmədən çağırışdan əvvəl məhv olunsa, yenidən emal oluna biləcəyini ifadə edirdi. poll(). Broker, mesajın neçə dəfə oxunduğu ilə bağlı heç bir vəziyyəti saxlamadığı üçün, bu mesajı çıxaran növbəti istehlakçı, pis bir şeyin baş verdiyini bilmir. Bu davranış təkrar-transaksion idi. Sıralama yalnız mesajın uğurla emal edildiyi halda təsdiqlənirdi, lakin müştəri işləmədikdə, broker eyni mesajı başqa bir müştəriyə yenidən göndərirdi. Bu davranış "bir dəfə olsun«.
ən azı bir dəfə Kafka 0.10-da müştəri kodu, sıralamanın "auto.commit.interval.ms" konfiqurasiyasına uyğun olaraq müştəri kitabxanası tərəfindən dövri olaraq başlandığı şəkildə dəyişdirildi.Bu davranış JMS AUTO_ACKNOWLEDGE və DUPS_OK_ACKNOWLEDGE rejimləri arasında bir yerdədir. Avtokomit istifadə edildikdə, mesajlar faktiki olaraq emal edilib-edilmədiyindən asılı olmayaraq təsdiqlənə bilər — bu yavaş istehlakçı zamanı baş verə bilər. İstehlakçı işləmədikdə, növbəti istehlakçı çıxarılmış mesajları, təsdiqlənmiş vəziyyətdən başlayaraq alırdı ki, bu da mesajın keçməsinə səbəb ola bilər. Bu müddətdə Kafka mesajları itirmir, oxu kodu sadəcə onları emal etmirdi.
Bu rejimin perspektivləri 0.9 versiyasındakı kimi eynidir: mesajlar emal oluna bilər, lakin uğursuzluq baş verərsə, sıralama təsdiqlənməyə bilər ki, bu da potensial olaraq çarpaz çatdırılma ilə nəticələnə bilər. Daha çox mesaj çıxardıqca poll(), bu problem daha da artır.
Bölüm "Sıra Halka Mesaj Okuma"da tartışıldığı gibi, mesajlaşma sisteminde tek seferlik mesaj teslimatı gibi bir kavram yoktur, arıza modları dikkate alındığında.
Kafka'da bir ofseti (offset) kaydetmenin (commit) iki yolu vardır: otomatik ve manuel. Her iki durumda da mesajlar, mesaj işlendi ancak committen önce bir arıza oluşursa birden fazla kez işlenebilir. Ayrıca, commit arka planda gerçekleştiyse ve kodunuz işlemeye başlamadan önce tamamlandıysa, bir mesajı hiç işlemeyebilirsiniz (bu muhtemelen Kafka 0.9 ve önceki sürümlerde geçerlidir).
Offset commit sürecini manuel olarak yönetmek için, Kafka tüketici API'sinde enable.auto.commit parametresini false olarak ayarlayıp, aşağıdaki yöntemlerden birini açıkça çağırmanız gerekir:
void commitSync();
void commitAsync();Eğer bir mesajı "en az bir kez" işlemeyi hedefliyorsanız, commitSync()komutunu mesajları işlemeyi tamamladıktan hemen sonra çağırmalısınız.
Bu yöntemler mesajlar işlenmeden önce onay (acknowledged) edilemez, fakat potansiyel çift işlemin ortadan kaldırılması için hiçbir şey yapmazken, işlem görünümünü yaratırlar. Kafka'da işlemler yoktur. Müşterinin bunu yapma şansı yoktur:
- Başarısız olan bir mesajı otomatik olarak geri almak (roll back). Tüketiciler, sorunlu payload'lar ve backend kesintileri nedeniyle oluşan istisnaları kendileri ele almalıdır, çünkü broker tarafından mesaj yeniden gönderimine güvenemezler.
- Bir atomik işlem içinde birden fazla konuya mesaj göndermek. Kısa sürede göreceğimiz gibi, çeşitli konular ve parçalar üzerinde kontrol farklı makinelerde olabilir ve bu işlemlerde işlemleri koordine etmezler. Bu makalenin yazıldığı sırada, bunu KIP-98 ile mümkün kılmak için bazı çalışmalar yapılmıştır.
- Bir konudan bir mesaj okumayı başka bir konuya başka bir mesaj gönderimi ile bağlamak. Gene, Kafka mimarisi bağımsız birçok makinenin tek bir hat gibi çalışmasına dayanır ve bunu gizlemek için hiçbir girişim yapılmamaktadır. Örneğin, transaksiyonlarda Tüketici və Üretici bağlayacak API bileşenleri yoktur. JMS’de bu, Sessioniçinden oluşturulan MessageProducers və MessageConsumers.
olarak sağlanır. Eğer işlemlere güvenemiyorsak, geleneksel mesajlaşma sistemleri tarafından sağlanan semantiğe daha yakın bir anlamı nasıl sağlayabiliriz?
Əgər istehlakçı yerinin mesajın emal edilməsindən əvvəl artması ehtimalı varsa, məsələn, istehlakçının uğursuzluğu zamanı, istehlakçının qrupunun mesajları itirdiyini bilmək imkanı yoxdur, ona görə də partiyaya təyin edildikdə. Bununla belə, strategiyalardan biri yerin əvvəlki vəziyyətə geri döndürülməsidir. Kafka istehlakçısının API-si bunun üçün aşağıdakı metodları təqdim edir:
void seek(TopicPartition partition, long offset);
void seekToBeginning(Collection partitions); Metod seek () metodu ilə istifadə edilə bilər
offsetsForTimes (Map timestampsToSearch) geçmişdəki müəyyən bir anda geri dönmək üçün.
Dolayısı ilə, bu yanaşmanın istifadəsi demək olar ki, artıq emal edilmiş bəzi mesajların yenidən oxunub emal ediləcəyi anlamına gəlir. Bunun qarşısını almaq üçün, 4-cü fəsildə təsvir edildiyi kimi, hərtərəfli oxuma metoda istifadə edərək, əvvəllər görülmüş mesajları izləyə və təkrarlananları istisna edə bilərik.
Alternativ olaraq, istehlakçı kodunuz sadə ola bilər, əgər mesajların itirilməsi və ya təkrarlanması icazə verilirsə. Kafka-yla adətən istifadə olunan istifadə senarilərini, məsələn, log hadisələrinin emalında, metriklərin, klik izlənməsində və s. nəzərdən keçirdikdə, narahat bir tətbiqinin mühiti üzərində ayrı mesajların itirilməsi ehtimalı çox azdır. Belə hallar üçün standart dəyərlər tamamilə qənaətbəxşdir. Digər tərəfdən, əgər tətbiqiniz ödəmələri ötürmək lazımdırsa, hər bir mesajı diqqətlə izləməlisiniz. Hamısı kontekstə bağlıdır.
Şəxsi müşahidələr göstərir ki, mesajların intensivliyinin artması ilə, hər bir ayrı mesajın dəyəri azalır. Böyük Həcmli mesajlar, adətən, toplanmış formada dəyərli hesab edilir.
Yüksək mövcudluq (High Availability)
Kafka'nın yüksək mövcudluq yanaşması ActiveMQ-nun yanaşmasından əhəmiyyətli dərəcədə fərqlənir. Kafka, bütün broker nümunələrinin eyni anda mesajları qəbul etdiyi və payladığı üfüqi miqyaslana bilən klasterlərdə inkişaf etdirilmişdir.
Kafka klasteri, müxtəlif serverlərdə işləyən bir neçə broker nümunəsindən ibarətdir. Kafka, hər düyünün özünə ayrılmış xatirəsi olan adi müstəqil avadanlıqlarda işləmək üçün yaradılmışdır. Şəbəkə xatirələrindən (SAN) istifadə etməyi tövsiyə edilmir, çünki bir çox hesablama düyünləri zaman intervallarını paylaşa bilər və münaqişələr yaradır.İnterval
Kafka - daima aktivdir. sistem. Bir çox iri Kafka istifadəçiləri heç vaxt klasterlərini dayandırmır və proqram təminatı həmişə ardıcıl yeniləmə təmin edir. Bu, mesajlar və brokerlər arasındakı qarşılıqlı əlaqələr üçün əvvəlki versiya ilə uyğunluğun təmin edilməsi sayəsində mümkün olur.
Brokerlər server klasterinə qoşulmuşdur , bu, konfiqurasiya məlumatlarının qeydiyyatı kimi çıxış edir və hər bir brokerin rollarının koordinasiyası üçün istifadə olunur. ZooKeeper özlüyündə yüksək əlçatanlığı təmin edən, məlumatların təkrarlanması yolu ilə fəaliyyəti olan paylanmış bir sistemdir. kvorum.
Əsas halda, Kafka klasterində topik aşağıdakı xüsusiyyətlərlə yaradılır:
- Partisyon sayı. Yuxarıda müzakirə edildiyi kimi, burada istifadə olunan dəqiq dəyər, arzu olunan paralel oxuma səviyyəsindən asılıdır.
- Təkrarlama faktoru, klasterdə bu partisiya üçün neçə broker nümunəsinin jurnalını saxlamalı olduğunu müəyyən edir.
ZooKeeper-lərin koordinasiyasını istifadə edərək, Kafka yeni partiyaların brokerlər arasında ədalətli paylanmasını təmin etməyə çalışır. Bu, rolunu İdarəçi kimi icra edən bir nümunə tərəfindən edilir.
İcraat zamanı hər bir topik partiyası üçün İdarəçi brokerə rollar təyin edir lider (lider, master, başçısı) və izin alanlar (followers, slaves, subyektlər). Müvafiq partisiya üçün lider olan broker, ona göndərilən bütün mesajların qəbul edilməsindən və konsumentlərə məlumatların yayılmasına cavabdehdir. Topik partiyasına mesaj göndərilərkən, bu, həmin partiya üçün lider olan bütün broker düyünlərinə təkrarlanır. Hər bir partiyanın jurnalını saxlayan düyünə təkrarlamadeyilir. Broker bəzi partiyalar üçün lider, bəziləri üçün isə izləyici kimi çıxış edə bilər.
Liderdə saxlanılan bütün mesajları ehtiva edən izləyici sinkronizasiya olunmuş təkrarlama (sinkronizasiya vəziyyətində olan təkrarlama, in-sync replica). Əgər bir partisiya üçün lider olan broker bağlanırsa, bu partiyası üçün aktual və ya sinkronizasiya olunmuş vəziyyətdə olan istənilən broker lider rolunu üzərinə götürə bilər. Bu, çox davamlı bir dizayndır.
Prodüserin konfiqurasiya hissəsi acks, mesajın qəbul edilməsini təsdiqləməsi lazım olan neçənin təkrarlanmalarının (acknowledge) olmasıdır, tətbiq axını 0, 1 və ya hamısı arasında davam edə bilər. Əgər hamısıdəyəri təyin edilibsə, mesajın lideri aldığı zaman cavab qaytarır (confirmation) prodüserə, o, özü də daxil olmaqla, partiyanın təyin edilmiş konfiqurasiyasında bir neçə təkrarlama(acknowledgements) tərəfindən təsdiqlənmə aldıqda, bununla birlikdə tövsiyyə olunan setting min.insync.replicas (şablon 1). Əgər mesaj uğurla çoxaldıla bilməzsə, istehsalçı tətbiq üçün istisna yaradacaq (YetərliNüsxələrYoxdur və MesajDaxilEdiləndənSonraYetərliNüsxələrYoxdur).
Tipik konfiqurasiyada 3 nüsxə ilə bir mövzu yaradılır (1 lider, hər bir partisiya üçün 2 izləyici) və parametr min.insync.replicas 2 dəyərinə təyin edilir. Bu halda, klaster, topik partisiya idarə edən brokerlərdən birinin söndürülməsinə icazə verəcək, müştəri tətbiqlərinə təsir etmədən.
Bu, artıq tanış olduğumuz performans və etibarlılıq arasında əməkdaşlığa dönüş edir. Nüsxələşmə, izləyicilərdən təsdiq (acknowledgments) gözləmə zamanı əlavə vaxt sərf edir. Ancaq, paralel icra edildiyi üçün, ən azı üç düyündə nüsxələşmə, iki düyündə olduğu kimi eyni performansa malikdir (şəbəkə bant genişliyinin artımını gözardı etsək).
Bu nüsxələşmə sxemindən istifadə edərək, Kafka, mesajların fiziki yazılmasını " sync ()" əməliyyatı vasitəsilə təmin etmə tələbindən ustalıqla qaçır. İstehsalçı tərəfindən göndərilən hər bir mesaj partisiya jurnalına yazılacaq, lakin 2-ci Fəsildə müzakirə edildiyi kimi, ilkin olaraq fayla yazma əməliyyatı əməliyyat sisteminin kuzeyində icra olunur. Əgər bu mesaj başqa bir Kafka instansiyasına nüsxələnibsə və onun yaddaşındadırsa, liderin itirilməsi, mesajın itirildiyi demək deyil - bunu sinxronizə olunmuş nüsxə üstlənə bilər.
Emal etmə tələbini yerinə yetirməmək sync () Kafka-nın mesajları, onları yaddaşa yaza biləcək sürətlə qəbul edə biləcəyini ifadə edir. Və tərsinə, yaddaşdan diska boşaltma (flushing) əməliyyatını mümkün olduğu qədər uzun müddət gecikdirmək daha yaxşıdır. Bu səbəbdən, Kafka brokerlərinə tez-tez 64 Gb və daha çox yaddaş ayrılması heç də nadir deyil. Bu cür yaddaş istifadəsi, bir Kafka instansiyasının ənənəvi mesaj brokerindən çox dəfələrlə daha sürətli işləyə biləcəyi deməkdir.
Kafka da " sync () " mesaj paketlərinə əməliyyat tətbiq olunması üçün konfiqurasiya oluna bilər. Kafka-da hər şey paketlərlə işləməyə yönəlmiş olduğu üçün, bu, əslində, bir çox istifadə senariləri üçün kifayət qədər yaxşı işləyir və çox güclü zəmanət tələb edən istifadəçilər üçün faydalı bir alətdir. Kafka-nın xalis performansının çoxu mesajların brokerə paketlər şəklində göndərilməsi və bu mesajların brokerdən ardıcıl bloklarla oxunması ilə bağlıdır. əməliyyatlar (əməliyyatlar, məlumatların bir yaddaş bölgəsindən digərinə köçürülməsi vəzifəsinin yerinə yetirilmədiyi hallarda). Sonuncu, performans və resurslar baxımından böyük bir qazancdır və yalnız arxasında yatan jurnal verilənlər strukturu sayəsində mümkün olur ki, bu da partion sxemasını müəyyən edir.
Kafka klasterində, mövzunun partionları ayrı-ayrı maşınlarda üfüqi şəkildə genişləndirilə bildiyi üçün bir Kafka brokerindən istifadə etməyə nisbətən çox daha yüksək performans mümkündür.
Yekunlar
Bu fəsildə biz Kafka arxitekturasının müştərilər və brokerlər arasında münasibətləri necə yenidən düşünməsini müzakirə etdik ki, bu da ənənəvi mesaj brokerlərindən qat-qat daha yüksək bant genişliyi olan inanılmaz davamlı bir mesajlaşma boru kəməri təmin edir. Bu məqsədə nail olmaq üçün istifadə etdiyi funksionallığı müzakirə etdik və bu funksionallığı təmin edən tətbiqin arxitekturasını qısa şəkildə nəzərdən keçirdik. Növbəti fəsildə biz mesajlaşma əsasında çalışan tətbiqlərin qarşılaşmalı olduğu ümumi problemləri müzakirə edəcəyik və onların həllinə dair strategiyaları müzakirə edəcəyik. Fəsili, mesajlaşma texnologiyaları haqqında düşünməyi necə əhatə edəcəyimizi təsvir edərək bitirəcəyik ki, siz onların istifadəyə yararlılığını qiymətləndirə biləsiniz.
Əvvəlki tərcümə olunmuş hissə:
Tərcümə edildi:
Davamı gəlir…
Yalnız qeydiyyatdan keçmiş istifadəçilər sorğuda iştirak edə bilərlər. , xahiş edirəm.
Sizin təşkilatda Kafka istifadə olunur?
Bəli
Xeyr
Əvvəl istifadə olunurdu, indi istifadə edilmir
İstifadə etməyi planlaşdırırıq
38 istifadəçi səs verdi. 8 istifadəçi hətta qaldı.
Mənbə: habr.com
