데이터 엔지니어링

Flink Checkpoint/Savepoint 실패: 데이터 유실 방지를 위한 5가지 핵심 진단 및 복구 전략

강코의 코딩 일기 2026. 7. 24. 21:23
반응형

Flink Checkpoint/Savepoint 실패로 인한 데이터 유실을 막고 싶으신가요? 초보자도 쉽게 따라 할 수 있는 진단 방법과 복구 전략으로 안정적인 스트림 처리를 유지하세요.

안녕하세요! 실시간 데이터 처리에 관심 있는 여러분, Flink(플링크)라는 이름을 한 번쯤 들어보셨을 거예요. Flink는 엄청나게 많은 데이터를 빠르게 처리하는 데 특화된 도구인데요. 그런데 이 똑똑한 Flink도 가끔 말썽을 부릴 때가 있답니다. 특히 Checkpoint(체크포인트)Savepoint(세이브포인트)가 실패하면 정말 골치 아파지죠. 왜냐하면 이게 바로 데이터 유실로 이어질 수 있거든요!

상상해보세요. 실시간으로 중요한 데이터를 처리하고 있는데, 갑자기 시스템이 멈추거나 오류가 나는 거예요. 이때 Checkpoint나 Savepoint가 제대로 동작하지 않으면, 어디까지 처리했는지 알 수 없어서 처음부터 다시 시작해야 할 수도 있고, 심지어 일부 데이터는 영영 사라질 수도 있습니다. 생각만 해도 아찔하죠?

그래서 오늘은 Flink를 이제 막 배우기 시작한 분들을 위해, Checkpoint/Savepoint 실패가 왜 일어나는지, 어떻게 진단하고, 데이터 유실 없이 복구할 수 있는지 쉽고 친절하게 설명해 드릴게요. 이 글만 읽어도 Flink 작업의 안정성을 한층 높일 수 있을 거예요!

Flink Checkpoint/Savepoint, 왜 그렇게 중요한가요?

먼저 Checkpoint와 Savepoint가 정확히 무엇이고, 왜 우리에게 그렇게 중요한지부터 알아볼까요? Flink스트림 처리 프레임워크라고 했잖아요? 스트림 처리란 물이 흐르듯이 실시간으로 끊임없이 들어오는 데이터를 처리하는 것을 말해요. 그런데 이렇게 연속적인 작업은 중간에 문제가 생겼을 때 어디서부터 다시 시작해야 할지 애매할 수 있겠죠?

이때 CheckpointSavepoint가 구원투수처럼 등장합니다. 쉽게 말해, 이 둘은 현재 Flink 작업의 상태를 저장해두는 스냅샷이에요. 마치 게임을 하다가 중요한 순간에 '저장하기'를 눌러두는 것과 비슷하죠.

Checkpoint와 Savepoint, 무엇이 다를까요?

둘 다 작업 상태를 저장하는 건 맞지만, 목적이 조금 달라요. 아래 표로 비교해 보면 이해하기 쉬울 거예요.

구분 Checkpoint (체크포인트) Savepoint (세이브포인트)
목적 시스템 장애 발생 시 자동 복구를 위해 주기적으로 생성 사용자가 수동으로 트리거하여 작업 변경(업데이트, 재배포) 시 사용
생성 주체 Flink 시스템이 자동으로 설정된 주기에 따라 생성 사용자가 명시적으로 명령을 내려 생성
사용 용도 예상치 못한 실패로부터 데이터 유실 없이 자동으로 복구 업그레이드, 코드 변경 후 중단 없는 재시작 (버전 관리)
관리 일반적으로 Flink가 오래된 것을 삭제하며 관리 사용자가 직접 관리하고 삭제해야 함

이 둘이 중요한 이유는 바로 장애 내성(Fault Tolerance)데이터 일관성(Data Consistency)을 보장해주기 때문이에요. 만약 Flink 작업 중 문제가 생기더라도, 가장 최근의 Checkpoint나 지정된 Savepoint로 돌아가서 데이터 유실 없이 작업을 이어서 처리할 수 있게 되는 거죠. 얼마나 든든한 보험 같은 존재인가요!

우리 Flink 작업, Checkpoint/Savepoint는 왜 자꾸 실패할까요?

