Apache Flink পরিচিতি
এই পাঠে যা শিখবেন
- Flink-এর core abstraction — DataStream, KeyedStream, ProcessFunction
- Event time, watermark, ও out-of-order event handling
- তিন ধরনের window — tumbling, sliding, session — কখন কোনটি
- Flink vs Spark Streaming — কখন কোনটি বাছবেন
১ · Flink কেন আলাদা?
Apache Flink-এর জন্ম Berlin Technical University-র Stratosphere project থেকে (২০১০)। Apache top-level হয় ২০১৪-তে। Flink-এর মূল প্রতিজ্ঞা — streaming first। Spark যেখানে batch engine-এ streaming "জুড়েছে" (micro-batch), Flink-এ কোনো batch ছিলই না — সব ইভেন্ট একে একে process হয়। Batch পরে যোগ হলো — bounded stream হিসেবে।
Spark: "batch on streams" — ছোট ছোট batch তৈরি করে process।
Flink: "streams are first-class" — প্রতিটি event আসামাত্র process; state নিজস্ব backend-এ; checkpoint asynchronously।
ফলাফল — Flink-এ latency ১০-৫০ ms; Spark-এ ১০০-৫০০ ms।
২ · DataStream API — Flink-এর ভাষা
Flink-এ data-র মূল abstraction হলো DataStream<T> — অসীম typed stream। Java বা Scala-তে লেখা হয়; Python (PyFlink) ও SQL API-ও আছে কিন্তু production-এ JVM API বেশি। প্রতিটি transformation (map, filter, keyBy, window, reduce) — DataStream → DataStream।
- Source: Kafka, Kinesis, file, socket, Pulsar।
- Transformation: map, flatMap, filter, keyBy, window, process।
- Sink: Kafka, Elasticsearch, JDBC, file, custom।
৩ · KeyedStream — partition by key
Stream-কে user_id, driver_id দিয়ে partition করতে keyBy() ব্যবহার হয়। একই key-এর সব event একই TaskManager-এ যায় — ফলে per-key state রাখা যায়। Pathao-র "প্রতি driver-এর last 5 ride" দেখার জন্য এটাই ভিত্তি।
৪ · Event time vs processing time — Flink-এর শক্তি
Flink event-time semantics-এ industry leader। প্রতিটি event-এর সাথে timestamp থাকে; out-of-order এলেও — watermark-এর মাধ্যমে Flink সঠিক window-এ assign করে। WatermarkWatermark in Flink"এই timestamp-এর আগের event আর আশা করছি না" — Flink-এর key concept। Periodic বা punctuated দু'ভাবে generate। Window/aggregation kicker হিসেবে কাজ করে। source বা assigner-এ generate হয় — periodic (every 200ms) বা per-event (punctuated)।
$$\text{watermark}(W_t) \;\Rightarrow\; \text{সব event with timestamp} < W_t \text{ already arrived}$$
৫ · Window — তিন ধরনের
- Tumbling window: ফিক্সড সাইজ, পরস্পর overlap নেই। "প্রতি ১ মিনিটে কত transaction" — Grameenphone-এর per-minute call counter-এ আদর্শ।
- Sliding window: ফিক্সড সাইজ + slide step। "শেষ ৫ মিনিটে কত click, প্রতি ১ মিনিটে recompute" — moving average-এর জন্য।
- Session window: সাইজ ফিক্সড নয় — gap-driven। User ১৫ মিনিট inactive হলে session শেষ। Daraz-এর shopping session detection-এ অপরিহার্য।
প্রতিটি window-এ reduce, aggregate, বা process apply করা যায়।
৬ · State — Flink-এর হৃদয়
Stateful streaming মানে — Flink মনে রাখে। প্রতি driver-এর last position, প্রতি card-এর last 10 transaction, প্রতি user-এর session start। State backend দু'টি প্রধান:
- HashMap (in-memory): দ্রুত, কিন্তু heap-এ সীমিত। ছোট state-এ।
- RocksDB: embedded LSM-tree, disk-backed। TB-scale state-ও সম্ভব। production default।
CheckpointFlink CheckpointAsynchronous distributed snapshot (Chandy-Lamport algorithm)। প্রতিটি operator-এর state একসাথে freeze করে durable storage-এ লেখা হয় — stream stop না করে। Failure-এ এই snapshot থেকে restore। asynchronous — stream থামানো ছাড়া S3/HDFS-এ snapshot। Failure-এ exact সেই snapshot থেকে recovery → exactly-once semantics।
৭ · Flink architecture সংক্ষেপে
- JobManager: coordinator — checkpoint trigger, scheduling, failover।
- TaskManager: worker — তাদের slot-এ operator-এর instance চলে।
- Slot: TaskManager-এর resource unit। parallelism N মানে N slot।
- Production-এ Flink Kubernetes operator বা YARN-এ deploy।
Kafka.Sink + transactional producer ব্যবহার করলে — end-to-end exactly-once guarantee। এটি financial system-এ critical (bKash, ব্যাংক ledger replication)।
৮ · কোডে — Flink Java DataStream API
Grameenphone-এর হাজার হাজার call event। Kafka topic gp_calls থেকে — প্রতি ১ মিনিট window-এ caller-প্রতি call count, ১০-এর বেশি হলে spam alert।
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.streaming.api.datastream.*;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import java.time.Duration;
public class GpSpamCheck {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60_000); // every 60s, async snapshot
KafkaSource<Call> source = KafkaSource.<Call>builder()
.setBootstrapServers("kafka:9092")
.setTopics("gp_calls")
.setGroupId("flink-spam")
.setDeserializer(new CallDeserializer())
.build();
DataStream<Call> calls = env.fromSource(
source,
WatermarkStrategy.<Call>forBoundedOutOfOrderness(Duration.ofMinutes(2))
.withTimestampAssigner((c, ts) -> c.eventTime),
"gp-calls");
DataStream<Alert> spam = calls
.keyBy(c -> c.callerMsisdn)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new CountAgg(), new EmitIfHigh(10));
spam.sinkTo(KafkaSinks.alertsTopic());
env.execute("gp-spam-check");
}
}
WatermarkStrategy.forBoundedOutOfOrderness(2 min) — সর্বোচ্চ ২ মিনিট out-of-order সহ্য। keyBy(callerMsisdn) partition by phone number। TumblingEventTimeWindows.of(1 min) ১-মিনিট bucket। aggregate-এর দু'টি অংশ — CountAgg per-record incremental count, EmitIfHigh window শেষে threshold check।
৯ · Flink SQL — সরল উপায়
Java code কঠিন মনে হলে Flink-এর SQL API আছে — ANSI SQL-এর মতো কিন্তু streaming-এ। ছোট pipeline-এ এটি সাশ্রয়ী।
-- Kafka topic 'pathao_rides' কে table হিসেবে রেজিস্টার
CREATE TABLE rides (
ride_id STRING,
driver_id STRING,
fare DECIMAL(10,2),
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' MINUTE
) WITH (
'connector' = 'kafka',
'topic' = 'pathao_rides',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json',
'scan.startup.mode' = 'latest-offset'
);
-- প্রতি ৫ মিনিট window-এ driver-প্রতি total fare
SELECT
driver_id,
TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS w_start,
SUM(fare) AS total_fare,
COUNT(*) AS ride_count
FROM rides
GROUP BY driver_id, TUMBLE(event_time, INTERVAL '5' MINUTE);
WATERMARK FOR event_time AS … — table definition-এই event time semantics ঘোষণা। TUMBLE built-in window function। SQL জানা থাকলে — Java/Scala না শিখেও Flink-এ production pipeline লেখা সম্ভব।
১০ · Flink vs Spark Streaming — সারমর্ম
- Latency: Flink ১০-৫০ ms · Spark ১০০-৫০০ ms (default)।
- Model: Flink event-by-event · Spark micro-batch।
- State: Flink RocksDB native, TB-scale OK · Spark সাধারণত ছোট state-এ ভালো।
- Watermark: Flink prefers explicit, mature handling · Spark সরল কিন্তু কম powerful।
- Ecosystem: Spark batch+ML+SQL+Streaming একসাথে · Flink pure streaming + Table API।
- Adoption: Bangladesh-এ Spark বেশি common; Alibaba, Netflix, Uber globally Flink-এ।
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ Flink-এর "exactly-once" exact কীভাবে কাজ করে? Chandy-Lamport snapshot কী, এবং Kafka-র সাথে end-to-end exactly-once কীভাবে অর্জিত হয়?
"Exactly-once" — distributed system-এর holy grail। প্রতিটি input event ঠিক একবার output-এ effect ফেলবে — না কম, না বেশি। Flink এটি অর্জন করে দু'টি পদ্ধতির মিশ্রণে।
Chandy-Lamport (১৯৮৫) — distributed snapshot:
- JobManager periodic checkpoint barrier inject করে — প্রতি source-এ একটি বিশেষ marker।
- Barrier downstream operator-এ পৌঁছালে — সেই operator তার state snapshot S3/HDFS-এ লেখে, তারপর barrier আগে পাঠায়।
- সব operator-এর barrier-aligned snapshot তৈরি হলে — পুরো job-এর consistent state image তৈরি।
- এই কাজ stream থামায় না — async, বেশিরভাগ throughput-এ ১০% impact।
Failure-এ recovery:
- সব operator গত successful checkpoint থেকে state restore করে।
- Source (Kafka) সেই checkpoint-এ রক্ষিত offset থেকে replay শুরু করে।
- কোনো event lost হয় না (durable Kafka), duplicate-ও তৈরি হয় না (state replay deterministic)।
End-to-end exactly-once — sink সমস্যা:
- Internal exactly-once সহজ। কিন্তু Kafka-তে output লেখা — checkpoint-এর মাঝে duplicate publish হলে?
- Solution: Kafka transactional producer। Flink
KafkaSinktwo-phase commit ব্যবহার করে — pre-commit (transaction open) → checkpoint complete → commit (transaction commit)। - Failure-এ uncommitted transaction abort, replay, পুনরায় commit। Consumer
isolation.level=read_committed-এ — শুধু committed দেখে।
Trade-offs:
- Latency বাড়ে — checkpoint interval-এর কাছাকাছি (১০-৬০ সেকেন্ড)। sub-second exactly-once কঠিন।
- Sink অবশ্যই idempotent বা transactional হতে হবে (Kafka, JDBC with PK, Iceberg)।
- Non-deterministic operation (random number, current time) — exact replay ভাঙতে পারে। সাবধানতা।
বাস্তব use case: bKash-এর ledger replication থেকে data warehouse-এ — ১ টাকাও duplicate হলে audit ব্যর্থ। Flink + Kafka transactional sink এই গ্যারান্টি দেয়।
মূল উপলব্ধি: "Exactly-once" magic না — কঠোর protocol। প্রতিটি stage (source, processing, sink) সঠিকভাবে cooperate করলেই সম্ভব। Flink সবগুলো সরবরাহ করে — engineer-এর কাজ এই pattern-এ নিজের code লেখা।
প্র ০২ Daraz-এ ব্যবহারকারীর "shopping session" detect করতে চান। কেন tumbling বা sliding window কাজ করবে না? Session window কীভাবে কাজ করে?
Online shopping সব user একই সময়ে শুরু-শেষ করে না। কেউ ৫ মিনিট browse করেন, কেউ ১ ঘণ্টা; কেউ একদিনে ৩ session, কেউ ১ সপ্তাহে একবার। ফিক্সড window-এ "session" ধারণা ভুল — window সীমা বাস্তব behavior-এর সাথে align করে না।
Tumbling-এর সমস্যা:
- "প্রতি ১ ঘণ্টায় কত click" — কিন্তু একজন user-এর session যদি ৪০ মিনিট চলে এবং ৫০ মিনিট-এ window boundary পার হয়, তাহলে একই session দুই window-এ ভাগ হয়ে যায়।
- "১ ঘণ্টা" পর কত-জন একটানা চলছিল — বলা যাবে না।
Sliding-এর সমস্যা:
- "শেষ ৩০ মিনিটে কত click" — দরকারী, কিন্তু "session" boundary নয়।
- একই click একাধিক window-এ counted — session count overcounted।
Session window — gap-driven:
- আপনি একটি gap timeout সেট করেন — যেমন ১৫ মিনিট inactivity।
- প্রতিটি event তার আগের event থেকে < ১৫ মিনিটে এলে — একই session-এ যুক্ত।
- ১৫+ মিনিট gap → পুরোনো session close, নতুন session শুরু।
- Window সাইজ data-নির্ভর — কারো ২ মিনিট, কারো ২ ঘণ্টা।
Flink-এ implementation:
EventTimeSessionWindows.withGap(Time.minutes(15))- প্রতিটি new event session window-এর end timestamp ১৫ মিনিট পরে set করে।
- আরেক event এসে এই window extend করতে পারে। Session merging internally complex — Flink সামলায়।
- Watermark gap timeout পার হলে — window finalize, downstream-এ session metadata (start, end, click count, total spent) emit।
Daraz-এর প্রকৃত metric:
- Avg session duration → user engagement।
- Pages per session → site quality।
- Cart-to-purchase conversion within session → checkout funnel।
- Bounce rate (single-page session) → landing page issue।
Trade-offs:
- Gap বড় (৩০ min) → bathroom break "একই session" ধরে; কিন্তু genuine session conflate হতে পারে।
- Gap ছোট (২ min) → পড়া বিরতিতেও session ভাঙে — overestimate session count।
- Industry standard: web ৩০ min, mobile app ১৫ min। Daraz-এ A/B test করে নির্ধারণ।
মূল উপলব্ধি: Window choice business semantics-এর সাথে directly tied। "User session" = behavior-defined, time-defined না। Session window সেই অস্পষ্টতা elegantly capture করে — এটি Flink-এর অন্যতম প্রিয় feature।
প্র ০৩ Bangladesh-এ একটি startup Flink ব্যবহার করতে চাইছে। Spark Streaming ইতিমধ্যে আছে। CTO হিসেবে কী advice দেবেন?
Tech migration decision শুধু technical না — team, cost, business risk সব মিশ্রণ। Bangladesh context-এ এই বিচার বিশেষ গুরুত্বপূর্ণ।
প্রথমে — Flink লাগে কিনা সেটা যাচাই:
- বর্তমান Spark Streaming-এ latency কত? Business requirement কত? যদি Spark ৫০০ ms-এ আছে, requirement ১ second — Flink শিফট করার মানে নেই।
- State size? Spark-এ যদি memory pressure বা frequent crash, Flink RocksDB মুক্তি দিতে পারে। না হলে — সমস্যা নেই।
- Out-of-order event severity? Pathao-র মতো mobile-heavy app-এ Flink-এর mature watermark handling কাজে আসতে পারে।
- Exactly-once requirement? Financial — Flink + Kafka + transactional sink সরল। Spark-এও সম্ভব কিন্তু sink-নির্ভর।
Bangladesh context-এর বাস্তবতা:
- Talent pool: Dhaka-তে Spark/PySpark engineer অনেক বেশি — bootcamp, freelance, BUET/IUT গ্র্যাজুয়েট। Flink-জানা engineer হাতে গোনা। Hiring time তিনগুণ।
- Training cost: বিদ্যমান team-কে Flink শেখাতে ৩-৬ মাস। সেই সময়ে business velocity কমবে।
- Vendor support: Databricks, AWS EMR — দু'টোতেই Spark first-class। Managed Flink (Amazon Kinesis Data Analytics, Aiven, Ververica Cloud) আছে কিন্তু region/cost issue।
- Community: Stack Overflow, blog, conference talk — Spark-এ ১০x বেশি resource। Debug কষ্ট কম।
Recommendation framework (CTO-র):
- Stay with Spark যদি: latency requirement ৫০০ ms-এর বেশি, batch+streaming mixed workload, team ছোট, AWS-এ EMR/Databricks ইতিমধ্যে।
- Migrate to Flink যদি: latency <১০০ ms hard requirement, large/complex state, financial-grade exactly-once, dedicated streaming team afford করা যায়।
- Hybrid (recommended): existing Spark workload-এ touch করবেন না। নতুন critical-latency workload — Flink-এ। দুই engine coexist করুক। Kafka common backbone।
Migration risk:
- Resume-driven design (engineer Flink চাচ্ছে CV-র জন্য) — এটা চিনে decline করুন।
- "Flink is faster" — context ছাড়া meaningless metric।
- Production migration ৬-১২ মাস। সেই সময়ে duplicate maintenance overhead।
Practical first step:
- Identify ONE specific use case যেখানে Flink সত্যিই value আনবে — যেমন bKash-এর fraud rule sub-100ms।
- ৩-মাস POC — production-এর shadow mode-এ। Real metric (latency, cost, accuracy) compare।
- Justify হলে — ধীরে ধীরে expansion। না হলে — Spark-এ থাকুন, লেখা ব্যর্থতা acknowledge করুন।
মূল উপলব্ধি: Tech excellence ≠ latest tool। CTO-র কাজ — engineering team-এর momentum, business outcome, এবং long-term cost-এর সর্বোত্তম balance। Flink চমৎকার — কিন্তু "চমৎকার" এবং "আপনার team-এর জন্য সঠিক" এক জিনিস না।
প্র ০৪ Watermark generation strategy — periodic vs punctuated কী, এবং Pathao-র মতো mobile-heavy app-এ কোনটি বাছবেন?
Watermark কীভাবে generate হয় — Flink performance ও correctness দু'টোর জন্যই critical। ভুল strategy-তে বা latency বাড়বে, বা data drop বাড়বে।
Periodic watermark:
- Flink নিয়মিত interval-এ (default ২০০ ms) source-এর সর্বোচ্চ-দেখা timestamp-এর ভিত্তিতে watermark inject করে।
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofMinutes(2))— সবচেয়ে কমন।- Pros: predictable, low overhead, লেখা সহজ।
- Cons: every event-এ watermark check হয় না — slight delay।
Punctuated watermark:
- প্রতিটি event check করে — watermark এই event-এ trigger করব কিনা।
- উপযোগী যখন stream-এ "end-of-batch" বা explicit signal থাকে।
- যেমন: e-commerce-এ "session_end" event — তখন punctuated watermark immediate emit।
- Pros: precise, signal-driven।
- Cons: per-event overhead, complex logic।
Pathao-র scenario:
- Driver GPS ping প্রতি ৫ সেকেন্ডে। মোবাইল নেটওয়ার্কে ১-২ মিনিট delay common, কখনো ১০-১৫ মিনিট (rural area, tunnel)।
- স্পষ্ট "end" signal নেই — driver app বন্ধ করলে nothing flushed।
- সঠিক choice: Periodic, bounded out-of-orderness ৫ মিনিট।
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofMinutes(5))
আরো সূক্ষ্ম strategy:
- Per-partition watermark: Kafka multi-partition থাকলে — প্রতি partition-এ আলাদা watermark; downstream-এ
min(all)। Pathao-র কিছু partition-এর driver hilly area-তে — সেগুলো overall watermark slow করবে। - Idle source detection:
.withIdleness(Duration.ofMinutes(2))— কোনো partition থেকে event না এলে সেটা ignore করে watermark advance। না হলে এক partition-এর driver phone বন্ধ হলে — পুরো job stalled। - Allowed lateness: watermark পার হওয়ার পরও নির্দিষ্ট সময় window রেখে দেওয়া — late event আসলে partial update। Result store-এ upsert চাই।
Real production pattern (Pathao live map):
- Periodic watermark, ৫ min out-of-orderness।
- Partition-aware watermark।
- Idleness detection ২ মিনিট।
- Window: tumbling ১ মিনিট event-time।
- Allowed lateness: ০ (real-time map-এ historical update কোনো মানে নেই)।
- Late event → side output → analytical store-এ ঠেলে দিন (post-hoc analysis)।
মূল উপলব্ধি: Watermark একটি "tuning knob" না — design decision। Mobile-network ভেদে strategy ভিন্ন হবে; bKash-এর in-app transaction (web+app) আর Pathao-র GPS ping ভিন্ন। প্রতিটি use case-এ trade-off measure করুন — কতটা data drop সহ্য, কতটা latency penalty সহ্য।
অনুশীলন
-
Window choose: Grameenphone-এর "প্রতি ৫ মিনিটে কোন cell tower-এ সর্বোচ্চ call drop" — কোন window?
Tumbling window (5 min)। প্রশ্নটি ফিক্সড সময়সীমার, একবার-ই গণনা চাই — overlap লাগবে না (sliding-এর প্রয়োজন নেই); session-এর behavioral concept এখানে অপ্রাসঙ্গিক।
calls.keyBy(c -> c.cellTower) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .reduce((a, b) -> new CallCount(a.tower, a.drops + b.drops)); -
Watermark calculation: stream-এ এ পর্যন্ত max event time = ১২:১০:০০। boundedOutOfOrderness = ৩ মিনিট। বর্তমান watermark কত? ১২:০৬:০০-এর event আসলে drop হবে?
Watermark = ১২:১০:০০ − ৩ min = ১২:০৭:০০।
১২:০৬:০০ event watermark-এর চেয়ে পুরোনো — drop (বা allowed lateness থাকলে late path-এ)।
সমাধান: out-of-orderness বাড়ান (৫ min), অথবা
.allowedLateness(Time.minutes(10))। -
Compare: Spark Streaming আর Flink-এর ৩টি মূল technical পার্থক্য লিখুন এবং কখন কোনটি বাছবেন।
- Execution: Spark micro-batch (১০০-৫০০ ms) · Flink event-by-event (১০-৫০ ms)। Lower latency দরকার → Flink।
- State: Spark সাধারণত memory-bound · Flink RocksDB native, TB-scale। Large/complex state → Flink।
- Ecosystem: Spark batch+ML+SQL+Streaming · Flink streaming-first। Mixed workload → Spark।
Choose Spark: Daraz analytics, ১+ second latency OK, batch ও ML একসাথে, Bangladesh team common skill। Choose Flink: bKash sub-100ms fraud, large stateful joins, financial exactly-once।
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ২১ · Change Data Capture (CDC) পরবর্তী পাঠ Database changes-কে stream হিসেবে capture — Flink/Spark-এর সঙ্গে natural fit।
- পাঠ ১৯ · Spark Structured Streaming আগের পাঠ Spark-এর streaming approach — Flink-এর সাথে comparison-এর ভিত্তি।
- পাঠ ১৭ · Apache Kafka পরিচিতি এই পাঠের সাথে সম্পর্কিত Flink-এর প্রায় সব production source — Kafka। Topic, partition, offset।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।
flink:latest image দিয়ে SQL Client শুরু করুন — Java code lab-environment-এ চালানো জটিল।