Savepoints

Savepoints

Savepoint이란 무엇인가?

Savepoint은 Flink의 checkpointing 메커니즘을 통해 생성된 스트리밍 작업의 실행 상태에 대한 일관된 이미지입니다. Savepoint을 사용해 Flink 작업을 중지-재개(stop-and-resume), 포크(fork), 또는 갱신(update)할 수 있습니다. Savepoint은 두 부분으로 구성됩니다: 안정적인 저장소(예: HDFS, S3 등)에 있는 (보통 큰) 이진 파일이 있는 디렉터리와 (비교적 작은) 메타데이터 파일. 안정적인 저장소의 파일은 작업의 실행 상태 이미지의 순수 데이터를 나타냅니다. Savepoint의 메타데이터 파일은 (주로) Savepoint의 일부인 안정적인 저장소의 모든 파일에 대한 포인터를 상대 경로 형태로 포함합니다.

프로그램과 Flink 버전 간 업그레이드를 허용하려면 연산자에 ID 할당에 관한 다음 섹션을 확인하는 것이 중요합니다.

savepoint을 제대로 사용하려면 checkpoints와 savepoint의 차이를 이해하는 것이 중요하며, 이는 checkpoints vs. savepoints에 설명되어 있습니다.

출처: 문서

본문

연산자 ID 할당

uid(String) 메서드로 연산자 ID를 지정하는 것이 매우 권장됩니다. 이 ID는 각 연산자의 상태를 범위 지정하는 데 사용됩니다.

DataStream<String> stream = env.
  // Stateful source (e.g. Kafka) with ID
  .addSource(new StatefulSource())
  .uid("source-id") // ID for the source operator
  .shuffle()
  // Stateful mapper with ID
  .map(new StatefulMapper())
  .uid("mapper-id") // ID for the mapper
  // Stateless printing sink
  .print(); // Auto-generated ID

ID를 수동으로 지정하지 않으면 자동으로 생성됩니다. 이 ID가 변경되지 않는 한 savepoint에서 자동으로 복원할 수 있습니다. 생성된 ID는 프로그램의 구조에 따라 달라지며 프로그램 변경에 민감합니다. 따라서 이 ID를 수동으로 할당하는 것이 매우 권장됩니다.

Savepoint 상태

savepoint은 각 상태 기반 연산자에 대해 Operator ID -> State 맵을 보유하는 것으로 생각할 수 있습니다:

Operator ID | State
------------+------------------------
source-id   | State of StatefulSource
mapper-id   | State of StatefulMapper

위 예시에서 print 싱크는 무상태이므로 savepoint 상태에 포함되지 않습니다. 기본적으로 우리는 savepoint의 각 항목을 새 프로그램에 다시 매핑하려고 시도합니다.

연산

savepoint을 트리거, [savepoint으로 작업 취소], [savepoint에서 재개], [savepoint 폐기]를 위해 커맨드라인 클라이언트를 사용할 수 있습니다.

webui를 사용해 savepoint에서 재개하는 것도 가능합니다.

Savepoint 트리거

savepoint을 트리거하면 데이터와 메타데이터가 저장되는 새 savepoint 디렉터리가 생성됩니다. 이 디렉터리의 위치는 기본 대상 디렉터리 구성 또는 트리거 명령으로 사용자 지정 대상 디렉터리를 지정하는 것(:targetDirectory 인수 참조)으로 제어할 수 있습니다.

주의: 대상 디렉터리는 JobManager와 TaskManager 모두가 접근할 수 있는 위치여야 합니다. 예: 분산 파일 시스템 또는 Object Store의 위치.

예를 들어 FsStateBackend 또는 RocksDBStateBackend를 사용하는 경우:

# Savepoint target directory
/savepoints/

# Savepoint directory
/savepoints/savepoint-:shortjobid-:savepointid/

# Savepoint file contains the checkpoint meta data
/savepoints/savepoint-:shortjobid-:savepointid/_metadata

