← Back to list

KAFKA A-Z

Xususi ile “monolith to microservices”, “distributed systems” trendlerinden sonra daha da populyarlasan Kafka haqqinda gelin danisaq.

Rasul Rzayev · 2021-06-30 14:17 · 20 claps · 9.8 min read
#kafka #kafka-streams #kafka-connect
Open on Medium ↗

KAFKA A-Z

Xususi ile “monolith to microservices”, “distributed systems” trendlerinden sonra daha da populyarlasan Kafka haqqinda gelin danisaq.

Ilk defe Linkedin hem yaratdi ve istifade etdi daha sonrasinda ise diger neheng sirketler Netflix, Uber, Slack, Spotify ve.s terefinden de Kafka mutemadi istifade edilmeye baslanildi ve bu populyarligina elave bir populyarliq geldi.

Evvelce gelin Kafka-nin isleme mentiqine goz ataq. Umumi arxitekturasina ve komponentlerine baxaq.

Daha umumi anlayisimiz Cluster-dir. Bu termin artiq zaten coxumuza bellidir. Cluster icerisinde de brokerler vardir. Datalarimiz bu brokerlerde saxlanilir (daha deqiq partitionlarda) ve onlar hard disk uzerinde saxlanilir in memory deyil. Bele bir yanlis micconception var ki, memory hemise diskden daha suretlidir lakin bu access pattern-den aslidir. Kafka sequential yeni ki, ardicil olaraq yazma oxuma ucun disk bloklarina muraciet etdiyi ucun , random etmediyi ucun bu cox suretle bas verir. Ve data transfer ucun de zero copy istifade etdiyi ucun daha suertle calisir. Bu zero-copy meselesini etrafli izah edecem bu mqalede irelide.

Kafkanin bu brokerlerinin idare edilmesi, fault broker healt check edilmesi, leader partition switchler bu kimi isler ise Zookeper terefinden edilir deye siz Zookeper install edirsiniz kafka install ederken. Brokerlerin sayi hemise tek eded olmalidir duzgun replication getmesi ucun buna diqqet edin.

*** Kafka 2.8.0-dan etibaren zookerper required deyil

Gelin indi brokerimizin icerisine baxaq: Kafka datalari topiclerde saxlayir dedik. Yeni toplic-lerimiz brokerlerin icerisindedir. Bir nov database table-lerimize benzede bilersiniz topicleri. Topic-in de icerisinde partitionlarimiz vardir. Yeni datamiz, eventlerimiz eslinde partitionlarin icierisindedir.

Bir topicde istediyiniz sayda partition ola biler. Partition sayini topicleri yaradarken ozumuz teyin edirik .Datalarimiz elave olunmas sira ile partitionlara dusur(single partition-dirsa). Ve bu datalar ya muyyen muddetden( default 7 gun) ve ya mueyyen size kecmediyi muddetce tutulur bu parititonlarda.

Her bir event offset ile bir yerde partitionda tutulur. Ve bu offset deyeri uniqie-dir. Bu partitionlar vasitesi ile sirali sekilde datalarimizi yaza ve store ede bilerik, event sourcing ede bilerik qisaca. Hemcinin data read ederken paritionlardan paralel oxumani temin ede bilerik daha da suret qazanmaq ucun.

Normalda topice publish olunan data partitionlara default olaraq raund ribbon paylanir random olaraq amma eger key istifade edilibse default staretegy key hash based range uzerinden paylanilir.

Amma biz record keyler vasitesile hansi fielde gore partitionlari ayri ayri aggregation olmasini temin ede bilerik. Eger record key olaraq sequence id-ni gonderirsiznizse demeli sizde sirali sekilde partionlara dusecek datalar ve event sorucing temin ede bilersiniz bu yolla. Bu arada partitionlarin da replication-ini ede bielrsiniz ve bu zaman mueyyen leader partition olur qalani replica partitinlar olur ve bu ayri ayri brokerlerde de tutula bilir.

Kafka yalniz 1 partition-a yazmada global ordering(siralama) qaranti edir. Multiple partitionda default olaraq ordering-i qaranti etmir!

Amma key istifade ederek siz event-leri key base partitionlara yonlendirmis olursunuz. Kafka oz daxilinde tekrar etme- retry mentiqi var internally:

Bele bir config-e baxaq: burada 30(sec) erzinde ACK ala bilmese broker bu zaman retry edecek 1 defe ve bu acknowledge oyrenmek ucun 3 paralel sorgu atacaq.

