Watermark — event-time stream processingda tizimning vaqt bo‘yicha qanchalik oldinga siljiganini bildiruvchi mantiqiy chegara. U odatda “watermarkdan eski timestampli eventlarning aksariyati allaqachon keldi” degan taxminni ifodalaydi. Watermark barcha event kelganini isbotlamaydi; window natijasini chiqarish, kech event siyosati va eski state ni tozalash uchun nazorat nuqtasi beradi.
Hisoblash
Oddiy bounded-out-of-orderness usuli kuzatilgan eng katta event timestampidan ruxsat etilgan kechikishni ayiradi:
watermark = max_seen_event_time - allowed_delay
Agar eng yangi event 12:10, delay ikki daqiqa bo‘lsa, watermark 12:08 bo‘ladi. 12:05–12:06 oynasi endi yopilishi mumkin. Keyin 12:05:30 eventi kelsa u late hisoblanadi.
Monoton timestampli source uchun watermark oxirgi eventga yaqin yuradi. Periodic generator ma’lum intervalda signal chiqaradi. Punctuated watermark maxsus event belgisi bilan keladi. Source o‘z partition progressini bilsa, watermarkni broker metadata si yoki domain markeridan ishlab chiqishi mumkin.
Parallel partitionlar
Bir nechta input partition bo‘lsa, operatorning watermarki ko‘pincha ularning minimumiga teng. Eng sust partition barcha windowlarni ushlab qoladi. Partition idle bo‘lib, event kelmayotgan bo‘lsa, tizim uni vaqtincha inactive deb belgilashi mumkin; aks holda watermark abadiy to‘xtaydi. Idle detection noto‘g‘ri bo‘lsa, keyin qaytgan eski eventlar late bo‘ladi.
Ikki stream joinida ikkala input watermarki state cleanupga ta’sir qiladi. Bir tomon tez, ikkinchisi kechiksa, mos kelishi mumkin bo‘lgan recordlar yetarli muddat saqlanadi. Watermarklarni ko‘r-ko‘rona tenglashtirish valid matchni o‘chirishi mumkin.
Window natijasi
Tumbling window tugash vaqti watermarkdan o‘tganda trigger final yoki dastlabki natijani chiqaradi. Allowed lateness davomida kech event kelganda agregat update qilinadi. Sink changelog, upsert yoki retractionni tushunishi kerak. Faqat append qiladigan sink bir oynaning bir nechta versiyasini alohida satr qilib yuborishi mumkin.
Watermark sekin bo‘lsa natija kechikadi va state, checkpoint hamda disk hajmi o‘sadi. Juda agressiv bo‘lsa ko‘p event late bo‘lib, dashboard kam ko‘rsatadi. Siyosat arrival delay taqsimoti va biznesning correctness–latency muvozanatiga asoslanadi.
Nosozlik va replay
Checkpoint watermark va operator state bilan izchil saqlanadi. Recoverydan keyin source offset qayta o‘qilganda watermark orqaga ketmasligi yoki window oldin tozalangan statega event qo‘shmasligi kerak. Framework watermarkni mantiqiy progressning bir qismi sifatida tiklaydi.
Backfill tarixiy eventlarni joriy streamga aralashtirsa, ularning timestampi watermarkdan ancha eski bo‘ladi. Bunday ish alohida pipeline, maxsus watermark yoki batch rejimida bajariladi. Monitoring current watermark, wall-clockdan farq, late event ulushi, idle partition va window state hajmini ko‘rsatadi.
Watermark alignment
Bir source partition tez, boshqasi juda sekin bo‘lsa tez tomon uzoq kelajak state ni yig‘ib, join xotirasini oshiradi. Watermark alignment tez sourceni vaqtincha throttle qilib, partitionlar progressini ma’lum farq ichida ushlaydi. Bu throughputni pasaytirishi mumkin, ammo checkpoint va state hajmini nazorat qiladi. Alignment faqat event-time progress farqini boshqaradi; data skew yoki sekin sinkning o‘zini hal qilmaydi.
Bo‘sh oqim
Event yo‘q bo‘lganda watermark yaratish domen qaroridir. Wall-clock asosida oldinga yurish kech offline eventni late qiladi, umuman yurmaslik esa window chiqishini to‘xtatadi. Tizim inactivity SLA va source xususiyatiga mos siyosat tanlaydi.
Bog‘liq tushunchalar
Event time, Window, Late data, Stream processing, Trigger, Stateful processing, Checkpoint