Kafka-ni qayta muvozanatlash va xabarlarni qayta ishlashdagi tor joylarni optimallashtirish

Kafka-ni qayta muvozanatlash va xabarlarni qayta ishlashdagi tor joylarni optimallashtirish

1. Muammoning yuzaga kelishi

Bizga Kafka xabarlari production muhitida to‘g‘ri iste’mol qilinmayotgani haqida so‘rov kelib tushdi.

Server jurnallarini tekshirganimizdan so‘ng, Consumer Group rebalancing jarayoni taxminan har 10 daqiqada takroran sodir bo‘layotganini aniqladik.

Rebalancing sababini aniqlash uchun Consumer sozlamalari va xabarlarni amalda qayta ishlash vaqtini tekshirdik.

O‘sha vaqtda max.poll.interval.ms 10 daqiqaga o‘rnatilgan edi, bitta Kafka xabarini qayta ishlash esa 10 daqiqadan ko‘proq vaqt olardi.

max.poll.interval.ms — Consumer poll() ni chaqirgan vaqt bilan poll() ni keyingi chaqiradigan vaqt oralig‘ida ruxsat etilgan maksimal intervaldir. Shu sababli, xabarni qayta ishlash vaqti ushbu sozlama qiymatidan oshib ketgani uchun Consumer odatdagi qayta ishlash siklini saqlab qola olmayotgani va natijada rebalancing yuz berayotganini aniqladik.

Muammo aniqlangan vaqtdagi sozlamalar quyidagicha edi.

max.poll.interval.ms=600000

Jurnallar va sozlamalarni birgalikda ko‘rib chiqish orqali bu muammo Kafka xabarlarni qabul qilayotganining o‘ziga bog‘liq emasligini, balki xabarlarni qayta ishlash vaqti haddan tashqari uzayib, Consumer’ning qayta ishlash intervali cheklovidan oshib ketgani sababli yuzaga kelganini tasdiqladik.

2. Dastlabki ekspluatatsion choralar

Sabab aniqlangandan so‘ng, production muhitida ishlayotgan Consumer’da rebalancing takroran sodir bo‘layotgan vaziyatni avvalo yumshatish zarur edi.

Shunga muvofiq, max.poll.interval.ms qiymatini avvalgi 10 daqiqadan 30 daqiqaga o‘zgartirdik.

max.poll.interval.ms=1800000

Ushbu konfiguratsiya o‘zgarishi joriy qayta ishlash vaqti 10 daqiqadan ko‘proq bo‘lgan taqdirda ham Consumer’ga xabarlarni qayta ishlash uchun yetarli vaqt berdi.

Biroq, biz bu chorani tub yechim deb hisoblamadik.

Joriy qayta ishlash vaqti 10 daqiqa bo‘lgani sababli, sozlamani 30 daqiqaga o‘zgartirish darhol yuzaga kelgan muammoni yumshatishi mumkin edi. Ammo kelajakda ma’lumotlar hajmi oshishi yoki DB qayta ishlash vaqti uzayishi sababli qayta ishlash vaqti 30 daqiqadan oshsa, xuddi shu muammo yana yuzaga kelishi mumkin edi.

Shu sababli, production muhitida avval Consumer xabarlarni normal qayta ishlay olishi uchun konfiguratsiyani o‘zgartirdik, so‘ngra xabarlarni qayta ishlash vaqti nega uzoq davom etayotganini qo‘shimcha ravishda tekshirdik.

3. Xabarlarni qayta ishlash kechikishi sababini tahlil qilish

Xabarlarni qayta ishlash mantiqini ko‘rib chiqqanimizdan so‘ng, bitta Kafka xabarini qayta ishlash jarayonida taxminan 46 000 ta yozuv MSSQL’da saqlanayotganini aniqladik.

Muammo yozuvlar sonining o‘zida — 46 000 ta ekanida — emas, balki bitta xabarni qayta ishlash uchun bir nechta DB saqlash amallari bajarilayotganida edi. Bu esa umumiy qayta ishlash vaqtining 10 daqiqadan oshishiga sabab bo‘lgan.

Kafka Consumer nuqtayi nazaridan bu bitta xabarni qayta ishlash vazifasi edi, biroq dastur ichkarida ushbu xabarni qayta ishlash uchun bir nechta DB saqlash amalini bajarayotgan edi.

Shu sababli, Consumer sozlamalarini yanada o‘zgartirishdan ko‘ra, xabarlarni qayta ishlash vaqtining o‘zini qisqartirish zarur degan xulosaga keldik.

Xabarlarni qayta ishlash mantiqini kuzatib bordik va har bir qayta ishlash bosqichini tekshirdik. Natijada DB saqlash qismi umumiy qayta ishlash vaqtining sezilarli qismini egallashini tasdiqladik.