“max.in.flight.requests.per.connection” deyeri ise nece dene paralel request unacknowledge \ acknowledge almasi ucun gonderilecek onu gosterir. Bu deyer 1 olmasi sizin ordered meselesinde destek olacaq. Amma 1 set olmasi latency yarada biler. Cunki default value 5-dir buna sebeb ise 5 paralel sorgu atildiqda hansindan result elde edilse diger netice gozlenilmeden dger eventin send edilmesine kecile biler artiq.

En esasi: bu retry sizin eventlerin orderiniz pozulmasina getire biler. Cunki eslinde 2-ci batch eventler once getmesine baxmayaraq retry zamani success olubsa batch-1 adlandirsaq fail olan eventleri bu eventler eslinde batch2-deki eventlerden evelde olmali idi indi ise daha sonra gelecek:

Gelek acknowledgment meselesine:

Burada bir nece cur yanasa bilerik. Kafkada bir nece delivery semanticler vardir ( at least once, at most once ve exactly once) ve bunlari ise ack deyerleri ve insync.replica deyerleri ile tenizmleye bilerik.

  • Eger producer ack = 0 veririkse bu o demekdirki mesaji gonder ve getdi getmedi cavab gozlemeden digerini gonder. Bu suretli amma riskli variantdir!
  • Deyerini 1 versek bu defe ancaq leader partition yazana qeder saxlayacaq diger mesaji, sonrasinda iseb davam edecek .
  • All deyerini versek bu defe mutleq sekilde butun partiitonlara dusmesini gozleyecek ve sonra diger mesaja kececek bu cox yavas usuldur! Ve elave olaraq “min.insync.replicas.” deyeri set olunmalidir.

Delivery semnaticleri de izahin verek ayriliqda nedir :

at most once — En coxu bir defe catdirila biler. Data itkisi goze alinandir bu case-de ack=0 olur. Offsete commit qebul edilen kimi bas verir proses sonunda yox. Ve Retry mexanizmi qurmuruq.

at least once — Burda bir event bir nece defe catdirila biler en azi 1 defe islenecek. Duplication ehtimnali var amma data itkisi qebul edilmir, ack =all . Process-i icra edib sonda commit edirsiniz. Xeta olarsa retry mentiqi qurursunuz.

exactly once — Event ancaq max bir defe catdirila biler ve data da itkisi de ola bilmez. Komplekslilik baximindan en komplekisidir. ack=all ve retry qurulur amma check , lock mentiqi qurulur ve ofset commit checkden sonra sonda edilir.

Kafkaya dusen data immutable-dir kenardan deyisdirile bilmez! Bunu ancaq kafka ozu ede biler.

Produce zamani xeta bas vererse ve broker acknowledge vermezse retry-dan sonra eyni event tekrar publish ola biler. Bunun qarsisin almaq ucun ise idempotency enable olmalidir. eventlere elave unique ID tikilir ve broker eyni event geldikde ignore edir artiq.

Kafka partition sayinda default configlerle getmek de eslinde duz olmazdi.Cunki by default single partition ve 1 replication factor ile olacaq. Partition sayi ve replication-i ozunuze uygun deyismeyiniz tovsiyye edilir.

Bes hemise cox partition vermek dogru addimdirmi ? Xeyir, cunki her partition ucun ayri buffer en azindan elave memory usage demekdir. Ve her bir partition brokerde file system-de directory-e map olur. Hemin directory-de log file-da index ve data tutulur. Emeliyyat sisteminizde max open file limit olur bezen hetta bu limiti kece biler. Bes en duzgun say nedir ? Bu use case-e gore deyisir ve hemcinin performans analizi ederek size en uygun say ve configleri vermelisiniz.

Gelek consumerlerimize, bu daha cox rabbit ve kafka bir biri ile qarisdiran ve ferqlerini ile bagli beyninde qarsiqiligi, suallari olanlar ucun de daha onemli hissedir. Kafkada da consumer publish olunan mesaji goturur eyni qaydada. Offset olaraq zookeperden harda qaldigi bilir ve ordan ve ya ilk defe ise dusubse basdan baslayaraq oxuyur partionlardaki datani.

Burada rabbitden de esas ferqli olaraq offset commit, offset reset anlayislari var. Rabbitde biz eventi consume edirdik ve artiq o queue-den itirdi, multiple consmer consume etse bele lock olurdu data digeri ucun ve process zamani event isleyerken xeta alsaq bele tutaqki bazaya yazanda rabbitde elave DL ve PL queue-ler yaradib mecbur fail olduqda ora yigirdiq eventlerimizi. Amma kafkada bir nece yanasma var ve biz bir nece defe commit etmesek eventi alib process ede bilirik.

