Apache Flink 2026: การประมวลผล Stream, Event Time และคำถามสัมภาษณ์

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

Apache Flink 2026: การประมวลผล Stream, Event Time และคำถามสัมภาษณ์

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 ทำงานบนสถาปัตยกรรมแบบกระจายด้วย 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 หลายตัว

FlinkStreamJob.javajava
// 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 ไม่น่าเชื่อถือสำหรับการดำเนินการที่อิงตามเวลา

EventTimeExample.javajava
// 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 แต่ละประเภทตอบสนองความต้องการด้านการวิเคราะห์ที่แตกต่างกัน

WindowingStrategies.javajava
// 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

StatefulProcessor.javajava
// 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 2.3 ขยายความสามารถ SQL ด้วย Materialized Tables สำหรับการบำรุงรักษา 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 อยู่เบื้องหลัง

ทั้งสอง framework ประมวลผลข้อมูล streaming แต่สถาปัตยกรรมแตกต่างกันโดยพื้นฐาน Flink ประมวลผล record ทีละรายการด้วย true streaming ในขณะที่ Spark ประมวลผล micro-batch สำหรับ การเปรียบเทียบ Apache Spark การแลกเปลี่ยนระหว่าง latency และ consistency มีความสำคัญใน production

ด้านFlinkSpark Structured Streaming
โมเดลการประมวลผลTrue streamingMicro-batch
Latencyมิลลิวินาทีวินาที (batch interval)
State BackendRocksDB, HashMapsIn-memory, HDFS
Exactly-OnceNative กับ checkpointต้องการ sink แบบ idempotent
Event Timeรองรับแบบ first-classรองรับตั้งแต่ 2.1
รองรับ SQLFull streaming SQLWindowing จำกัด
ข้อมูลเชิงลึกสำหรับการสัมภาษณ์

เมื่อถูกถามเกี่ยวกับ Flink vs Spark สำหรับ streaming ให้มุ่งเน้นที่ความเหมาะสมกับ use case Flink เป็นเลิศในการประมวลผล event ที่ latency ต่ำและรูปแบบ event ที่ซับซ้อน Spark Streaming เหมาะสำหรับองค์กรที่รัน Spark สำหรับ batch อยู่แล้วและต้องการการประมวลผล batch-stream แบบรวม

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

Flink Kubernetes Operator 1.15 ทำให้การ 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

การ deploy production ต้องให้ความสนใจกับการกำหนดค่า parallelism, memory และ state backend ดู รูปแบบ ETL และ data pipeline สำหรับข้อพิจารณาด้านการผสานรวม

ProductionConfig.javajava
// 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

เริ่มฝึกซ้อมเลย!

ทดสอบความรู้ของคุณด้วยตัวจำลองสัมภาษณ์และแบบทดสอบเทคนิคครับ

  • 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

เขียนโดย

Anthony Fillion-Maillet

ผู้ก่อตั้ง SharpSkill

เป็นนักพัฒนาฟูลสแตกมากว่า 10 ปี ดูแล SharpSkill และรับผิดชอบทุกสิ่งที่เผยแพร่ที่นี่

อัปเดตเมื่อ 28 สิงหาคม 2569

แชร์

บทความที่เกี่ยวข้อง

Apache Beam vs Spark 2026 เปรียบเทียบสำหรับ data engineering

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: สถาปัตยกรรม ประสิทธิภาพ และคำถามสัมภาษณ์

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

Delta Lake vs Apache Iceberg 2026: สถาปัตยกรรม Lakehouse และคำถามสัมภาษณ์

Delta Lake vs Apache Iceberg 2026: สถาปัตยกรรม Lakehouse และคำถามสัมภาษณ์

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