Shunga muvofiq, Consumer’ning qayta ishlash vaqtini qisqartirish uchun ushbu DB saqlash qismini takomillashtirishga kirishdik.

Ushbu muammoning sababi va yechim yo‘nalishini quyidagicha umumlashtirish mumkin.

  • Bitta Kafka xabarini qayta ishlash 10 daqiqadan ko‘proq vaqt oldi

  • max.poll.interval.ms 10 daqiqaga o‘rnatilgan edi

  • Xabarni qayta ishlash vaqtida rebalancing takroran sodir bo‘ldi

  • DB saqlash qismi xabarni qayta ishlash vaqtida sezilarli miqdorda vaqt sarfladi

  • DB saqlash usulini takomillashtirish orqali xabarni qayta ishlashning umumiy vaqtini qisqartirishga qaror qilindi

4. DB saqlash usulini o‘zgartirish

Xabarlarni qayta ishlash mantiqini ko‘rib chiqqanimizdan so‘ng, MSSQL’da taxminan 46 000 ta yozuvni saqlash sezilarli vaqt olayotganini aniqladik.

Mavjud dastur JPA’dan foydalanardi. Ushbu takomillashtirish uchun muayyan xabarlarni qayta ishlash oqimining faqat DB saqlash qismini Batch asosidagi usulga o‘zgartirish zarur edi.

JPA’ning Batch sozlamalari ushbu sozlamalardan foydalanadigan Hibernate amallariga ta’sir qilishi mumkinligi sababli, mavjud JPA sozlamalarini saqlab qoldik va faqat tegishli saqlash mantiqiga JdbcTemplate.batchUpdate() ni qo‘lladik.

jdbcTemplate.batchUpdate(
   sql,
   data,
   batchSize,
   (PreparedStatement ps, T obj) -> {
      // 데이터 바인딩
   }
 );

Avval bir nechta yozuvni saqlash jarayonida DB saqlash amallari takroran bajarilardi. O‘zgarishdan so‘ng yozuvlar Batch birliklariga guruhlandi va birgalikda kiritildi.

Bu mavjud JPA’dan foydalanadigan boshqa saqlash mantiqlariga ta’sir qilmagan holda, muammoli DB saqlash qismining qayta ishlash vaqtini yaxshiladi.

5. Mavjud ID generatsiya usulini saqlab qolish

DB saqlash usulini o‘zgartirish davomida mavjud ID generatsiya usulini ham ko‘rib chiqdik.

Entity Hibernate’ning guid strategiyasidan foydalanardi, mavjud muhitda esa ID’lar MSSQL’ning NEWID() funksiyasi yordamida generatsiya qilinardi.

JdbcTemplate yordamida insert amallarini bevosita bajarishda avval JPA/Hibernate tomonidan boshqarilgan ID generatsiyasi jarayoni endi bevosita boshqarilishi kerak. Shu sababli, ID’lar avvalgidek usulda generatsiya qilinishini ta’minlaydigan jarayonni sozlash zarur edi.

Java’da UUID generatsiya qilib, uni uzatish ham mumkin edi, biroq ushbu vazifa uchun mavjud tizimning ID generatsiya usulini o‘zgartirmaslikni ustuvor deb bildik.

Shu sababli, Batch Insert ham MSSQL’ning NEWID() funksiyasidan foydalanadigan qilib sozlandi.

INSERT INTO TABLE_NAME ( ID, COLUMN_A, COLUMN_B)
VALUES (NEWID() , ?, ?)

Bu mavjud tizimda ishlatilayotgan ID generatsiya usulini saqlab qolgan holda, faqat DB saqlash usulini Batch Insert’ga o‘zgartirish imkonini berdi.

Ushbu vazifada muhim mezon nafaqat qayta ishlash unumdorligini yaxshilash, balki mavjud ma’lumotlarni generatsiya qilish usuliga keraksiz o‘zgarishlar kiritmaslik ham edi.

6. Takomillashtirish natijalari

Batch Insert qo‘llanilgandan so‘ng, taxminan 46 000 ta yozuvni saqlagan DB qayta ishlash qismi uchun zarur vaqt qisqardi. Bu Kafka xabarlarini qayta ishlashning umumiy vaqtini ham yaxshiladi.

Muammo yuzaga kelganida, xabarni qayta ishlash vaqti max.poll.interval.ms qiymatidan oshib ketgani sababli rebalancing takroran sodir bo‘layotgan edi.

Production muhitidagi vaziyatni yumshatish uchun avval max.poll.interval.ms qiymatini 30 daqiqaga o‘zgartirdik, so‘ngra qayta ishlash vaqtining uzayishiga sabab bo‘lgan DB saqlash qismini Batch Insert usulini qo‘llash orqali takomillashtirdik.

Natijada DB saqlash uchun zarur vaqt qisqardi, xabarni qayta ishlashning umumiy vaqti ham kamaydi va takroran yuz berayotgan rebalancing hodisasi ham yumshatildi.