이렇게 중요한 Checkpoint/Savepoint인데, 왜 실패할까요? 원인은 정말 다양할 수 있지만, 초보자들이 자주 겪는 몇 가지 대표적인 이유들을 살펴볼게요.

1. 리소스 부족 (메모리, 디스크, 네트워크)

Flink 작업은 데이터를 처리하면서 상태를 저장하는데, 이때 메모리(RAM), 디스크(저장 공간), 네트워크 대역폭 같은 리소스가 많이 필요해요. 만약 이 리소스들이 충분하지 않으면 Checkpoint/Savepoint가 제때 완료되지 못하고 실패할 수 있습니다. 예를 들어:

  • 메모리 부족: Flink 작업자가 너무 많은 상태를 메모리에 들고 있는데, Checkpoint를 만들려고 하면 메모리가 부족해서 터져버리는 경우가 있어요.
  • 디스크 I/O 병목: Checkpoint나 Savepoint는 작업 상태를 디스크(주로 HDFS나 S3 같은 분산 스토리지)에 저장해야 하는데, 이때 디스크 쓰기 속도가 너무 느리거나 다른 작업으로 인해 과부하가 걸리면 실패할 수 있습니다.
  • 네트워크 문제: Flink는 분산 시스템이라 작업자들끼리, 그리고 저장소와 데이터를 주고받아요. 네트워크가 불안정하거나 대역폭이 부족하면 데이터 전송이 지연되거나 실패해서 Checkpoint가 타임아웃될 수 있습니다.

2. 외부 시스템 문제

Flink는 보통 Kafka(카프카) 같은 메시지 큐에서 데이터를 받아와서 HDFSS3 같은 분산 스토리지에 결과를 저장하거나, Checkpoint/Savepoint를 저장하기도 해요. 그런데 이 외부 시스템에 문제가 생기면 Flink 작업도 영향을 받아서 Checkpoint/Savepoint가 실패할 수 있습니다.

  • 스토리지 시스템 장애: Checkpoint 데이터를 저장하는 HDFS나 S3 버킷에 접근할 수 없거나 장애가 발생하면 당연히 Checkpoint가 실패하겠죠.
  • 접근 권한 문제: Flink가 저장소에 데이터를 쓸 권한이 없어서 실패하는 경우도 의외로 많답니다.

3. 사용자 코드 오류

가장 흔한 원인 중 하나인데요, 여러분이 작성한 Flink 코드 자체에 문제가 있을 때입니다.

  • 예외 발생: 특정 데이터가 들어왔을 때 처리 로직에서 예상치 못한 NullPointerException 같은 예외가 발생하면 작업이 중단되고, 이로 인해 Checkpoint가 실패할 수 있습니다.
  • 직렬화 문제: Flink는 작업 상태를 저장할 때 객체를 바이트 형태로 바꾸는 직렬화(Serialization) 과정을 거치는데요, 만약 상태로 저장하려는 객체가 직렬화 불가능한 객체라면 Checkpoint가 실패하게 됩니다.

4. Flink 설정 문제

Flink의 Checkpoint 관련 설정이 적절하지 않을 때도 실패할 수 있어요.

  • Checkpoint 타임아웃: Checkpoint가 완료되어야 하는 시간 제한(timeout)이 너무 짧게 설정되어 있으면, 리소스가 약간만 부족해도 타임아웃으로 실패할 수 있습니다.
  • 비동기 Checkpoint 비활성화: 대용량 상태를 처리할 때는 Checkpoint가 메인 작업에 영향을 주지 않도록 비동기(Asynchronous)로 설정하는 것이 좋은데요, 이 설정이 제대로 안 되어 있으면 Checkpoint 중 작업이 멈춰 성능 저하나 실패로 이어질 수 있습니다.

실패 원인을 어떻게 찾아낼 수 있을까요? 진단 도구 활용법!

문제가 생겼다면, 어디가 문제인지 정확히 아는 것이 중요하겠죠? Flink는 친절하게도 문제를 진단할 수 있는 여러 도구를 제공합니다.

1. Flink Web UI 활용하기

