Apache Flink 작업 재시작 시 데이터 유실 원인 분석 및 해결 전략

1. 데이터 유실의 근본적 원인

Apache Flink 잡(Job)이 재시작될 때 데이터 유실이 발생하는 핵심적인 원인은 처리 과정에서 발생하는 상태(State) 불일치입니다. 장애 발생 시 메모리 상에 존재하던 처리 중인 데이터가 영속화되지 않았거나, 데이터 소스의 컨슘 오프셋(Consumption Offset)이 정상적으로 저장되지 않은 상태로 복구가 이루어지면 필연적으로 데이터 누락이 발생하게 됩니다.

2. 체크포인트(Checkpoint) 메커니즘 심층 분석

2.1. 체크포인트 최적화 설정

StreamExecutionEnvironment flinkEnv = StreamExecutionEnvironment.getExecutionEnvironment();

// 기본 체크포인트 활성화 (10초 주기, 정확히 한 번 처리)
flinkEnv.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);

// 고급 체크포인트 최적화 설정
CheckpointConfig cpConfig = flinkEnv.getCheckpointConfig();
cpConfig.setMinPauseBetweenCheckpoints(1000);         // 체크포인트 간 최소 대기 시간
cpConfig.setCheckpointTimeout(120000);                // 체크포인트 타임아웃
cpConfig.setMaxConcurrentCheckpoints(1);              // 동시 실행 체크포인트 제한
cpConfig.setTolerableCheckpointFailureNumber(5);      // 허용 가능한 체크포인트 실패 횟수
cpConfig.enableExternalizedCheckpoints(
    ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); // 잡 취소 시 체크포인트 유지

// 상태 백엔드 구성 (RocksDB)
flinkEnv.setStateBackend(new RocksDBStateBackend("hdfs:///flink/state-backends/", true));

3. End-to-End Exactly-Once 보장 체계

3.1. 2단계 커밋(2PC) Sink 구현 예시

// 정확히 한 번(JDBC) 싱크 구현 예시
JdbcSink.exactlyOnceSink(
    "INSERT INTO order_metrics (customer_id, event_type, total_amount, created_at) " +
    "VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE " +
    "total_amount = total_amount + VALUES(total_amount), created_at = VALUES(created_at)",
    (preparedStatement, metric) -> {
        preparedStatement.setString(1, metric.getCustomerId());
        preparedStatement.setString(2, metric.getEventType());
        preparedStatement.setInt(3, metric.getTotalAmount());
        preparedStatement.setTimestamp(4, new Timestamp(metric.getCreatedAt()));
    },
    JdbcExecutionOptions.builder()
        .withBatchSize(2000)
        .withBatchIntervalMs(500)
        .withMaxRetries(5)
        .build(),
    new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
        .withUrl("jdbc:postgresql://db-host:5432/warehouse")
        .withDriverName("org.postgresql.Driver")
        .withUsername("flink_app")
        .withPassword("secret_token")
        .build()
);

4. 장애 복구 전략 및 데이터 무결성 확보

4.1. 고가용성 및 재시작 설정

# flink-conf.yaml 고가용성(HA) 설정
high-availability: zookeeper
high-availability.zookeeper.quorum: zookeeper1:2181,zookeeper2:2181,zookeeper3:2181
high-availability.storageDir: hdfs:///flink/ha-storage/
high-availability.cluster-id: /production-cluster

# 상태 및 체크포인트 스토리지 경로
state.checkpoints.dir: hdfs:///flink/cp-storage
state.savepoints.dir: hdfs:///flink/sp-storage

# 재시작 전략 (지수 백오프)
restart-strategy: exponential-delay
restart-strategy.exponential-delay.initial-backoff: 5s
restart-strategy.exponential-delay.max-backoff: 5min
restart-strategy.exponential-delay.backoff-multiplier: 3.0
restart-strategy.exponential-delay.reset-backoff-threshold: 15min
restart-strategy.exponential-delay.jitter-factor: 0.2

5. 실무 운영을 위한 체크리스트 및 이슈 대응

5.1. 데이터 무결성 점검 항목

점검 항목권장 설정핵심 모니터링 지표
체크포인트 활성화enableCheckpointing(10000)체크포인트 성공률
상태 백엔드RocksDBStateBackend상태 크기 및 접근 지연 시간
소스 재생 가능성Kafka Offset 커밋 연동Consumer Lag
Sink 멱등성멱등 쓰기 또는 2PC 적용쓰기 실패율
고가용성 구성ZooKeeper HAJobManager 가용성
모니터링 및 알림체크포인트 소요 시간 추적임계치 초과 알림

5.2. 주요 트러블슈팅

이슈 1: 체크포인트 타임아웃 빈발

# 해결: 타임아웃 시간 및 최소 간격 조정
setCheckpointTimeout(180000)  # 타임아웃 연장
setMinPauseBetweenCheckpoints(2000)  # 최소 대기 시간 증가

이슈 2: 관리 상태 크기 과다

// 해결: RocksDB 상태 백엔드 및 증분 체크포인트 적용
flinkEnv.setStateBackend(new RocksDBStateBackend("hdfs:///flink/rocks-state/", true));

이슈 3: 소스 시스템 오프셋 리셋 미지원

// 해결: 오프셋 리셋을 지원하는 소스로 교체 또는 버퍼링 큐 추가
kafkaSource.setStartFromGroupOffsets();
kafkaSource.setCommitOffsetsOnCheckpoints(true);

태그: Apache Flink Checkpoint Exactly-Once StateBackend RocksDB

8월 28일 23:41에 게시됨