Ushbu takomillashtirishdagi muhim jihat shundaki, biz max.poll.interval.ms qiymatini 30 daqiqaga oshirish orqali muammodan shunchaki qochmadik. Buning o‘rniga, konfiguratsiya o‘zgarishi orqali erishilgan ekspluatatsion yumshatishni haqiqiy qayta ishlash vaqtini qisqartirishga qaratilgan kod takomillashtirishidan ajratib ko‘rib chiqdik.

Konfiguratsiya o‘zgarishi production muhitida ishlayotgan Consumer xabarlarni qayta ishlashi uchun yetarli vaqtga ega bo‘lishini ta’minlash chorasi edi, Batch Insert’ni qo‘llash esa xabarlarni qayta ishlashning haqiqiy vaqtini qisqartirishga qaratilgan takomillashtirish edi.

Bu bizga avval ekspluatatsion hodisani yumshatish va shu bilan birga qayta ishlash kechikishining haqiqiy sababini bartaraf etish imkonini berdi.

7. Ushbu muammo orqali tasdiqlaganlarimiz

Ushbu muammoni tekshirish jarayonida Kafka Consumer’da yuz beradigan Rebalancing hodisasiga oddiygina Kafka konfiguratsiyasi muammosi sifatida yondashmaslik kerakligini tasdiqladik.

Dastlab muammo "Kafka xabarlari iste’mol qilinmayapti" tarzida xabar qilingan edi, biroq loglarda Rebalancing takroran yuz berayotganini tasdiqlaganimizdan so‘ng, Consumer konfiguratsiyasini xabarlarni amalda qayta ishlash vaqti bilan taqqoslash orqali sababni toraytira oldik.

Xususan, max.poll.interval.ms va xabarlarni amalda qayta ishlash vaqti deyarli bir xil darajada edi. Xabarlarni qayta ishlash vaqtida ma’lumotlarni DB’ga saqlash ancha vaqt olayotganini tasdiqlab, Consumer muammosiga o‘xshab ko‘ringan hodisa aslida xabarlarni qayta ishlash mantiqining bajarilish vaqti bilan bog‘liqligini aniqladik.

Shuningdek, faqat max.poll.interval.ms qiymatini oshirish muammoni hal qila olmasligini ham tasdiqladik.

Ushbu sozlamani oshirish Consumer’ga xabarlarni qayta ishlash uchun ko‘proq vaqt beradi, biroq amaldagi qayta ishlash vaqti o‘sishda davom etsa, oxir-oqibat ayni muammo yana yuz berishi mumkin.

Shu sababli, production muhitida shunga o‘xshash muammo yuzaga kelganda, muayyan konfiguratsiya qiymatlaridan tashqariga qarash va qayta ishlash vaqti qayerda oshayotganini aniqlash uchun Rebalancing qaysi nuqtada yuz berganini amaldagi xabarlarni qayta ishlash oqimi bilan birgalikda tekshirish muhim deb hisoblaymiz.

Bu holatda asosiy bottleneck DB’ga saqlash bosqichi ekani aniqlandi, biroq xuddi shu hodisa yuz bersa ham, uning haqiqiy sababi DB bo‘lmasligi mumkin.

Muhimi, muayyan texnologiya yoki sozlamani sabab deb oldindan xulosa qilmasdan, loglar, konfiguratsiya va amaldagi qayta ishlash oqimini birgalikda tekshirish orqali sababni bosqichma-bosqich toraytirishdir.

8. Xulosa

Dastlab bu muammo Kafka xabarlari to‘g‘ri iste’mol qilinmayotgan muammo sifatida namoyon bo‘ldi, biroq loglar, Consumer konfiguratsiyasi va amaldagi qayta ishlash vaqtini birgalikda tekshirish orqali Rebalancing sababini toraytira oldik.

Tekshiruvni shunchaki max.poll.interval.ms qiymatini sozlash bilan yakunlamasdan, amaldagi xabarlarni qayta ishlash oqimida qayta ishlash uzoq davom etgan DB’ga saqlash bosqichini ham aniqladik va qayta ishlash usulini yaxshiladik.

Xususan, operatsion hodisa yuz berganda joriy alomatlarni tezda yumshatish bilan haqiqiy qayta ishlash oqimini tahlil qilib, asosiy sababni yaxshilashni bir-biridan ajratish zarurligini tasdiqladik.

Kelajakda shunga o‘xshash muammo yuzaga kelsa, xato yuz bergan nuqta yoki konfiguratsiya qiymatlarini shunchaki tekshirish o‘rniga, loglar va amaldagi qayta ishlash oqimini birgalikda ko‘rib chiqamiz, sababni bosqichma-bosqich toraytiramiz va aniqlangan bottleneck’ni yaxshilaymiz.

Foydalanilgan manbalar

yeop

Site footer