Apache Flink 2026: การประมวลผล Stream, Event Time และคำถามสัมภาษณ์
คู่มือ Apache Flink 2.3 สำหรับการประมวลผล stream ด้วย event time semantics, watermark และ windowing พร้อมคำถามสัมภาษณ์และตัวอย่างโค้ด production

Apache Flink 2.3 จัดการการประมวลผล stream ในระดับใหญ่ด้วย exactly-once semantics และ latency ต่ำกว่าวินาที แตกต่างจากระบบ batch-oriented ตรงที่ Flink ประมวลผลข้อมูลอย่างต่อเนื่องเมื่อข้อมูลมาถึง ทำให้เป็น framework ที่เลือกใช้สำหรับ real-time analytics, การตรวจจับการฉ้อโกง และสถาปัตยกรรม event-driven
Flink แตกต่างจาก Spark Streaming ผ่านการประมวลผล stream แบบแท้จริง: Flink ประมวลผล event ทีละรายการด้วย event time semantics ในขณะที่ Spark Streaming ประมวลผล micro-batch โดยใช้ processing time เป็นค่าเริ่มต้น
สถาปัตยกรรม Flink 2.3 สำหรับการประมวลผล Stream
Flink ทำงานบนสถาปัตยกรรมแบบกระจายด้วย JobManager ที่ประสานงานระหว่าง TaskManager หลายตัว TaskManager แต่ละตัวรัน task slot ที่ดำเนินการส่วนต่างๆ ของ parallel operator ของ job การแยกนี้ช่วยให้ Flink scale ในแนวนอนพร้อมรักษา fault tolerance ผ่าน distributed checkpoint
โมเดล dataflow ใน Flink แสดงการคำนวณเป็น directed acyclic graph (DAG) ข้อมูลไหลจาก source ผ่าน transformation ไปยัง sink โดยแต่ละ operator สามารถทำงานบน parallel instance หลายตัว
// Basic Flink streaming application setup
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Enable checkpointing every 10 seconds for fault tolerance
env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);
// Configure state backend for large state
env.setStateBackend(new EmbeddedRocksDBStateBackend());
// Define the data source - Kafka in production scenarios
DataStream<String> rawStream = env.addSource(
new FlinkKafkaConsumer<>("events", new SimpleStringSchema(), kafkaProps)
);
// Parse and transform the stream
DataStream<Event> events = rawStream
.map(json -> objectMapper.readValue(json, Event.class))
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
);การกำหนดค่า checkpoint ด้านบนจะบันทึก snapshot ของ distributed state ทุกๆ 10 วินาที หากเกิดความล้มเหลว Flink จะกู้คืนจาก checkpoint ล่าสุดที่เสร็จสมบูรณ์และ replay record จาก Kafka
Event Time vs Processing Time Semantics
Event time หมายถึงเวลาที่ event เกิดขึ้นจริง ซึ่งฝังอยู่ในข้อมูลเอง Processing time คือเวลาที่ Flink ประมวลผล record ความแตกต่างนี้สำคัญเพราะ network delay, การส่งมอบที่ไม่เรียงลำดับ และ backlog การประมวลผลทำให้ processing time ไม่น่าเชื่อถือสำหรับการดำเนินการที่อิงตามเวลา
// Configuring event time with watermarks
public class EventTimeProcessor {
public DataStream<AggregatedMetric> processWithEventTime(
DataStream<SensorReading> readings) {
return readings
// Extract timestamp from the event payload
.assignTimestampsAndWatermarks(
WatermarkStrategy
.<SensorReading>forBoundedOutOfOrderness(Duration.ofMinutes(2))
.withTimestampAssigner((reading, ts) -> reading.getEventTime())
.withIdleness(Duration.ofMinutes(5)) // Handle idle partitions
)
// Window by event time, not wall clock
.keyBy(SensorReading::getSensorId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new AverageAggregator());
}
}กลยุทธ์ forBoundedOutOfOrderness บอก Flink ว่า event อาจมาถึงช้าได้ถึง 2 นาที Watermark จะเลื่อนไปข้างหน้าเมื่อ Flink ตัดสินว่าจะไม่มี event ที่มี timestamp ก่อน watermark มาถึงอีก
จะเกิดอะไรขึ้นกับ late event ใน Flink? โดยค่าเริ่มต้น event ที่มาถึงหลังจาก watermark ผ่านเวลาสิ้นสุดของ window จะถูกทิ้ง กำหนดค่า allowed lateness ด้วย .allowedLateness(Time.minutes(10)) เพื่อประมวลผลรายการที่มาช้า หรือใช้ side output เพื่อจับไว้สำหรับการจัดการแยก
กลยุทธ์ Windowing สำหรับ Real-Time Analytics
Flink มี window สี่ประเภท: tumbling, sliding, session และ global window แต่ละประเภทตอบสนองความต้องการด้านการวิเคราะห์ที่แตกต่างกัน
// Different windowing approaches for stream processing
public class WindowingStrategies {
// Tumbling windows: fixed-size, non-overlapping
// Use case: hourly aggregations, daily summaries
public DataStream<Summary> tumblingAggregation(DataStream<Transaction> txns) {
return txns
.keyBy(Transaction::getAccountId)
.window(TumblingEventTimeWindows.of(Time.hours(1)))
.sum("amount");
}
// Sliding windows: fixed-size, overlapping
// Use case: moving averages, rolling metrics
public DataStream<Double> slidingAverage(DataStream<Metric> metrics) {
return metrics
.keyBy(Metric::getCategory)
.window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
.aggregate(new AverageAggregator());
}
// Session windows: activity-based, variable size
// Use case: user sessions, conversation analysis
public DataStream<Session> sessionAnalysis(DataStream<ClickEvent> clicks) {
return clicks
.keyBy(ClickEvent::getUserId)
.window(EventTimeSessionWindows.withGap(Time.minutes(30)))
.process(new SessionBuilder());
}
}Session window จะปิดหลังจากช่วงเวลาที่ไม่มีกิจกรรมที่กำหนดค่าได้ รูปแบบนี้เหมาะสำหรับการวิเคราะห์พฤติกรรมผู้ใช้ที่ความยาว session แตกต่างกันตามการมีส่วนร่วม
พร้อมที่จะพิชิตการสัมภาษณ์ Data Engineering แล้วหรือยังครับ?
ฝึกฝนด้วยตัวจำลองแบบโต้ตอบ, flashcards และแบบทดสอบเทคนิคครับ
การจัดการ State และ Checkpointing
Flink รักษา operator state และ keyed state ตลอดการประมวลผล Keyed state แบ่งข้อมูลตาม key ทำให้สามารถประมวลผลแบบขนานในขณะที่เก็บ record ที่เกี่ยวข้องไว้ด้วยกัน Operator state ใช้กับทั้ง operator instance
// Managing state in a Flink KeyedProcessFunction
public class FraudDetector extends KeyedProcessFunction<String, Transaction, Alert> {
// Keyed state: one value per key (account)
private ValueState<Double> lastAmountState;
private ValueState<Long> lastTransactionTimeState;
private MapState<String, Integer> merchantCountState;
@Override
public void open(Configuration parameters) {
// Initialize state descriptors
lastAmountState = getRuntimeContext().getState(
new ValueStateDescriptor<>("lastAmount", Double.class));
lastTransactionTimeState = getRuntimeContext().getState(
new ValueStateDescriptor<>("lastTime", Long.class));
merchantCountState = getRuntimeContext().getMapState(
new MapStateDescriptor<>("merchantCounts", String.class, Integer.class));
}
@Override
public void processElement(Transaction txn, Context ctx, Collector<Alert> out)
throws Exception {
Double lastAmount = lastAmountState.value();
Long lastTime = lastTransactionTimeState.value();
// Detect suspicious patterns
if (lastAmount != null && lastTime != null) {
long timeDelta = txn.getTimestamp() - lastTime;
// Flag transactions 10x larger than previous within 1 minute
if (txn.getAmount() > lastAmount * 10 && timeDelta < 60000) {
out.collect(new Alert(txn.getAccountId(), "SUSPICIOUS_SPIKE", txn));
}
}
// Update state for next transaction
lastAmountState.update(txn.getAmount());
lastTransactionTimeState.update(txn.getTimestamp());
// Track merchant frequency
Integer count = merchantCountState.get(txn.getMerchantId());
merchantCountState.put(txn.getMerchantId(), (count == null ? 0 : count) + 1);
}
}Stateful processor นี้ติดตามรูปแบบธุรกรรมต่อบัญชี State ยังคงอยู่ตลอด checkpoint รอดพ้นจากความล้มเหลวโดยไม่สูญเสียบริบทการตรวจจับการฉ้อโกง
Flink SQL และ Table API สำหรับการประมวลผล Stream
Flink 2.3 ขยายความสามารถ SQL ด้วย Materialized Tables สำหรับการบำรุงรักษา view แบบ incremental Table API ให้อินเทอร์เฟซรวมสำหรับการประมวลผล batch และ stream
-- flink_sql_streaming.sql
-- Create a streaming source table from Kafka
CREATE TABLE orders (
order_id STRING,
customer_id STRING,
product_id STRING,
amount DECIMAL(10, 2),
order_time TIMESTAMP(3),
-- Define watermark for event time processing
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
-- Streaming aggregation with tumbling window
SELECT
customer_id,
TUMBLE_START(order_time, INTERVAL '1' HOUR) AS window_start,
COUNT(*) AS order_count,
SUM(amount) AS total_amount
FROM orders
GROUP BY
customer_id,
TUMBLE(order_time, INTERVAL '1' HOUR);วิธีการ SQL ทำให้การพัฒนาง่ายขึ้นสำหรับนักวิเคราะห์ที่คุ้นเคยกับ SQL ในขณะที่ Flink จัดการความซับซ้อนของ stream processing อยู่เบื้องหลัง
Flink vs Spark Structured Streaming
ทั้งสอง framework ประมวลผลข้อมูล streaming แต่สถาปัตยกรรมแตกต่างกันโดยพื้นฐาน Flink ประมวลผล record ทีละรายการด้วย true streaming ในขณะที่ Spark ประมวลผล micro-batch สำหรับ การเปรียบเทียบ Apache Spark การแลกเปลี่ยนระหว่าง latency และ consistency มีความสำคัญใน production
| ด้าน | Flink | Spark Structured Streaming |
|---|---|---|
| โมเดลการประมวลผล | True streaming | Micro-batch |
| Latency | มิลลิวินาที | วินาที (batch interval) |
| State Backend | RocksDB, HashMaps | In-memory, HDFS |
| Exactly-Once | Native กับ checkpoint | ต้องการ sink แบบ idempotent |
| Event Time | รองรับแบบ first-class | รองรับตั้งแต่ 2.1 |
| รองรับ SQL | Full streaming SQL | Windowing จำกัด |
เมื่อถูกถามเกี่ยวกับ Flink vs Spark สำหรับ streaming ให้มุ่งเน้นที่ความเหมาะสมกับ use case Flink เป็นเลิศในการประมวลผล event ที่ latency ต่ำและรูปแบบ event ที่ซับซ้อน Spark Streaming เหมาะสำหรับองค์กรที่รัน Spark สำหรับ batch อยู่แล้วและต้องการการประมวลผล batch-stream แบบรวม
คำถามสัมภาษณ์ Flink ที่พบบ่อยและคำตอบ
Flink บรรลุ exactly-once semantics ได้อย่างไร?
Flink รวม checkpointing กับ two-phase commit สำหรับ sink ที่รองรับ transaction ระหว่าง checkpoint Flink จะ snapshot operator state และบันทึก source offset สำหรับ Kafka sink Flink จะ pre-commit record ไปยัง Kafka เสร็จสิ้น checkpoint แล้วจึง commit transaction หากเกิดความล้มเหลวก่อน checkpoint เสร็จสมบูรณ์ record ที่ยังไม่ commit จะถูกยกเลิกและการประมวลผลจะดำเนินต่อจาก checkpoint ล่าสุด
อธิบายการแพร่กระจาย watermark ใน topology หลาย source
เมื่อ job อ่านจากหลาย partition หรือ source แต่ละอันจะสร้าง watermark ของตัวเองตาม event ที่เข้ามา Watermark ของ downstream operator จะเท่ากับ watermark ต่ำสุดในทุก input channel สิ่งนี้ทำให้มั่นใจว่าไม่มี window ใดปิดก่อนเวลาเพราะ partition หนึ่งเร็วกว่า partition ที่ช้ากว่า กำหนดค่า withIdleness() เพื่อเลื่อน watermark เมื่อบาง partition หยุดส่งข้อมูล
อะไรทำให้เกิด backpressure ใน Flink และวินิจฉัยอย่างไร?
Backpressure เกิดขึ้นเมื่อ downstream operator ไม่สามารถตามทันอัตราข้อมูล upstream Flink Web UI แสดงสถานะ backpressure ต่อ operator สาเหตุทั่วไปรวมถึง:
- การเรียกระบบภายนอกที่ช้า (database query, API call)
- การคำนวณที่แพงในฟังก์ชัน map/process
- Parallelism ไม่เพียงพอสำหรับปริมาณข้อมูล
- การดำเนินการ state ขนาดใหญ่ที่บล็อกการประมวลผล
แก้ไขโดยเพิ่ม parallelism ปรับแต่งการดำเนินการที่ช้า หรือใช้ async I/O สำหรับการเรียกภายนอก
Savepoint แตกต่างจาก checkpoint อย่างไร?
Checkpoint เป็นแบบอัตโนมัติ incremental และปรับให้เหมาะสมสำหรับการกู้คืนจากความล้มเหลว Flink จัดการ lifecycle ของมัน ลบอันเก่าโดยอัตโนมัติ Savepoint ถูกเรียกใช้โดยผู้ใช้ เป็น snapshot สมบูรณ์ที่มีไว้สำหรับงานปฏิบัติการ: deploy โค้ดใหม่ rescale job หรือ migrate ระหว่าง cluster Savepoint คงอยู่จนกว่าจะถูกลบอย่างชัดเจนและรองรับ schema evolution
การ Deploy Flink บน Kubernetes
Flink Kubernetes Operator 1.15 ทำให้การ deploy ง่ายขึ้นด้วย FlinkDeployment custom resource มันจัดการ job lifecycle, upgrade และ scaling
# flink-deployment.yaml
# Kubernetes deployment for a Flink application
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: fraud-detection-job
spec:
image: flink:1.20
flinkVersion: v1_20
flinkConfiguration:
taskmanager.numberOfTaskSlots: "4"
state.backend.type: rocksdb
state.checkpoints.dir: s3://flink-checkpoints/fraud-detection
execution.checkpointing.interval: "30s"
serviceAccount: flink
jobManager:
resource:
memory: "2048m"
cpu: 1
taskManager:
resource:
memory: "4096m"
cpu: 2
replicas: 3
job:
jarURI: s3://flink-artifacts/fraud-detection-1.0.jar
entryClass: com.example.FraudDetectionJob
parallelism: 12
upgradeMode: savepointการตั้งค่า upgradeMode: savepoint ทำให้มั่นใจว่า operator จะสร้าง savepoint ก่อน upgrade รักษา state ตลอดการ deployment
การปรับแต่งแอปพลิเคชัน Flink สำหรับ Production
การ deploy production ต้องให้ความสนใจกับการกำหนดค่า parallelism, memory และ state backend ดู รูปแบบ ETL และ data pipeline สำหรับข้อพิจารณาด้านการผสานรวม
// Production-ready Flink configuration
public class ProductionConfig {
public static void configureForProduction(StreamExecutionEnvironment env) {
// Checkpoint configuration
CheckpointConfig checkpointConfig = env.getCheckpointConfig();
checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
checkpointConfig.setMinPauseBetweenCheckpoints(30000); // 30 seconds
checkpointConfig.setCheckpointTimeout(600000); // 10 minutes
checkpointConfig.setMaxConcurrentCheckpoints(1);
checkpointConfig.setExternalizedCheckpointRetention(
ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
// State backend with incremental checkpoints
EmbeddedRocksDBStateBackend rocksDB = new EmbeddedRocksDBStateBackend(true);
rocksDB.setDbStoragePath("/tmp/rocksdb");
env.setStateBackend(rocksDB);
// Restart strategy with exponential backoff
env.setRestartStrategy(RestartStrategies.exponentialDelayRestart(
Duration.ofSeconds(1), // Initial delay
Duration.ofMinutes(5), // Max delay
2.0, // Backoff multiplier
Duration.ofHours(1), // Reset backoff after
0.1 // Jitter
));
}
}Checkpoint แบบ incremental ลดขนาด checkpoint โดยเขียนเฉพาะ state ที่เปลี่ยนแปลงตั้งแต่ checkpoint ล่าสุด การปรับแต่งนี้กลายเป็นสิ่งสำคัญเมื่อจัดการ keyed state ขนาด gigabyte
เริ่มฝึกซ้อมเลย!
ทดสอบความรู้ของคุณด้วยตัวจำลองสัมภาษณ์และแบบทดสอบเทคนิคครับ
ประเด็นสำคัญสำหรับการประมวลผล Stream Apache Flink
- Flink 2.3 ประมวลผล event ทีละรายการด้วย latency ระดับมิลลิวินาที ไม่เหมือนระบบ micro-batch
- Event time semantics กับ watermark จัดการข้อมูลที่ไม่เรียงลำดับได้อย่างถูกต้อง ตอบคำถามสัมภาษณ์ทั่วไปเกี่ยวกับ late event
- Keyed state แบ่งข้อมูลสำหรับการประมวลผลแบบขนานในขณะที่เก็บ record ที่เกี่ยวข้องไว้ด้วยกัน
- Checkpoint ให้การรับประกัน exactly-once ผ่าน distributed snapshot และ two-phase commit
- Kubernetes Operator ทำให้การ deploy, scaling และ upgrade เป็นอัตโนมัติพร้อมการรักษา state ด้วย savepoint
- เลือก Flink แทน Spark Streaming เมื่อ latency ต่ำกว่าวินาทีหรือรูปแบบการประมวลผล event ที่ซับซ้อนมีความสำคัญ
- กำหนดค่า RocksDB ด้วย checkpoint แบบ incremental สำหรับ workload production ที่มี state ขนาดใหญ่
คุณหาบั๊กใน Data Engineering เจอไหม
โค้ดจริงหนึ่งชิ้น บั๊กที่ซ่อนอยู่หนึ่งจุด วันละหนึ่งครั้ง ลองได้โดยไม่ต้องมีบัญชี

เขียนโดย
Anthony Fillion-Mailletผู้ก่อตั้ง SharpSkill
เป็นนักพัฒนาฟูลสแตกมากว่า 10 ปี ดูแล SharpSkill และรับผิดชอบทุกสิ่งที่เผยแพร่ที่นี่
อัปเดตเมื่อ 28 สิงหาคม 2569
แชร์
บทความที่เกี่ยวข้อง

Apache Beam vs Spark 2026: เปรียบเทียบ Pipeline แบบรวมและคำถามสัมภาษณ์
คู่มือครบถ้วนเปรียบเทียบ Apache Beam 2.76 และ Spark 4.2 สำหรับ data engineering เรียนรู้ความแตกต่างด้านสถาปัตยกรรม windowing ประสิทธิภาพ และคำถามสัมภาษณ์ที่พบบ่อย

Apache Spark 4.2 vs Databricks 2026: สถาปัตยกรรม ประสิทธิภาพ และคำถามสัมภาษณ์
เปรียบเทียบเชิงลึก Apache Spark 4.2 vs Databricks สำหรับปี 2026 เรียนรู้ความแตกต่างด้านสถาปัตยกรรม การแลกเปลี่ยนด้านประสิทธิภาพ ฟีเจอร์ล่าสุด และเตรียมคำถามสัมภาษณ์ data engineering

Delta Lake vs Apache Iceberg 2026: สถาปัตยกรรม Lakehouse และคำถามสัมภาษณ์
คู่มือฉบับสมบูรณ์เปรียบเทียบ Delta Lake และ Apache Iceberg สำหรับสถาปัตยกรรม lakehouse รวมถึงตัวอย่างโค้ด แนวปฏิบัติที่ดี และคำถามสัมภาษณ์ data engineering 2026