# Savepoint state
/savepoints/savepoint-:shortjobid-:savepointid/...

Savepoint은 일반적으로 전체 savepoint 디렉터리를 다른 위치로 이동(또는 복사)하여 이동할 수 있으며, Flink는 이동된 savepoint에서 복원할 수 있습니다.

두 가지 예외가 있습니다:

  1. 엔트로피 주입(entropy injection)이 활성화된 경우: 이 경우 savepoint 디렉터리가 모든 savepoint 데이터 파일을 포함하지 않습니다. 주입된 경로 엔트로피가 파일을 많은 디렉터리에 분산시키기 때문입니다. 공통 savepoint 루트 디렉터리가 없으므로 savepoint은 절대 경로 참조를 포함하게 되어 디렉터리 이동을 방지합니다.
  2. 작업이 GenericWriteAhreadLog sink 같은 task-owned state를 포함하는 경우.

savepoint과 달리 checkpoint는 일반적으로 다른 위치로 이동할 수 없습니다. checkpoint에는 일부 절대 경로 참조가 포함될 수 있기 때문입니다.

statebackend: jobmanager를 사용하면 메타데이터 savepoint 상태가 _metadata 파일에 저장되므로 추가 데이터 파일이 없다는 것에 혼동하지 마세요.

Flink 1.15부터 중간 savepoint(stop-with-savepoint로 생성된 것 외의 savepoint)는 복구에 사용되지 않으며 어떤 부수 효과도 커밋하지 않습니다.

이는 특히 같은 checkpointing 타임라인에서 여러 작업을 실행할 때 고려해야 합니다. 그 솔루션에서 원래 작업이 (savepoint을 찍은 후) 실패하면, savepoint 이전의 checkpoint로 폴백할 수 있습니다. 그러나 savepoint에서 작업을 재개하면 savepoint 이전의 checkpoint로 폴백하여 실제로 발생하지 않았을 수 있는 트랜잭션을 커밋할 수 있습니다(비결정성 가정).

이러한 시나리오에서 안전하려면 sink의 uids를 변경하여 트랜잭션 sink의 상태를 버리는 것을 권장합니다.

같은 checkpointing 타임라인에서 단일 작업만 실행된다면 추가 단계가 필요하지 않아야 합니다. 즉, savepoint에서 새 작업을 실행하기 전에 원래 작업을 중지한다는 뜻입니다.

Savepoint 포맷

savepoint의 두 가지 이진 포맷 중 하나를 선택할 수 있습니다.

  • canonical 포맷 - 모든 상태 백엔드에 걸쳐 통합된 포맷으로, 한 상태 백엔드로 savepoint을 찍고 다른 백엔드로 복원할 수 있습니다. 이것은 가장 안정적인 포맷이며 이전 버전, 스키마, 수정 등과의 최대 호환성을 유지하는 것을 목표로 합니다.
  • native 포맷 - canonical 포맷의 단점은 찍고 복원하는 것이 종종 느리다는 것입니다. native 포맷은 사용된 상태 백엔드에 특정한 포맷으로 스냅샷을 생성합니다(예: RocksDB의 SST 파일).

native 포맷으로 savepoint을 트리거하는 기능은 Flink 1.15에서 도입되었습니다. 그때까지 savepoint은 canonical 포맷으로 생성되었습니다.

Savepoint 트리거
$ bin/flink savepoint :jobId [:targetDirectory]

이 명령은 ID :jobId인 작업에 대해 savepoint을 트리거하고 생성된 savepoint의 경로를 반환합니다. savepoint을 복원하고 폐기하려면 이 경로가 필요합니다. savepoint을 찍을 타입을 전달할 수도 있습니다. 기본적으로 savepoint은 canonical 포맷으로 찍힙니다.

$ bin/flink savepoint --type [native/canonical] :jobId [:targetDirectory]

