Apache Flink 2026 완벽 가이드: 스트림 처리, 이벤트 시간, 면접 질문

Apache Flink 2.3의 스트림 처리 아키텍처, 이벤트 시간 시맨틱스, 윈도우 처리, 상태 관리를 상세히 설명합니다. 데이터 엔지니어링 면접에서 자주 나오는 질문과 답변도 포함되어 있습니다.

Apache Flink 스트림 처리 아키텍처 다이어그램

Apache Flink 2.3은 Exactly-Once 시맨틱스와 밀리초 단위 지연 시간으로 대규모 스트림 처리를 수행한다. 배치 지향 시스템과 달리 Flink는 데이터가 도착하는 즉시 연속적으로 처리하므로, 실시간 분석, 사기 탐지, 이벤트 기반 아키텍처에 최적의 프레임워크이다.

면접 핵심 포인트

Flink가 Spark Streaming과 다른 점은 진정한 스트림 처리에 있다. Flink는 이벤트 시간 시맨틱스로 이벤트를 하나씩 처리하지만, Spark Streaming은 처리 시간을 기본으로 마이크로 배치를 처리한다.

Flink는 분산 아키텍처에서 실행되며, JobManager가 여러 TaskManager 간의 작업을 조율한다. 각 TaskManager는 작업의 병렬 연산자 일부를 실행하는 태스크 슬롯을 보유한다. 이러한 분리를 통해 분산 체크포인트로 내결함성을 유지하면서 수평 확장이 가능하다.

Flink의 데이터플로우 모델은 계산을 방향성 비순환 그래프(DAG)로 표현한다. 데이터는 소스에서 변환을 거쳐 싱크로 흐르며, 각 연산자는 여러 병렬 인스턴스에서 실행될 수 있다.

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())
    );

위의 체크포인트 설정은 10초마다 분산 상태의 스냅샷을 저장한다. 장애가 발생하면 Flink는 마지막으로 완료된 체크포인트에서 복원하고 Kafka에서 레코드를 재생한다.

이벤트 시간과 처리 시간 시맨틱스

이벤트 시간은 이벤트가 실제로 발생한 시점으로, 데이터 자체에 포함되어 있다. 처리 시간은 Flink가 레코드를 처리하는 시점이다. 네트워크 지연, 순서 뒤바뀜, 처리 백로그로 인해 처리 시간은 시간 기반 연산에서 신뢰할 수 없으므로 이 구분이 중요하다.

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 전략은 이벤트가 최대 2분까지 늦게 도착할 수 있음을 Flink에 알린다. 워터마크는 해당 워터마크 이전 타임스탬프를 가진 이벤트가 더 이상 도착하지 않을 것으로 Flink가 판단할 때 진행된다.

자주 나오는 면접 질문

Flink에서 지연 이벤트는 어떻게 되는가? 기본적으로 워터마크가 윈도우 종료 시간을 지난 후 도착한 이벤트는 삭제된다. 지연 도착을 처리하려면 .allowedLateness(Time.minutes(10))로 허용 지연을 설정하거나, 사이드 출력을 사용하여 별도로 처리한다.

실시간 분석을 위한 윈도우 전략

Flink는 네 가지 윈도우 유형을 제공한다: 텀블링, 슬라이딩, 세션, 글로벌 윈도우. 각각 다른 분석 요구사항에 적합하다.

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());
    }
}

세션 윈도우는 설정 가능한 비활성 기간 후에 닫힌다. 이 패턴은 참여도에 따라 세션 길이가 달라지는 사용자 행동 분석에 적합하다.

Data Engineering 면접 준비가 되셨나요?

인터랙티브 시뮬레이터, flashcards, 기술 테스트로 연습하세요.

상태 관리와 체크포인팅

Flink는 처리 전반에 걸쳐 연산자 상태와 키 상태를 유지한다. 키 상태는 키별로 데이터를 분할하여 관련 레코드를 함께 유지하면서 병렬 처리를 가능하게 한다. 연산자 상태는 전체 연산자 인스턴스에 적용된다.

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);
    }
}

이 상태 저장 프로세서는 계정별 트랜잭션 패턴을 추적한다. 상태는 체크포인트 간에 지속되어 사기 탐지 컨텍스트를 잃지 않고 장애를 극복한다.

Flink 2.3은 증분 뷰 유지를 위한 Materialized Tables로 SQL 기능을 확장한다. Table API는 배치 처리와 스트림 처리를 위한 통합 인터페이스를 제공한다.

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가 내부적으로 스트림 처리의 복잡성을 처리한다.

두 프레임워크 모두 스트리밍 데이터를 처리하지만, 아키텍처가 근본적으로 다르다. Flink는 진정한 스트리밍으로 레코드를 개별 처리하고, Spark는 마이크로 배치를 처리한다. Apache Spark 비교에서 프로덕션 환경의 지연 시간과 일관성 트레이드오프가 중요하다.

측면FlinkSpark Structured Streaming
처리 모델진정한 스트리밍마이크로 배치
지연 시간밀리초초 (배치 간격)
상태 백엔드RocksDB, HashMap인메모리, HDFS
Exactly-Once체크포인트로 네이티브멱등 싱크 필요
이벤트 시간일급 지원2.1부터 지원
SQL 지원완전한 스트리밍 SQL제한적 윈도우
면접 팁

