# Apache Flink 2026: การประมวลผล Stream, Event Time และคำถามสัมภาษณ์ > คู่มือ Apache Flink 2.3 สำหรับการประมวลผล stream ด้วย event time semantics, watermark และ windowing พร้อมคำถามสัมภาษณ์และตัวอย่างโค้ด production - Published: 2026-08-28 - Updated: 2026-08-28 - Author: Anthony Fillion-Maillet - Reading time: 5 min --- 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 หลายตัว ```java // FlinkStreamJob.java // 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 rawStream = env.addSource( new FlinkKafkaConsumer<>("events", new SimpleStringSchema(), kafkaProps) ); // Parse and transform the stream DataStream events = rawStream .map(json -> objectMapper.readValue(json, Event.class)) .assignTimestampsAndWatermarks( WatermarkStrategy.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 ไม่น่าเชื่อถือสำหรับการดำเนินการที่อิงตามเวลา ```java // EventTimeExample.java // Configuring event time with watermarks public class EventTimeProcessor { public DataStream processWithEventTime( DataStream readings) { return readings // Extract timestamp from the event payload .assignTimestampsAndWatermarks( WatermarkStrategy .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 แต่ละประเภทตอบสนองความต้องการด้านการวิเคราะห์ที่แตกต่างกัน ```java // WindowingStrategies.java // Different windowing approaches for stream processing public class WindowingStrategies { // Tumbling windows: fixed-size, non-overlapping // Use case: hourly aggregations, daily summaries public DataStream tumblingAggregation(DataStream 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 slidingAverage(DataStream 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 sessionAnalysis(DataStream clicks) { return clicks .keyBy(ClickEvent::getUserId) .window(EventTimeSessionWindows.withGap(Time.minutes(30))) .process(new SessionBuilder()); } } ``` Session window จะปิดหลังจากช่วงเวลาที่ไม่มีกิจกรรมที่กำหนดค่าได้ รูปแบบนี้เหมาะสำหรับการวิเคราะห์พฤติกรรมผู้ใช้ที่ความยาว session แตกต่างกันตามการมีส่วนร่วม ## การจัดการ State และ Checkpointing Flink รักษา operator state และ keyed state ตลอดการประมวลผล Keyed state แบ่งข้อมูลตาม key ทำให้สามารถประมวลผลแบบขนานในขณะที่เก็บ record ที่เกี่ยวข้องไว้ด้วยกัน Operator state ใช้กับทั้ง operator instance ```java // StatefulProcessor.java // Managing state in a Flink KeyedProcessFunction public class FraudDetector extends KeyedProcessFunction { // Keyed state: one value per key (account) private ValueState lastAmountState; private ValueState lastTransactionTimeState; private MapState 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 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](https://flink.apache.org/downloads/) สำหรับการบำรุงรักษา view แบบ incremental Table API ให้อินเทอร์เฟซรวมสำหรับการประมวลผล batch และ stream ```sql -- 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](/blog/data-engineering/apache-spark-4-new-features-structured-streaming-interview) การแลกเปลี่ยนระหว่าง 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](https://flink.apache.org/2026/05/26/apache-flink-kubernetes-operator-1.15.0-release-announcement/) ทำให้การ deploy ง่ายขึ้นด้วย FlinkDeployment custom resource มันจัดการ job lifecycle, upgrade และ scaling ```yaml # 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](/technologies/data-engineering/interview-questions/etl-elt-patterns) สำหรับข้อพิจารณาด้านการผสานรวม ```java // ProductionConfig.java // 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 ขนาดใหญ่ --- Source: SharpSkill (https://sharpskill.dev), tech interview preparation for your real stack. HTML version of this page: https://sharpskill.dev/th/blog/data-engineering/apache-flink-stream-processing-event-time-interview-2026