위 명령으로 savepoint을 트리거할 때 클라이언트는 savepoint이 완료될 때까지 기다려야 합니다. 따라서 태스크의 상태 크기가 크면 클라이언트가 타임아웃될 수 있습니다. 이 경우 분리(detached) 모드로 savepoint을 트리거할 수 있습니다.

$ bin/flink savepoint :jobId [:targetDirectory] -detached

이 명령을 사용하면 클라이언트는 savepoint의 트리거 id를 얻는 즉시 반환합니다. savepoint의 상태는 REST API rest api를 통해 모니터링할 수 있습니다.

YARN으로 Savepoint 트리거
$ bin/flink savepoint :jobId [:targetDirectory] -yid :yarnAppId

이 명령은 ID :jobId와 YARN 애플리케이션 ID :yarnAppId인 작업에 대해 savepoint을 트리거하고 생성된 savepoint의 경로를 반환합니다.

Savepoint으로 작업 중지
$ bin/flink stop --type [native/canonical] --savepointPath [:targetDirectory] :jobId

이 명령은 ID :jobid인 작업에 대해 savepoint을 원자적으로 트리거하고 작업을 중지합니다. 또한 savepoint을 저장할 대상 파일 시스템 디렉터리를 지정할 수 있습니다. 디렉터리는 JobManager와 TaskManager가 접근할 수 있어야 합니다. savepoint을 찍을 타입을 전달할 수도 있습니다. 기본적으로 savepoint은 canonical 포맷으로 찍힙니다.

분리 모드로 savepoint을 트리거하려면 명령에 -detached 옵션을 추가하세요.

Savepoint에서 재개

$ bin/flink run -s :savepointPath [:runArgs]

이 명령은 작업을 제출하고 재개할 savepoint을 지정합니다. savepoint의 디렉터리 또는 _metadata 파일의 경로를 줄 수 있습니다.

비복원 상태 허용

기본적으로 재개 작업은 savepoint의 모든 상태를 복원 중인 프로그램에 다시 매핑하려고 시도합니다. 연산자를 버린 경우 --allowNonRestoredState(짧게: -n) 옵션으로 새 프로그램에 매핑할 수 없는 상태를 건너뛰는 것을 허용할 수 있습니다.

이 기능을 잘못 사용하면 애플리케이션의 정확성에 심각한 문제가 발생할 수 있습니다. 남아 있는 상태가 적절한 연산자에 정확하게 매핑될 수 있는지 확인하는 것이 중요합니다. 기본적으로 연산자 UID는 위상적 순서에 따라 재할당되므로 상태와 연산자 사이에 잘못된 연결이 발생할 수 있으며, 따라서 상태가 원하는 대로 올바르게 복원되지 않을 수 있다는 점에 유의하세요. 이러한 불일치를 방지하려면 DataStream 작업의 모든 연산자에 UID를 명시적으로 할당하는 것이 좋습니다.

Claim 모드

Claim Mode는 복원 후 Savepoint 또는 외부화된 checkpoint를 구성하는 파일의 소유권을 누가 가지는지 결정합니다. Savepoint과 외부화된 checkpoint는 이 맥락에서 유사하게 동작합니다. 여기서는 명시적으로 달리 언급하지 않는 한 "snapshot"이라고 부릅니다.

언급했듯이 claim 모드는 우리가 복원하는 스냅샷 파일의 소유권을 누가 넘겨받는지 결정합니다. 스냅샷은 사용자 또는 Flink 자체가 소유할 수 있습니다. 스냅샷이 사용자가 소유하면 Flink는 그 파일을 삭제하지 않으며, 또한 그런 스냅샷의 파일 존재에 의존할 수 없습니다. Flink의 통제 밖에서 삭제될 수 있기 때문입니다.

각 claim 모드는 특정 목적을 제공합니다. 그럼에도 우리는 기본 NO_CLAIM 모드가 대부분의 상황에서 좋은 균형이라고 믿습니다. 복원 후 첫 checkpoint에 대해 작은 비용으로 명확한 소유권을 제공하기 때문입니다.