스트리밍에서 Flink vs Spark에 대해 질문받으면 사용 사례 적합성에 집중한다. Flink는 저지연 이벤트 처리와 복잡한 이벤트 패턴에 탁월하다. Spark Streaming은 이미 배치 처리에 Spark를 사용하고 배치-스트림 통합 처리가 필요한 조직에 적합하다.

Flink는 어떻게 Exactly-Once 시맨틱스를 달성하는가?

Flink는 트랜잭션을 지원하는 싱크에 대해 체크포인팅과 2단계 커밋을 결합한다. 체크포인트 동안 Flink는 연산자 상태의 스냅샷을 찍고 소스 오프셋을 기록한다. Kafka 싱크의 경우, Flink는 레코드를 Kafka에 사전 커밋하고, 체크포인트를 완료한 다음 트랜잭션을 커밋한다. 체크포인트 완료 전에 장애가 발생하면 커밋되지 않은 레코드는 삭제되고 마지막 체크포인트에서 처리가 재개된다.

다중 소스 토폴로지에서 워터마크 전파를 설명하라

작업이 여러 파티션이나 소스에서 읽을 때, 각각 수신 이벤트를 기반으로 자체 워터마크를 생성한다. 다운스트림 연산자의 워터마크는 모든 입력 채널의 워터마크 중 최소값과 같다. 이렇게 하면 빠른 파티션이 느린 파티션보다 앞서 진행해도 윈도우가 조기에 닫히지 않는다. 일부 파티션이 데이터 전송을 중단할 때 워터마크를 진행시키려면 withIdleness()를 구성한다.

Flink에서 백프레셔의 원인은 무엇이며 어떻게 진단하는가?

백프레셔는 다운스트림 연산자가 업스트림 데이터 속도를 따라잡지 못할 때 발생한다. Flink Web UI는 연산자별 백프레셔 상태를 표시한다. 일반적인 원인:

  • 느린 외부 시스템 호출 (데이터베이스 쿼리, API 호출)
  • map/process 함수의 비싼 계산
  • 데이터 볼륨에 비해 불충분한 병렬성
  • 큰 상태 연산이 처리를 차단

병렬성 증가, 느린 연산 최적화, 외부 호출에 비동기 I/O 사용으로 해결한다.

세이브포인트와 체크포인트의 차이점은?

체크포인트는 자동, 증분적이며 장애 복구에 최적화되어 있다. Flink가 수명 주기를 관리하고 오래된 것은 자동으로 삭제된다. 세이브포인트는 사용자가 트리거하는 완전한 스냅샷으로, 새 코드 배포, 작업 재조정, 클러스터 간 마이그레이션과 같은 운영 작업을 위한 것이다. 세이브포인트는 명시적으로 삭제될 때까지 지속되며 스키마 진화를 지원한다.

Flink Kubernetes Operator 1.15는 FlinkDeployment 커스텀 리소스로 배포를 단순화한다. 작업 수명 주기, 업그레이드, 스케일링을 처리한다.

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 설정으로 오퍼레이터는 업그레이드 전에 세이브포인트를 취하여 배포 간 상태를 보존한다.

프로덕션 배포는 병렬성, 메모리, 상태 백엔드 구성에 주의가 필요하다. 통합 고려사항은 ETL 및 데이터 파이프라인 패턴을 참조한다.

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
        ));
    }
}

증분 체크포인트는 마지막 체크포인트 이후 변경된 상태만 기록하여 체크포인트 크기를 줄인다. 이 최적화는 기가바이트 단위의 키 상태를 관리할 때 중요해진다.

연습을 시작하세요!

면접 시뮬레이터와 기술 테스트로 지식을 테스트하세요.

  • Flink 2.3은 밀리초 지연 시간으로 이벤트를 개별 처리하며, 마이크로 배치 시스템과 다르다
  • 워터마크가 포함된 이벤트 시간 시맨틱스는 순서가 뒤바뀐 데이터를 올바르게 처리하여 지연 이벤트에 관한 일반적인 면접 질문에 답한다
  • 키 상태는 관련 레코드를 함께 유지하면서 병렬 처리를 위해 데이터를 분할한다
  • 체크포인트는 분산 스냅샷과 2단계 커밋을 통해 Exactly-Once 보장을 제공한다
  • Kubernetes Operator는 세이브포인트 기반 상태 보존으로 배포, 스케일링, 업그레이드를 자동화한다
  • 밀리초 미만 지연 시간이나 복잡한 이벤트 처리 패턴이 중요할 때 Spark Streaming보다 Flink를 선택한다
  • 대용량 상태의 프로덕션 워크로드에는 증분 체크포인트가 포함된 RocksDB를 구성한다
오늘의 챌린지

Data Engineering 코드의 버그를 찾을 수 있나요

실제 코드 한 조각, 숨은 버그 하나, 하루 한 번. 계정 없이 바로 도전할 수 있습니다.

Anthony Fillion-Maillet

작성자

Anthony Fillion-Maillet

SharpSkill 창업자

10년 이상 풀스택 개발을 해왔습니다. SharpSkill을 운영하며 이곳에 게시되는 모든 내용에 책임을 집니다.

2026년 8월 28일 업데이트

태그

#Apache Flink
#스트림 처리
#이벤트 시간
#데이터 엔지니어링
#면접 준비

공유

관련 기사