Bosh sahifa Wiki Stream processing

Stream processing

Stream processing — chegarasi oldindan ma’lum bo‘lmagan eventlar oqimini kelishi bilan uzluksiz qayta ishlash modelidir. Batch tizimi yopiq datasetni bir ish sifatida ko‘rsa, stream engine uzoq yashovchi operatorlar orqali filter, transform, join, window va agregatsiya bajaradi. Natija millisekundlardan daqiqalargacha kechikishda yangilanishi mumkin.

Event oqimi

Event biznes hodisasi, sensor o‘lchovi, log yoki database o‘zgarishi bo‘lishi mumkin. Broker eventlarni topic va partitionlarda saqlaydi. Consumer offset orqali o‘qish joyini belgilaydi. Partition ichida tartib saqlanishi mumkin, ammo turli partitionlar orasida global tartib kafolatlanmaydi.

Operatorlar stateless yoki stateful bo‘ladi. Filter bitta eventni mustaqil tekshiradi. Running count, session yoki stream join oldingi eventlar holatini saqlaydi. State hajmi kalitlar soni va window umriga bog‘liq; backend checkpoint hamda incremental snapshot bilan uni tiklaydi.

Vaqt semantikasi

Processing time event operatorga kelgan tizim vaqtiga asoslanadi. Event time hodisa manbada sodir bo‘lgan vaqtni ishlatadi va tarmoq kechikishi hamda qayta ishlashdan kamroq ta’sirlanadi. Ingestion time oraliq variantdir. Event-time natija tartibsiz va kechikkan eventlarni hisobga olish uchun watermark talab qiladi.

Tumbling window kesishmaydigan intervallar, sliding window esa qadamdan katta uzunlik sabab ustma-ust intervallar yaratadi. Session window bir kalitdagi hodisalarni inactivity gap bilan ajratadi. Window yopilgandan keyin kelgan late event drop qilinishi, alohida oqimga yuborilishi yoki oldingi natijani update qilishi mumkin.

Yetkazish kafolati

At-most-once yo‘qotish mumkin, lekin takrorlamaydi. At-least-once retry sabab duplicate berishi mumkin. Exactly-once state update va source offsetni izchil checkpointga bog‘laydi. External APIga yuborilgan side effect engine checkpointi bilan atomar bo‘lmasa, idempotent sink yoki transactional connector kerak.

Backpressure sink yoki operator sekinlashganda upstream oqimini cheklaydi. Broker retention vaqtincha buffer beradi, ammo backlog davom etsa disk to‘ladi yoki eventlar muddati tugaydi. Consumer lag, watermark delay, checkpoint duration va output latency kuzatiladi.

Parallelizm va qayta taqsimlash

Keyed state bir xil key eventlarini ayni logical operator instancega yo‘naltiradi. Parallelizm o‘zgarsa state partitionlari rescale qilinadi. Skewli key bitta taskni band qiladi. Local pre-aggregation, salting yoki maxsus heavy-key yo‘li yukni kamaytirishi mumkin.

Stream joinning ikki tomoni turli vaqtda keladi. Tizim mos eventni kutish uchun state saqlaydi va watermark asosida endi match kelmasligini taxmin qiladi. Chegarasiz join retention bo‘lmasa state cheksiz o‘sadi.

Qo‘llanish

Fraud aniqlash, real-time dashboard, monitoring, personalization va CDC pipeline oqimli qayta ishlashdan foydalanadi. “Real-time” qat’iy son emas; biznes deadline bilan belgilanadi. Ba’zi tizimda besh daqiqalik micro-batch yetarli, boshqasida 100 ms kerak. Arxitektura correctness, replay, operatsion murakkablik va xarajat bilan birga baholanadi.

Schema evolyutsiyasi

Uzoq yashovchi streamda producer va consumer bir vaqtda yangilanmasligi mumkin. Field qo‘shish optional va default bilan backward-compatible bo‘lishi, type almashtirish esa yangi field yoki versiya talab qilishi mumkin. Schema registry compatibilityni tekshiradi. Consumer noma’lum fieldni e’tiborsiz qoldirishi, ammo required field yo‘qligini jim qabul qilmasligi kerak. Replay eski event versiyalarini ham o‘qiydi, shuning uchun decoder migrationdan keyin tarixiy schema ni qo‘llab turadi.

Dead-letter oqimi

Parse yoki biznes validatsiyasidan o‘tmagan event sabab, schema va raw payload bilan alohida oqimga yuboriladi. Dead-letter queue monitoring, retention va qayta ishlash egasiga ega bo‘lmasa, jim data yo‘qotish omboriga aylanadi.

Bog‘liq tushunchalar

Event time, Watermark, Window, Kafka, Stateful processing, Exactly-once, Backpressure