claim 모드를 다음과 같이 전달할 수 있습니다.

$ bin/flink run -s :savepointPath -claimMode :mode -n [:runArgs]

NO_CLAIM (기본값)

NO_CLAIM 모드에서 Flink는 스냅샷의 소유권을 가정하지 않습니다. 파일을 사용자 통제 아래 두고 어떤 파일도 삭제하지 않습니다. 이 모드에서는 같은 스냅샷에서 여러 작업을 시작할 수 있습니다.

Flink가 해당 스냅샷의 어떤 파일에도 의존하지 않도록 하기 위해, 첫 (성공한) checkpoint를 증분이 아닌 전체(full) checkpoint로 강제합니다. 이는 state.backend: rocksdb에서만 차이가 있습니다. 다른 모든 상태 백엔드는 항상 전체 checkpoint를 찍기 때문입니다.

첫 전체 checkpoint가 완료되면 이후 모든 checkpoint는 평소/구성대로 찍힙니다. 결과적으로 checkpoint가 성공하면 원본 스냅샷을 수동으로 삭제할 수 있습니다. 완료된 checkpoint가 없으면 Flink가 실패 시 초기 스냅샷에서 복구하려고 시도하므로 그보다 일찍 삭제할 수 없습니다.

NO_CLAIM claim mode

CLAIM

다른 사용 가능한 모드는 CLAIM 모드입니다. 이 모드에서 Flink는 스냅샷의 소유권을 주장하고 본질적으로 checkpoint처럼 취급합니다. 수명주기를 제어하고 더 이상 복구에 필요하지 않으면 삭제할 수 있습니다. 따라서 스냅샷을 수동으로 삭제하거나 같은 스냅샷에서 두 작업을 시작하는 것은 안전하지 않습니다. Flink는 구성된 수의 checkpoint를 유지합니다.

CLAIM mode

주의:

  1. 유지된 checkpoint는 <checkpoint_dir>/<job_id>/chk-<x> 같은 경로에 저장됩니다. Flink는 <checkpoint_dir>/<job_id> 디렉터리의 소유권을 취하지 않고 chk-<x>만 취합니다. 이전 작업의 디렉터리는 Flink에 의해 삭제되지 않습니다.
  2. Native 포맷은 증분 RocksDB savepoint을 지원합니다. 이러한 savepoint의 경우 Flink는 모든 SST 파일을 savepoint 디렉터리 안에 넣습니다. 이는 그러한 savepoint이 자체 포함(self-contained)되고 재배치 가능함을 의미합니다. CLAIM 모드로 복원하면 이후 checkpoint가 일부 SST 파일을 재사용할 수 있어 savepoint 디렉터리 삭제를 지연시킬 수 있음에 유의하세요.

LEGACY (deprecated)

레거시 모드는 Flink가 1.15까지 작동한 방식입니다. 이 모드에서 Flink는 초기 checkpoint를 절대 삭제하지 않습니다. 동시에 사용자가 그것을 삭제할 수 있는지도 명확하지 않습니다. 여기서 문제는 Flink가 복원된 checkpoint 위에 증분 checkpoint를 즉시 구축할 수 있다는 것입니다. 따라서 이후 checkpoint는 복원된 checkpoint에 의존합니다. 전반적으로 소유권이 잘 정의되어 있지 않습니다.

LEGACY claim mode

주의: LEGACY 모드는 deprecated이며 Flink 2.0에서 제거될 예정입니다. CLAIM 또는 NO_CLAIM 모드를 대신 사용하세요.

Savepoint 폐기

$ bin/flink savepoint -d :savepointPath

이 명령은 :savepointPath에 저장된 savepoint을 폐기합니다.

일반 파일 시스템 연산으로 savepoint을 수동으로 삭제하는 것도 가능하며, 다른 savepoint이나 checkpoint에는 영향을 주지 않습니다(각 savepoint은 자체 포함임을 기억하세요).

