# Apache Flink 2026: Pemrosesan Stream, Event Time, dan Pertanyaan Interview > Pelajari Apache Flink 2.3 untuk pemrosesan stream dengan semantik event time, watermark, dan windowing. Panduan lengkap dengan pertanyaan interview dan contoh kode produksi. - Published: 2026-08-28 - Updated: 2026-08-28 - Author: Anthony Fillion-Maillet - Reading time: 5 min --- Apache Flink 2.3 menangani pemrosesan stream dalam skala besar dengan semantik exactly-once dan latensi sub-detik. Berbeda dengan sistem berbasis batch, Flink memproses data secara kontinu saat data tiba, menjadikannya framework pilihan utama untuk analitik real-time, deteksi fraud, dan arsitektur event-driven. > **Penting untuk Interview** > > Flink membedakan dirinya dari Spark Streaming melalui pemrosesan stream sejati: Flink memproses event satu per satu dengan semantik event time, sementara Spark Streaming memproses micro-batch dengan processing time sebagai default. ## Arsitektur Flink 2.3 untuk Pemrosesan Stream Flink berjalan pada arsitektur terdistribusi dengan JobManager yang mengkoordinasikan pekerjaan di beberapa TaskManager. Setiap TaskManager menjalankan task slot yang mengeksekusi bagian dari operator paralel job. Pemisahan ini memungkinkan Flink untuk scaling secara horizontal sambil mempertahankan fault tolerance melalui checkpoint terdistribusi. Model dataflow di Flink merepresentasikan komputasi sebagai directed acyclic graph (DAG). Data mengalir dari source melalui transformasi ke sink, dengan setiap operator berpotensi berjalan pada beberapa instance paralel. ```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()) ); ``` Konfigurasi checkpoint di atas menyimpan snapshot state terdistribusi setiap 10 detik. Jika terjadi kegagalan, Flink memulihkan dari checkpoint terakhir yang selesai dan memutar ulang record dari Kafka. ## Semantik Event Time vs Processing Time Event time mengacu pada kapan event sebenarnya terjadi, yang tertanam dalam data itu sendiri. Processing time adalah kapan Flink memproses record. Perbedaan ini penting karena delay jaringan, pengiriman tidak berurutan, dan backlog pemrosesan membuat processing time tidak dapat diandalkan untuk operasi berbasis waktu. ```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()); } } ``` Strategi `forBoundedOutOfOrderness` memberi tahu Flink bahwa event mungkin tiba terlambat hingga 2 menit. Watermark maju ketika Flink menentukan bahwa tidak ada lagi event dengan timestamp sebelum watermark yang akan tiba. > **Pertanyaan Interview Umum** > > Apa yang terjadi pada late event di Flink? Secara default, event yang tiba setelah watermark melewati waktu akhir window akan dibuang. Konfigurasikan allowed lateness dengan `.allowedLateness(Time.minutes(10))` untuk memproses kedatangan terlambat, atau gunakan side output untuk menangkapnya untuk penanganan terpisah. ## Strategi Windowing untuk Analitik Real-Time Flink menyediakan empat tipe window: tumbling, sliding, session, dan global window. Masing-masing melayani kebutuhan analitis yang berbeda. ```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 ditutup setelah jeda aktivitas yang dapat dikonfigurasi. Pola ini berfungsi untuk analisis perilaku pengguna di mana durasi sesi bervariasi berdasarkan engagement. ## Manajemen State dan Checkpointing Flink memelihara operator state dan keyed state di seluruh pemrosesan. Keyed state mempartisi data berdasarkan key, memungkinkan pemrosesan paralel sambil menjaga record terkait tetap bersama. Operator state berlaku untuk seluruh instance operator. ```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); } } ``` Prosesor stateful ini melacak pola transaksi per akun. State bertahan di seluruh checkpoint, bertahan dari kegagalan tanpa kehilangan konteks deteksi fraud. ## Flink SQL dan Table API untuk Pemrosesan Stream Flink 2.3 memperluas kemampuan SQL-nya dengan [Materialized Tables](https://flink.apache.org/downloads/) untuk pemeliharaan view inkremental. Table API menyediakan antarmuka terpadu untuk pemrosesan batch dan 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); ``` Pendekatan SQL menyederhanakan pengembangan untuk analis yang familiar dengan SQL sementara Flink menangani kompleksitas pemrosesan stream di balik layar. ## Flink vs Spark Structured Streaming Kedua framework memproses data streaming, tetapi arsitekturnya berbeda secara fundamental. Flink memproses record secara individual dengan streaming sejati, sementara Spark memproses micro-batch. Untuk [perbandingan Apache Spark](/blog/data-engineering/apache-spark-4-new-features-structured-streaming-interview), tradeoff latensi dan konsistensi penting dalam produksi. | Aspek | Flink | Spark Structured Streaming | |--------|-------|---------------------------| | Model Pemrosesan | Streaming sejati | Micro-batch | | Latensi | Milidetik | Detik (interval batch) | | State Backend | RocksDB, HashMaps | In-memory, HDFS | | Exactly-Once | Native dengan checkpoint | Memerlukan sink idempotent | | Event Time | Dukungan first-class | Didukung sejak 2.1 | | Dukungan SQL | Full streaming SQL | Windowing terbatas | > **Wawasan Interview** > > Ketika ditanya tentang Flink vs Spark untuk streaming, fokus pada kesesuaian use case. Flink unggul dalam pemrosesan event latensi rendah dan pola event kompleks. Spark Streaming cocok untuk organisasi yang sudah menjalankan Spark untuk batch yang membutuhkan pemrosesan batch-stream terpadu. ## Pertanyaan Interview Flink Umum dan Jawabannya **Bagaimana Flink mencapai semantik exactly-once?** Flink menggabungkan checkpointing dengan two-phase commit untuk sink yang mendukung transaksi. Selama checkpoint, Flink menyimpan snapshot operator state dan mencatat source offset. Untuk Kafka sink, Flink melakukan pre-commit record ke Kafka, menyelesaikan checkpoint, lalu commit transaksi. Jika kegagalan terjadi sebelum checkpoint selesai, record yang tidak di-commit dibuang dan pemrosesan dilanjutkan dari checkpoint terakhir. **Jelaskan propagasi watermark dalam topologi multi-source.** Ketika job membaca dari beberapa partisi atau source, masing-masing menghasilkan watermarknya sendiri berdasarkan event yang masuk. Watermark operator downstream sama dengan watermark minimum di semua channel input. Ini memastikan tidak ada window yang ditutup terlalu dini karena satu partisi cepat maju lebih dahulu dari yang lebih lambat. Konfigurasikan `withIdleness()` untuk memajukan watermark ketika beberapa partisi berhenti mengirim data. **Apa yang menyebabkan backpressure di Flink dan bagaimana mendiagnosisnya?** Backpressure terjadi ketika operator downstream tidak dapat mengikuti laju data upstream. Flink Web UI menampilkan status backpressure per operator. Penyebab umum meliputi: - Panggilan sistem eksternal yang lambat (query database, panggilan API) - Komputasi mahal dalam fungsi map/process - Paralelisme tidak mencukupi untuk volume data - Operasi state besar yang memblokir pemrosesan Atasi dengan meningkatkan paralelisme, mengoptimalkan operasi lambat, atau menggunakan async I/O untuk panggilan eksternal. **Bagaimana savepoint berbeda dari checkpoint?** Checkpoint bersifat otomatis, inkremental, dan dioptimalkan untuk recovery dari kegagalan. Flink mengelola lifecycle-nya, menghapus yang lama secara otomatis. Savepoint dipicu oleh pengguna, snapshot lengkap yang dimaksudkan untuk tugas operasional: deploy kode baru, rescaling job, atau migrasi antar cluster. Savepoint bertahan sampai dihapus secara eksplisit dan mendukung evolusi schema. ## Deploy Flink di Kubernetes [Flink Kubernetes Operator 1.15](https://flink.apache.org/2026/05/26/apache-flink-kubernetes-operator-1.15.0-release-announcement/) menyederhanakan deployment dengan FlinkDeployment custom resource. Operator ini menangani lifecycle job, upgrade, dan 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 ``` Pengaturan `upgradeMode: savepoint` memastikan operator mengambil savepoint sebelum upgrade, mempertahankan state di seluruh deployment. ## Mengoptimalkan Aplikasi Flink untuk Produksi Deployment produksi memerlukan perhatian pada konfigurasi paralelisme, memori, dan state backend. Lihat [pola ETL dan data pipeline](/technologies/data-engineering/interview-questions/etl-elt-patterns) untuk pertimbangan integrasi. ```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 inkremental mengurangi ukuran checkpoint dengan hanya menulis state yang berubah sejak checkpoint terakhir. Optimisasi ini menjadi kritis ketika mengelola state keyed berukuran gigabyte. ## Poin Utama untuk Pemrosesan Stream Apache Flink - Flink 2.3 memproses event secara individual dengan latensi milidetik, tidak seperti sistem micro-batch - Semantik event time dengan watermark menangani data tidak berurutan dengan benar, menjawab pertanyaan interview umum tentang late event - Keyed state mempartisi data untuk pemrosesan paralel sambil menjaga record terkait tetap bersama - Checkpoint memberikan jaminan exactly-once melalui snapshot terdistribusi dan two-phase commit - Kubernetes Operator mengotomatisasi deployment, scaling, dan upgrade dengan preservasi state berbasis savepoint - Pilih Flink daripada Spark Streaming ketika latensi sub-detik atau pola pemrosesan event kompleks penting - Konfigurasikan RocksDB dengan checkpoint inkremental untuk workload produksi dengan state besar --- Source: SharpSkill (https://sharpskill.dev), tech interview preparation for your real stack. HTML version of this page: https://sharpskill.dev/id/blog/data-engineering/apache-flink-stream-processing-event-time-interview-2026