Hemcinin en gozel xususiyyeti menim fikrimce consumer group-lar olmasidir ki bu bize one message to multiple consumer imkani verir. Ve kafkada consumer groupaki consumer ve onlarin oxusudqlari topic saylarimizi kafka ozu bolusdurur, Meselem 2 topic 2 consumer varsa heresi birin. 2 consumer 4 topic arsa heresi 2-sini. Bunu ozu bolur kafka elave setting ehtiyac yoxdur.

Artiq elave consumer varsa o ise bosda qalir. Kafka race condition meselelerine gore bir partition-a birden cox consumer baglamamiza izin vermir.

Consumerler by default range stategy uzre consume edir.

Spring kafka ile consumerlerimiz ucun bir nece configler vardir bunlar kafka oz documentationinda qeyd edilib. Bunlardan bezilerine baxaq :

concurrency: Bir consumer daxilinde olan consumer-ler by default single (1) thread isleyirler. Amma configuration ile concurrency deyerini artira bilersiz. Bu deyer her bir consumere aid olacaq consumer group daxilindeki. Ve onlarin hamisi active statusda olacaq deye bir sey de yoxdur. idle-da qala biler bir coxu ve imkan yarandiqda (meselem rebalance olduqda ve.s) active statusa kece biler.

session.timeout.ms : CG daixlindeki comsumerin broklerle elaqenin kesilmesinin maximal vaxtidir.Bu muddet kecildikde reblalans bas verir aktiv CG-e

enable.auto.commit : Consumerin offseti auto commit edilib edilmemesi ucun istifde edilir. Default true-dur amma false verilse commit edilme intervali control edilmelidir ve duplicate consumingden qacilmalidir. True secildiyi halda her 5 saniyeden bir (default commit interval) pool()-dan gelen last offset id-ni commit edir. commit.interval.ms ile bu intervali deyise bilersiniz ancaq qisa qoysaz network traffic yukleneceyini nezere alin!

auto.offser.reset : Bele bir ssenari dusunek. 1 topicim var ve icerisinde 2 partition ve ilk partitionda 100 message. 2 dene de consumer groupumuz var ve heresinde 2 consumer var. Ve her ikisinin consumeri 1-ci topici consume edir. Bu zaman duplication bas vermesin deye CG1 ve CG2 consumerlerler offsetlerle islemeldiir. Buna gore kafkada offset reset anlayisi vardir. Burada optionlar earliest, latest, datetime, shift by ve none ola biler ve default latest-dir. Earliest secdikde en evvel basdaki offsetden read edecek, latestde en sondan oxuyur. None zamani ise previous offset tapmasa CG ucun xeta verir.

partition.assignment.strategy: Kafkada consumerlere konkret spesifik uneven/even falan spesifik partitionlardan oxunmani temin ede bilerik. Default ise round Robbin-dir yeni partitionlara muraciet paylanir, ancaq spesifik paritionlardan oxunma isteseydik bunun ucun range deyeri veririk.

Producerler ucun de mueyyen maraqli configler var :

batch.size : eventler gonderilmezden evvel max batch seklinde ne qeder byte yigila biler onu gösteririz

linger.ms : Bu producerin eventi kafkaya atmaq ucun gozleme muddetidir. Default eyeri 0-dir. Bele bir example deyek avtobus daynacagi dusunun burda 0 deyeri versek avtobus adam goren kimi dayanacaga gelib onu aparacaq eve ( topic), amma 5 versem meselem 5 deqiqeden bir gelecek ve hetta batch.size = 10 versem her defe max 10 nefer alacaq avtobusa ))

linger.ms = 0

linger.ms = 0

linger.ms = 5 ve batch.size=10

linger.ms = 5 ve batch.size=10

compression.type: Kafkada high throughput elde etmek ucun compression configle de oynaya bilersiniz, hansiki size datani gondermezden evvel compress edib kafka servere atmaga imkan verir. Ve consumer teref ucun her hansi elace decompress-e ehtiyac olmur:

compression type-larin xaraktersitiklari

compression type-larin xaraktersitiklari

Kafka Schema Registry