구성

execution.checkpointing.savepoint-dir 키 또는 StreamExecutionEnvironment로 기본 savepoint 대상 디렉터리를 구성할 수 있습니다. savepoint을 트리거할 때 이 디렉터리가 savepoint을 저장하는 데 사용됩니다. 트리거 명령으로 사용자 지정 대상 디렉터리를 지정하여 기본값을 덮어쓸 수 있습니다(:targetDirectory 인수 참조).

# Default savepoint target directory
execution.checkpointing.savepoint-dir: hdfs:///flink/savepoints
env.setDefaultSavepointDir("hdfs:///flink/savepoints");

기본값을 구성하지도 않고 사용자 지정 대상 디렉터리를 지정하지도 않으면 savepoint 트리거가 실패합니다.

대상 디렉터리는 JobManager와 TaskManager 모두가 접근할 수 있는 위치여야 합니다. 예: 분산 파일 시스템의 위치.

자주 묻는 질문

작업의 모든 연산자에 ID를 할당해야 하나요?

경험칙으로는 그렇습니다. 엄밀히 말하면 작업의 상태 기반 연산자에만 uid 메서드로 ID를 할당하는 것으로 충분합니다. savepoint은 이러한 연산자의 상태만 포함하며 무상태 연산자는 savepoint의 일부가 아닙니다.

실무적으로는 모든 연산자에 할당하는 것이 권장됩니다. Window 연산자 같은 Flink의 일부 내장 연산자도 상태 기반이며 어떤 내장 연산자가 실제로 상태 기반인지 아닌지가 명확하지 않기 때문입니다. 연산자가 확실히 무상태라고 확신한다면 uid 메서드를 건너뛸 수 있습니다.

작업에 상태를 요구하는 새 연산자를 추가하면 어떻게 되나요?

작업에 새 연산자를 추가하면 상태 없이 초기화됩니다. Savepoint은 각 상태 기반 연산자의 상태를 포함합니다. 무상태 연산자는 savepoint의 일부가 아닙니다. 새 연산자는 무상태 연산자와 유사하게 동작합니다.

작업에서 상태가 있는 연산자를 삭제하면 어떻게 되나요?

기본적으로 savepoint 복원은 모든 상태를 복원된 작업에 다시 매칭하려고 시도합니다. 삭제된 연산자에 대한 상태를 포함하는 savepoint에서 복원하면 실패합니다.

다음 run 명령으로 --allowNonRestoredState(짧게: -n)를 설정하여 비복원 상태를 허용할 수 있습니다.

$ bin/flink run -s :savepointPath -n [:runArgs]

작업에서 상태 기반 연산자를 재정렬하면 어떻게 되나요?

이 연산자에 ID를 할당했다면 평소대로 복원됩니다.

ID를 할당하지 않았다면 상태 기반 연산자의 자동 생성 ID가 재정렬 후 변경될 가능성이 높습니다. 이는 이전 savepoint에서 복원할 수 없게 됩니다.

상태가 없는 연산자를 추가/삭제/재정렬하면 어떻게 되나요?

상태 기반 연산자에 ID를 할당했다면 무상태 연산자는 savepoint 복원에 영향을 주지 않습니다.

ID를 할당하지 않았다면 상태 기반 연산자의 자동 생성 ID가 재정렬 후 변경될 가능성이 높습니다. 이는 이전 savepoint에서 복원할 수 없게 됩니다.

복원할 때 프로그램의 병렬도를 변경하면 어떻게 되나요?

savepoint에서 프로그램을 복원하고 새 병렬도를 지정하기만 하면 됩니다.

안정적인 저장소의 Savepoint 파일을 이동할 수 있나요?

이 질문에 대한 빠른 대답은 현재 "예"입니다. Savepoint은 자체 포함되고 재배치 가능합니다. 파일을 이동하고 어느 위치에서든 복원할 수 있습니다.

더 알아보기 (Learn more)