Bosh sahifa Wiki Flink

Flink

Flink — cheksiz va cheklangan data oqimlarini stateful tarzda qayta ishlashga mo‘ljallangan distributed stream processing engine. U event time, window, checkpoint, state va exactly-oncega yaqin processing semantikasiga urg‘u beradi.

Flink batch data’ni ham chegaralangan stream sifatida qayta ishlashi mumkin.

Stream

Stream ketma-ket keladigan eventlar oqimi.

Event:

ga ega bo‘lishi mumkin.

Source message broker, file, socket, database change log yoki custom connector bo‘lishi mumkin.

Job graph

Application transformationlari logical graph hosil qiladi.

Masalan:

source
→ parse
→ filter
→ keyBy
→ window
→ aggregate
→ sink

Flink graphni parallel tasklarga aylantiradi.

Operator

Operator stream ustida amal bajaradi.

Misollar:

Har operator parallelismga ega bo‘lishi mumkin.

Keyed stream

keyBy eventlarni key bo‘yicha logical partitionlarga ajratadi.

Bir xil keyga tegishli eventlar bitta parallel operator instance’iga tushadi.

Bu per-user, per-device yoki per-order state saqlashga yordam beradi.

Hot key load skew yaratadi.

State

State operator oldingi eventlar haqida ma’lumot saqlaydi.

Turlari:

State memory yoki persistent backend’da saqlanishi mumkin.

State backend

State qayerda va qanday saqlanishini backend belgilaydi.

Talablar:

Katta keyed state uchun diskka asoslangan backend ishlatilishi mumkin.

Checkpoint

Flink ma’lum intervalda operator state va source positionini barqaror storage’ga snapshot qiladi.

Failure bo‘lsa job oxirgi successful checkpointdan tiklanadi.

Source eventlar qayta o‘qilishi mumkin.

Sink idempotent yoki transactional bo‘lsa duplicate side effect kamayadi.

Barrier

Checkpoint barrier stream ichida tarqalib, operatorlar snapshot uchun consistent nuqtani belgilaydi.

Ko‘p inputli operatorda barrier alignment talab qilinishi mumkin.

Bitta input sekin bo‘lsa alignment backpressure yaratadi.

Unaligned checkpoint ayrim holatda in-flight data’ni ham snapshot qiladi.

Savepoint

Savepoint — boshqariladigan, uzoqroq saqlanadigan application state snapshoti.

U:

uchun ishlatiladi.

Savepoint compatibility operator ID va state schema’ga bog‘liq.

Event time

Event time hodisa manbada sodir bo‘lgan vaqt.

Processing time esa engine eventni qayta ishlagan vaqt.

Event kech yoki tartibsiz kelishi mumkin.

Analitik natija real hodisa vaqtiga asoslanishi kerak bo‘lsa event time ishlatiladi.

Watermark

Watermark event time oqimining taxminan qayergacha yetganini bildiradi.

Masalan, watermark 10:00 bo‘lsa tizim 10:00dan oldingi eventlarning katta qismi kelgan deb hisoblaydi.

Juda tez watermark late eventlarni yo‘qotadi.

Juda sekin watermark window natijasini kechiktiradi.

Window

Cheksiz stream cheklangan guruhlarga ajratiladi.

Turlari:

  • tumbling;
  • sliding;
  • session;
  • global;
  • custom.

Window key va vaqt asosida state saqlaydi.

Window tugaganda aggregate sinkga yuboriladi.

Late event

Watermark’dan keyin eski timestamp bilan kelgan event late hisoblanadi.

Policy:

  • tashlab yuborish;
  • allowed lateness ichida yangilash;
  • side output;
  • correction event.

Sink oldingi natijani update qila olishi kerak.

Backpressure

Downstream operator incoming event tezligiga yetisha olmasa queue to‘ladi.

Bosim upstream’ga tarqaladi.

Sabablar:

Backpressure metriclari bottleneckni ko‘rsatadi.

Parallelism

Har operator bir nechta task instance bilan ishlaydi.

Keyed state key guruhlari bo‘yicha taqsimlanadi.

Parallelism o‘zgarganda state rescale qilinadi.

Source partition soni maksimal foydali parallelismni cheklashi mumkin.

Exactly-once

Checkpoint source position va operator state’ni consistent saqlaydi.

Failure’dan keyin event qayta o‘qilsa state aynan checkpoint holatiga qaytadi.

Tashqi sink uchun exactly-once transactional commit yoki idempotent upsert talab qiladi.

Engine ichidagi kafolat tashqi emailni avtomatik exactly-once qilmaydi.

Connector

Connector source yoki sink bilan integratsiya qiladi.

U:

ni boshqaradi.

Connector capability umumiy processing kafolatiga ta’sir qiladi.

Table va SQL API

Stream va batch ma’lumot table sifatida query qilinishi mumkin.

SQL:

bilan ishlaydi.

Dynamic table update, insert va delete eventlari bilan o‘zgaradi.

Timer

Keyed process function event time yoki processing time timer o‘rnatishi mumkin.

Timer ma’lum vaqtda callback chaqirib:

bajaradi.

Timerlar checkpoint state’ining bir qismi bo‘lishi mumkin.

Rescaling

Job parallelism o‘zgarganda keyed state key-group’lar bo‘yicha yangi tasklarga taqsimlanadi.

Savepoint yoki checkpointdan restore orqali rescale qilinadi.

Operator UID o‘zgarsa eski state bilan bog‘lanish yo‘qolishi mumkin.

Operator chaining

Ketma-ket yengil operatorlar bitta task threadida chain qilinishi mumkin.

Bu serialization va network almashuvini kamaytiradi.

Failure isolation yoki alohida parallelism kerak bo‘lsa chaining uziladi.

Idle source

Bir source partition event yubormasa umumiy watermark ortda qolishi mumkin.

Idle timeout shu partitionni vaqtincha watermark hisobidan chiqaradi.

Partition yana faol bo‘lganda kech event siyosati qo‘llanadi.

State TTL

Keyed state ma’lum muddat ishlatilmasa avtomatik tozalanishi mumkin.

TTL state hajmini cheklaydi, ammo cleanup va expired qiymatning ko‘rinish semantikasi backendga bog‘liq.

Bog‘liq tushunchalar

Stream processing, Event time, Watermark, Window, Stateful processing, Checkpoint, Savepoint, Backpressure, Exactly-once, Connector