가장 먼저 확인해야 할 곳은 바로 Flink Web UI입니다. Flink 작업을 실행하면 웹 브라우저를 통해 접속할 수 있는 대시보드인데요, 여기에 아주 유용한 정보들이 많아요.

  • Job Graph: 작업의 흐름을 시각적으로 보여줍니다. 여기서 어떤 연산자(Operator)에서 문제가 발생했는지 대략적인 위치를 파악할 수 있어요.
  • Checkpoints 탭: 여기에 들어가면 모든 Checkpoint 시도와 그 결과(성공/실패)가 기록되어 있습니다. 실패한 Checkpoint를 클릭하면 실패 원인(Failure Cause)스택 트레이스(Stack Trace)를 확인할 수 있어요. 이 스택 트레이스에 보통 어떤 파일의 몇 번째 줄에서 문제가 발생했는지 나와있답니다.
  • Task Managers 탭: 각 작업자의 리소스 사용량(CPU, 메모리)과 로그를 볼 수 있습니다. 메모리 부족이나 GC(Garbage Collection) 문제가 의심될 때 유용해요.

2. Flink 로그 분석하기

Web UI에서 대략적인 원인을 파악했다면, 더 자세한 내용은 로그(Log) 파일을 확인해야 합니다. Flink 작업이 실행되는 서버나 컨테이너의 로그 디렉토리에 가면 .out이나 .log 확장자를 가진 파일들을 찾을 수 있을 거예요.

특히 TaskManager 로그에서 'Checkpoint failed', 'Exception', 'OutOfMemoryError' 같은 키워드를 검색해보세요. 실패 원인과 관련된 자세한 에러 메시지를 찾을 수 있을 겁니다. 예를 들어, 아래와 같은 메시지를 발견할 수도 있죠.


2023-10-27 10:30:05,123 WARN  org.apache.flink.runtime.checkpoint.CheckpointCoordinator     - Checkpoint 123 for job 456 failed.
org.apache.flink.runtime.checkpoint.CheckpointException: Could not materialize checkpoint 123 for job 456.
    at org.apache.flink.runtime.checkpoint.CheckpointCoordinator.failCheckpoint(CheckpointCoordinator.java:1122)
    at org.apache.flink.runtime.checkpoint.CheckpointCoordinator.lambda$triggerCheckpoint$12(CheckpointCoordinator.java:999)
    ...
Caused by: java.io.IOException: Could not write to file 's3://your-bucket/flink-checkpoints/chk-123/shared/...'
    at org.apache.flink.runtime.fs.hdfs.HadoopFileSystem.create(HadoopFileSystem.java:319)
    at org.apache.flink.core.fs.Path.create(Path.java:189)
    ...
