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 HA | JobManager 가용성 |
| 모니터링 및 알림 | 체크포인트 소요 시간 추적 | 임계치 초과 알림 |
5.2. 주요 트러블슈팅
이슈 1: 체크포인트 타임아웃 빈발
# 해결: 타임아웃 시간 및 최소 간격 조정
setCheckpointTimeout(180000) # 타임아웃 연장
setMinPauseBetweenCheckpoints(2000) # 최소 대기 시간 증가
이슈 2: 관리 상태 크기 과다
// 해결: RocksDB 상태 백엔드 및 증분 체크포인트 적용
flinkEnv.setStateBackend(new RocksDBStateBackend("hdfs:///flink/rocks-state/", true));
이슈 3: 소스 시스템 오프셋 리셋 미지원
// 해결: 오프셋 리셋을 지원하는 소스로 교체 또는 버퍼링 큐 추가
kafkaSource.setStartFromGroupOffsets();
kafkaSource.setCommitOffsetsOnCheckpoints(true);