Niye buna ehtiyac yarandi evvelce o suala cavab tapaq. Cunki Kafka datani inputu byte seklinde alir qebul edir ve publish edir. Poctalyon kimi mektubu aparir ama icinde ne var bilmir)) Her hansi verification yoxdur. Eger field name deyiserse , bad data ginderilerse ve.s consumerlerimiz deyisikliklerde break ola biler!! Bu riskdir! Bes nece hell edecik ?

Kafkada biz verify etsek her defe duz olmazdi, Kafkan-i suretli eden seyler ne idi (zero copy) idi, yeni ki lazimsiz diskden data copy evezine byte array seklinde datani qebul etmesi ve hemin datani parse ve ya read etmemesi (no CPU usage)ve page cache-e yazmasi direct ve oradan network controllere sendFile() ile oturulmesi ve NIC -den de uygun destination-a gonderilmesi flowu.

Zero copy logic

Zero copy logic

Biz data publish edende kafka onu byte seklinde 0 ve 1-lerden ibaret gorur json, string ve ya ne oldugu icinde ona maraqli deyil :)

O zaman publisher ve consumer bir biri ile danisarken schema registry lazim geldi. Hansiki bu bad datalari reject etme imkani olmalidir onda. lightweight olmalidir, schema destekli olmalidir ve.s

Burada ferq odur ki arada schema regstry olur ve data orada evvelce register edilir sonra consume edilende retrieve edilir.

Avro haqqinda danisaq. Avro shcema define zamani qeyd edilir json daxilinde. Compress istifade edir CPU az istifade etsin deye , yeni field name uzun olsa bele compress edir zaten. Documentation schema icerisine embedded olur. Hadoop based texnologiyalara, meselem Hive kimi yaxsi destek verir. Avro-da mueyyen key-ler field namesler var: name, doc, namespace, fields, aliases, type, default ve.s kimi

Avro-da hemcinin complex type-lar: enum, map, array, unions ve.s da teyin ede bilirik.

Schema Evolution

Eger bugun menim schemamda name ve surname isteyirem amma sabah phone number elave etsem ne bas verecek? Business-im qirilacaqmi ?

Schema evolution type-lar : Backward: yenide yeni field-e default deyer verilerek hell edilir Forward: Avro yeni field-i ignore edir Full: hem back hem forward

Calisin gelecekde silinme ehtimali olan field-lere default valualer set edin. Field raname etmek evezine aileslara ustunluk verin. Required filedlerle ehtiyatli olun silmemeye calisin.

Bu arada UI uzerinden de rahatliqla schema add ede, compabilty type sece bilersiniz.

Kafka registered schema-lari _schemas adli topicde tutur. Ve kafka elave local cache istifade edir bu schema-lar ucun. Producer her defe scema gondermir onun evezine schemaId gonderilir.

Schema registrasyasinin iki usulu vardir: 1. by default produce zamani register edilir . Amma bu productionda aciq qalmasi duz deyil sondurmek ucun asagidaki deyeri false etmek lazimdir: auto.register.schemas=false

  1. Daha secure way ile kenardan register etmek de olur plugin ve ya api yollari ile.

Confluent-in kafka avro serializer ve deseiralize ucun hazir classlari movcuddur. Register schema enable olubsa io.confluent.kafka.serializers.AbstractKafkaAvroSerializer by default istifade edilecek her publishde. Ve bu perfomansa menfi tesir edir.

KafkaAvroSerializer/Des.. , KafkaProtobufSerializer/Des.., KafkaJsonSchemaSerializer/Des… , GenericAvroSerde/Des.. bunlardan bir necesidir.

Kafka Connect

Kafka ozellikleri tekce bunla bitmir Connect Api ve streaming ucun api da verir.

Kafka connect ile meselem siz hanisa bazada table-a qosulub deye bilersizki bu datalar kafkaya replika olunsun dussun o eventler ora.

Kafka streaming ile ise kafkadaki eventler uzerinde aggregationlar apara, join ede bilersiz bir nece topicde olan datani ve.s. Ve sonda harasa sink ede bilersiniz neticeni.

Bunlar haqqinda etrafli novbeti meqalelerde yazacagam …

Top 20 common Kafka problems with solutions:


메타데이터
post_id
b99e5d77bccd
slug
kafka-b99e5d77bccd
url
https://medium.com/@rzaeeff/kafka-b99e5d77bccd
canonical_url
https://medium.com/@rzaeeff/kafka-b99e5d77bccd
author_url
https://medium.com/@rzaeeff
status
ok
fetched_at
2026-06-15 20:49:13