Caused by: com.amazonaws.services.s3.model.AmazonS3Exception: Access Denied (Service: Amazon S3; Status Code: 403; Error Code: AccessDenied; Request ID: ...

위 로그를 보면 "Access Denied" 메시지와 "Status Code: 403"이 보이시죠? 이는 Flink가 S3 스토리지에 Checkpoint를 저장하려 했지만 접근 권한이 없어서 실패했다는 것을 명확히 알려줍니다. 이렇게 로그는 문제 해결의 가장 중요한 실마리를 제공해줘요.

3. 지표(Metrics) 모니터링

Flink는 Grafana나 Prometheus 같은 모니터링 시스템과 연동하여 다양한 지표를 확인할 수 있습니다. 특히 Checkpoint 관련 지표는 실패 원인을 미리 감지하거나 사후 분석하는 데 아주 유용해요.

  • numberOfFailedCheckpoints: 실패한 Checkpoint 수를 보여줍니다. 이 수치가 계속 올라간다면 심각한 문제겠죠.
  • lastCheckpointSize, lastCheckpointDuration: 마지막 Checkpoint의 크기와 걸린 시간을 보여줍니다. 이 값들이 갑자기 커지거나 길어진다면 상태 크기가 비정상적으로 커졌거나 리소스 문제가 발생했을 가능성이 높아요.
  • 백프레셔(Backpressure): Flink UI에서 Backpressure 탭을 확인해보세요. 만약 어느 오퍼레이터에서 백프레셔가 발생하고 있다면, 그 오퍼레이터가 데이터를 충분히 빠르게 처리하지 못하고 있다는 뜻이고, 이는 Checkpoint 실패로 이어질 수 있습니다.

데이터 유실 없이 안전하게 복구하고, 재발을 막는 전략은?

문제를 진단했다면 이제 해결하고 재발을 막는 방법을 알아봐야겠죠?

1. 실패한 작업 복구하기

가장 중요한 건 데이터 유실 없이 작업을 복구하는 거예요. Flink는 기본적으로 가장 최근에 성공한 Checkpoint에서 자동으로 작업을 재시작하려고 시도합니다. 만약 이 자동 복구가 실패하거나, 특정 시점으로 돌아가고 싶다면 Savepoint를 활용할 수 있어요.

만약 작업이 계속 실패한다면, 실패 원인을 해결한 후 Flink CLI (Command Line Interface)를 사용해 특정 Checkpoint나 Savepoint에서 작업을 재시작할 수 있습니다.


# 특정 Savepoint에서 작업 재시작
flink run -s s3://your-bucket/flink-savepoints/savepoint-000001 -d your-flink-job.jar

여기서 -s 옵션 뒤에 복구하고 싶은 Savepoint 경로를 넣어주면 됩니다. Checkpoint의 경우, Flink는 기본적으로 실패 시 가장 최근 성공한 Checkpoint를 사용하려고 하지만, 특정 Checkpoint를 지정하려면 해당 Checkpoint의 메타데이터 경로를 알아야 해요. 보통 Savepoint가 더 유연하게 특정 시점으로 돌아가는 데 사용됩니다.

2. 데이터 유실 방지를 위한 예방책

a) 충분한 리소스 확보 및 모니터링

가장 기본적이면서도 중요한 부분입니다. Flink 작업에 필요한 메모리, CPU, 디스크, 네트워크 대역폭을 충분히 확보해야 해요. 그리고 앞서 설명한 지표 모니터링을 통해 리소스 사용량을 꾸준히 감시하고, 임계치를 넘으면 알림을 받을 수 있도록 설정해두세요. 백프레셔가 발생하거나, GC 시간이 비정상적으로 길어지는지 주시하는 것이 중요합니다.

b) 견고한 코드 작성 및 직렬화 고려

여러분 코드에서 예외가 발생하지 않도록 예외 처리(Error Handling)를 꼼꼼히 해주세요. 예를 들어, null 값이나 예상치 못한 데이터 형식에 대비해야 합니다. 또한, 상태(State)로 저장되는 모든 객체는 직렬화 가능(Serializable)해야 한다는 것을 항상 기억하고, 필요하다면 KryoSerializer 같은 커스텀 직렬화기를 사용하는 것도 고려해볼 수 있습니다.

c) Flink Checkpoint 설정 튜닝

Flink의 Checkpoint 관련 설정을 적절히 튜닝하는 것이 중요해요. 몇 가지 중요한 설정들을 살펴볼까요?

  • checkpointing.interval: Checkpoint를 얼마나 자주 생성할지 설정합니다. (예: 60초마다) 너무 길면 복구 시 잃을 수 있는 데이터가 많아지고, 너무 짧으면 시스템 부하가 커질 수 있어요.
  • checkpointing.timeout: Checkpoint가 완료되어야 하는 최대 시간을 설정합니다. (예: 10분) 작업 상태가 크거나 네트워크가 느리다면 이 값을 늘려줄 필요가 있습니다.
  • state.backend: 상태를 어디에 저장할지 결정합니다. RockDB, FileSystem, Memory 백엔드 등이 있는데, 대용량 상태를 다룬다면 RocksDB State Backend를 사용하는 것이 일반적입니다. 이는 상태를 디스크에 저장하여 메모리 부담을 줄여줍니다.
  • state.checkpoints.dir: Checkpoint 데이터를 저장할 경로를 지정합니다. S3나 HDFS처럼 영속적인 분산 스토리지를 사용해야 합니다.
  • checkpointing.mode: EXACTLY_ONCE 또는 AT_LEAST_ONCE 중 하나를 선택합니다. EXACTLY_ONCE는 데이터 중복 없이 정확히 한 번만 처리됨을 보장하지만, 오버헤드가 더 큽니다.

이러한 설정들은 여러분의 Flink 작업 특성과 시스템 환경에 맞춰 신중하게 결정해야 합니다. 초기에는 기본값을 사용하더라도, 문제가 발생하면 로그와 모니터링 지표를 보면서 조금씩 조절해나가는 것이 좋아요.

코드 예시 (Python Flink API - PyFlink):


from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.checkpointing_mode import CheckpointingMode
from pyflink.common.restart_strategy import RestartStrategies
from pyflink.common.time import Time

env = StreamExecutionEnvironment.get_execution_environment()

# Checkpoint 활성화 및 설정
env.enable_checkpointing(10000, CheckpointingMode.EXACTLY_ONCE) # 10초마다 EXACTLY_ONCE 모드로 Checkpoint
env.get_checkpoint_config().set_checkpoint_timeout(60000) # Checkpoint 타임아웃 1분
env.get_checkpoint_config().set_min_pause_between_checkpoints(5000) # Checkpoint 간 최소 간격 5초
env.get_checkpoint_config().set_tolerable_checkpoint_failure_number(3) # 3번까지 Checkpoint 실패 허용

# Checkpoint 저장 경로 설정 (S3 예시)
# 반드시 Flink 클러스터가 S3에 접근할 권한이 있어야 합니다.
env.get_checkpoint_config().set_checkpoint_storage("s3://your-bucket/flink/checkpoints")

# 재시작 전략 설정 (Fixed Delay)
# 작업 실패 시 3번 재시도, 각 시도 사이에 10초 대기
env.set_restart_strategy(RestartStrategies.fixed_delay_restart(
    restart_attempts=3,
    delay_between_restarts=Time.seconds(10)
))

# ... 데이터 스트림 및 연산 로직 ...

env.execute("my_flink_job")

위 코드처럼 Python Flink API를 사용해도 Checkpoint 관련 설정을 쉽게 할 수 있어요. JVM 기반 Flink 애플리케이션에서는 StreamExecutionEnvironment 객체를 통해 비슷한 메서드들을 호출하여 설정합니다.

마무리하며

오늘은 Flink Checkpoint/Savepoint 실패라는 조금 무섭게 들릴 수 있는 주제를 다뤄봤는데요, 어떠셨나요? Flink를 처음 접하는 분들에게는 다소 복잡하게 느껴질 수도 있지만, Checkpoint와 Savepoint가 왜 중요하고, 실패했을 때 어떻게 대처해야 하는지 기본적인 내용을 이해하는 것이 정말 중요하답니다. 이들이 바로 여러분의 소중한 데이터를 유실로부터 보호해주는 핵심 방어막이기 때문이죠.

핵심은 다음과 같아요. Flink Web UI와 로그를 통해 문제의 원인을 정확히 파악하고, 충분한 리소스 확보와 견고한 코드 작성, 그리고 적절한 Checkpoint 설정을 통해 실패를 예방하는 것이죠. 만약 실패하더라도 당황하지 않고 Savepoint를 활용하여 복구할 수 있다는 점도 기억해 주세요.

꾸준히 Flink를 사용하고 익숙해지다 보면, 이런 문제들은 더 이상 두렵지 않을 거예요! 혹시 이 글을 읽으시면서 궁금한 점이 생기거나, 여러분만의 Flink 트러블슈팅 경험이 있다면 댓글로 자유롭게 나눠주세요. 함께 배우고 성장해나가요!

📌 함께 읽으면 좋은 글

  • [보안] 보안 테스트 리포트의 치명적 함정: False Positive/Negative에 갇혀 본질을 놓치지 마세요
  • [보안] Stateful Firewall, 그 깊이를 파헤쳐보니: 세션 트래킹부터 NAT까지, 제가 직접 경험한 내부 동작 원리
  • [테스트 QA] 높은 코드 커버리지에도 버그가 쏟아진다면? 찐 개발자를 위한 테스트 커버리지 실전 개선 가이드

이 글이 도움이 되셨다면 공감(♥)댓글로 응원해 주세요!
궁금한 점이나 다루었으면 하는 주제가 있다면 댓글로 남겨주세요.

반응형