애플리케이션 및 Flink 버전 업그레이드

Flink DataStream 프로그램은 일반적으로 몇 주, 몇 달, 심지어 몇 년과 같은 장기간 실행되도록 설계됩니다. 모든 장기 실행 서비스와 마찬가지로 Flink 스트리밍 애플리케이션도 유지보수가 필요합니다. 버그 수정, 개선 구현, 또는 애플리케이션을 이후 버전의 Flink 클러스터로 마이그레이션하는 것이 포함됩니다.

출처: 문서

본문

이 문서는 Flink 스트리밍 애플리케이션을 업데이트하는 방법과 실행 중인 스트리밍 애플리케이션을 다른 Flink 클러스터로 마이그레이션하는 방법을 설명합니다.

API 호환성 보장 (API compatibility guarantees)

사용자용으로 설계된 Java API의 클래스와 멤버는 다음과 같은 안정성 어노테이션으로 표시됩니다:

  • Public
  • PublicEvolving
  • Experimental

클래스의 어노테이션은 달리 명시되지 않는 한 그 클래스의 모든 멤버에도 적용됩니다.

이런 어노테이션이 없는 API는 Flink 내부용으로 간주되며 보장이 제공되지 않습니다.

source 호환 API는 그 API에 대해 작성된 코드가 이후 버전에서도 계속 컴파일된다는 뜻입니다. binary 호환 API는 그 API에 대해 컴파일된 코드가 이후 버전에서도 계속 실행된다는 뜻입니다.

다음 표는 특정 릴리스로 업그레이드할 때 각 어노테이션에 대한 source/binary 호환성 보장을 나열합니다:

어노테이션 Major 릴리스(Source / Binary) Minor 릴리스(Source / Binary) Patch 릴리스(Source / Binary)
Public / / /
PublicEvolving / / /
Experimental / / /

예시 1.15.2에서 Public API에 대해 작성된 코드를 고려해보세요:

  • 1.15.3으로 업그레이드하면 재컴파일 없이 계속 실행될 수 있습니다. Public API의 patch 버전 업그레이드는 binary 호환성을 보장하기 때문입니다.
  • 1.15.x에서 1.16.0으로 업그레이드하면 재컴파일해야 할 수 있습니다. Public API의 minor 버전 업그레이드는 binary 호환성이 아닌 source 호환성만 제공하기 때문입니다.
  • 1.x에서 2.x로 업그레이드하면 코드 변경이 필요할 수 있습니다. Public API의 major 버전 업그레이드는 sourcebinary 호환성 모두 제공하지 않기 때문입니다.

1.15.2에서 PublicEvolving API에 대해 작성된 코드를 고려해보세요:

  • 1.15.3으로 업그레이드하면 재컴파일 없이 계속 실행될 수 있습니다. PublicEvolving API의 patch 버전 업그레이드는 binary 호환성을 보장하기 때문입니다.
  • 1.15.x에서 1.16.0으로 업그레이드하면 코드 변경이 필요할 수 있습니다. PublicEvolving API의 minor 버전 업그레이드는 sourcebinary 호환성을 모두 제공하지 않기 때문입니다.

Deprecated API 마이그레이션 기간 (Deprecated API Migration Period)

API가 deprecated되면 @Deprecated 어노테이션으로 표시되고 deprecation 메시지가 Javadoc에 추가됩니다. FLIP-321에 따라 1.18 릴리스부터 각 deprecated API는 API 안정성 수준에 따라 보장된 마이그레이션 기간을 갖습니다:

어노테이션 보장된 마이그레이션 기간 마이그레이션 기간 후 제거 가능
Public 2 minor 릴리스 다음 major 버전
PublicEvolving 1 minor 릴리스 다음 minor 버전
Experimental 해당 minor 릴리스의 1 patch 릴리스 다음 patch 버전

deprecated API의 소스 코드는 최소한 보장된 마이그레이션 기간 동안 유지되며, 마이그레이션 기간이 지난 후에는 언제든지 제거될 수 있습니다.

예시 릴리스 순서가 1.18, 1.19, 1.20, 2.0, 2.1, …, 3.0이라고 가정할 때,

  • 1.18에서 Public API가 deprecated되면 2.0까지 제거되지 않습니다.
  • 1.20에서 Public API가 deprecated되면 마이그레이션 기간이 2 minor 릴리스이므로 소스 코드가 2.0에서 유지됩니다. 또한 Public API는 major 버전 전체에서 source 호환성을 유지해야 하므로 2.x 버전 전체에서 소스 코드가 유지되고 최소 3.0에서 제거됩니다.
  • 1.18에서 PublicEvolving API가 deprecated되면 최소 1.20에서 제거됩니다.
  • 1.20에서 PublicEvolving API가 deprecated되면 마이그레이션 기간이 1 minor 릴리스이므로 소스 코드가 2.0에서 유지됩니다. 최소 2.1에서 소스 코드가 제거될 수 있습니다.
  • 1.18.0에서 Experimental API가 deprecated되면 1.18.1 동안 소스 코드가 유지되고 최소 1.18.2에서 제거됩니다. 또한 1.19.0에서 소스 코드가 제거될 수 있습니다.

자세한 내용은 FLIP-321 위키를 확인하세요.

스트리밍 애플리케이션 재시작 (Restarting Streaming Applications)

스트리밍 애플리케이션을 업그레이드하거나 애플리케이션을 다른 클러스터로 마이그레이션하는 행동 지침은 Flink의 Savepoint 기능에 기반합니다. savepoint는 특정 시점의 애플리케이션 상태에 대한 일관된 스냅샷입니다.

실행 중인 스트리밍 애플리케이션에서 savepoint를 만드는 방법은 두 가지입니다.

  • savepoint를 만들고 계속 처리하기:
> ./bin/flink savepoint <jobID> [pathToSavepoint]

이전 시점에서 애플리케이션을 재시작할 수 있도록 주기적으로 savepoint를 만드는 것이 권장됩니다. detached 모드에서 savepoint를 트리거하려면 -detached 옵션을 추가하면 됩니다.

  • savepoint를 만들고 애플리케이션을 단일 작업으로 중지하기:
> ./bin/flink cancel -s [pathToSavepoint] <jobID>

이것은 savepoint 완료 직후 애플리케이션이 취소된다는 뜻입니다. 즉 savepoint 이후 다른 checkpoint는 찍히지 않습니다.

애플리케이션에서 만든 savepoint가 주어지면 동일하거나 호환되는 애플리케이션(애플리케이션 상태 호환성 섹션 참고)을 그 savepoint에서 시작할 수 있습니다. savepoint에서 애플리케이션을 시작한다는 것은 연산자의 상태가 savepoint에 영속화된 연산자 상태로 초기화된다는 뜻입니다. 이는 savepoint를 사용해 애플리케이션을 시작함으로써 수행됩니다.

> ./bin/flink run -d -s [pathToSavepoint] ~/application.jar

시작된 애플리케이션의 연산자는 savepoint가 찍힌 시점에 원래 애플리케이션(savepoint를 만든 애플리케이션)의 연산자 상태로 초기화됩니다. 시작된 애플리케이션은 정확히 그 시점부터 처리를 계속합니다.

참고: Flink는 애플리케이션 상태를 일관되게 복원하지만, 외부 시스템에 대한 쓰기는 되돌릴 수 없습니다. 애플리케이션을 중지하지 않고 만든 savepoint에서 재개하면 문제가 될 수 있습니다. 이 경우 애플리케이션은 savepoint 이후 데이터를 이미 방출했을 수 있습니다. 재시작된 애플리케이션은(애플리케이션 로직을 변경했는지에 따라) 같은 데이터를 다시 방출할 수 있습니다. 이 동작의 정확한 영향은 SinkFunction과 저장 시스템에 따라 크게 다를 수 있습니다. 중복 방출되는 데이터는 Cassandra 같은 key-value 저장소에 대한 멱등 쓰기에서는 괜찮을 수 있지만 Kafka 같은 내구성 있는 로그에 대한 append에서는 문제가 될 수 있습니다. 어떤 경우든 재시작된 애플리케이션의 동작을 신중히 확인하고 테스트해야 합니다.

애플리케이션 상태 호환성 (Application State Compatibility)

버그를 수정하거나 개선하기 위해 애플리케이션을 업그레이드할 때, 보통 실행 중인 애플리케이션의 상태를 보존하면서 애플리케이션 로직을 교체하는 것이 목표입니다. 이는 원래 애플리케이션에서 만든 savepoint에서 업그레이드된 애플리케이션을 시작함으로써 수행합니다. 하지만 이는 두 애플리케이션이 *상태 호환(state compatible)*일 때만 가능합니다. 즉 업그레이드된 애플리케이션의 연산자가 원래 애플리케이션 연산자의 상태로 자신의 상태를 초기화할 수 있어야 합니다.

이 섹션에서는 애플리케이션이 상태 호환이 되도록 수정하는 방법을 논의합니다.

DataStream API

연산자 상태 매칭 (Matching Operator State)

애플리케이션이 savepoint에서 재시작될 때 Flink는 savepoint에 저장된 연산자 상태를 시작된 애플리케이션의 상태 저장 연산자와 매칭합니다. 매칭은 savepoint에도 저장된 연산자 ID를 기준으로 수행됩니다. 각 연산자에는 애플리케이션의 연산자 토폴로지에서 연산자의 위치에서 파생되는 기본 ID가 있습니다. 따라서 수정되지 않은 애플리케이션은 항상 자신의 savepoint 중 하나에서 재시작될 수 있습니다. 하지만 애플리케이션이 수정되면 연산자의 기본 ID가 바뀔 가능성이 높습니다. 따라서 수정된 애플리케이션은 연산자 ID가 명시적으로 지정된 경우에만 savepoint에서 시작할 수 있습니다. 연산자에 ID를 할당하는 것은 매우 간단하며 uid(String) 메서드를 사용합니다:

DataStream<String> mappedEvents = events
  .map(new MyStatefulMapFunc()).uid("mapper-1");

참고: savepoint에 저장된 연산자 ID와 시작할 애플리케이션의 연산자 ID가 같아야 하므로, 향후 업그레이드될 수 있는 애플리케이션의 모든 연산자에 고유 ID를 할당하는 것이 강력히 권장됩니다. 이 조언은 모든 연산자(명시적으로 선언된 연산자 상태가 있든 없든)에 적용됩니다. 일부 연산자는 사용자에게 보이지 않는 내부 상태를 갖기 때문입니다. 연산자 ID를 할당하지 않은 애플리케이션을 업그레이드하는 것은 훨씬 어려우며 setUidHash() 메서드를 사용하는 저수준 우회 방법으로만 가능할 수 있습니다.

중요: 1.3.x부터 이는 체인의 일부인 연산자에도 적용됩니다.

기본적으로 savepoint에 저장된 모든 상태는 시작되는 애플리케이션의 연산자와 매칭되어야 합니다. 하지만 사용자는 savepoint에서 애플리케이션을 시작할 때 연산자와 매칭할 수 없는 상태를 건너뛰어(그리고 버려)도 된다고 명시적으로 동의할 수 있습니다. savepoint에서 상태를 찾지 못한 상태 저장 연산자는 기본 상태로 초기화됩니다. 사용자는 ExecutionConfig#disableAutoGeneratedUIDs를 호출해 모범 사례를 강제할 수 있습니다. 이는 어떤 연산자에도 사용자 지정 고유 ID가 없으면 작업 제출을 실패시킵니다.

상태 저장 연산자와 사용자 함수 (Stateful Operators and User Functions)

애플리케이션을 업그레이드할 때 사용자 함수와 연산자는 한 가지 제한과 함께 자유롭게 수정할 수 있습니다. 연산자 상태의 데이터 타입을 변경하는 것은 불가능합니다. 이는 savepoint의 상태가 연산자에 로드되기 전에 다른 데이터 타입으로 변환될 수 없기(현재) 때문에 중요합니다. 따라서 업그레이드 시 연산자 상태의 데이터 타입을 변경하면 애플리케이션 상태 일관성이 깨지고 업그레이드된 애플리케이션이 savepoint에서 재시작되지 못합니다.

연산자 상태는 사용자 정의 또는 내부 상태일 수 있습니다.

  • 사용자 정의 연산자 상태: 사용자 정의 연산자 상태를 갖는 함수에서 상태 타입은 사용자가 명시적으로 정의합니다. 연산자 상태의 데이터 타입을 변경할 수는 없지만, 다른 데이터 타입의 두 번째 상태를 정의하고 원래 상태를 새 상태로 마이그레이션하는 로직을 구현하는 우회 방법이 있습니다. 이 접근 방식은 좋은 마이그레이션 전략과 키 분할 상태 동작에 대한 확실한 이해를 요구합니다.

  • 내부 연산자 상태: window나 join 연산자 같은 연산자는 사용자에게 노출되지 않는 내부 연산자 상태를 보유합니다. 이 연산자들의 내부 상태 데이터 타입은 연산자의 입력 또는 출력 타입에 의존합니다. 결과적으로 해당 입력 또는 출력 타입을 변경하면 애플리케이션 상태 일관성이 깨지고 업그레이드를 막습니다. 다음 표는 내부 상태를 가진 연산자를 나열하고 상태 데이터 타입이 입력/출력 타입과 어떻게 관련되는지 보여줍니다. 키드 스트림에 적용되는 연산자의 경우 키 타입(KEY)도 항상 상태 데이터 타입의 일부입니다.

연산자 내부 연산자 상태의 데이터 타입
ReduceFunction[IOT] IOT (입력 및 출력 타입) [, KEY]
WindowFunction[IT, OT, KEY, WINDOW] IT (입력 타입), KEY
AllWindowFunction[IT, OT, WINDOW] IT (입력 타입)
JoinFunction[IT1, IT2, OT] IT1, IT2 (1번·2번 입력 타입), KEY
CoGroupFunction[IT1, IT2, OT] IT1, IT2 (1번·2번 입력 타입), KEY
내장 집계 (sum, min, max, minBy, maxBy) 입력 타입 [, KEY]
애플리케이션 토폴로지 (Application Topology)

기존 연산자 하나 이상의 로직을 변경하는 것 외에도, 애플리케이션 토폴로지를 변경(연산자 추가/제거, 연산자 병렬도 변경, 연산자 체이닝 동작 수정)하여 애플리케이션을 업그레이드할 수 있습니다.

토폴로지를 변경해 애플리케이션을 업그레이드할 때 애플리케이션 상태 일관성을 보존하려면 몇 가지를 고려해야 합니다.

  • 무상태 연산자 추가/제거: 아래 경우 중 하나가 아니라면 문제가 없습니다.
  • 상태 저장 연산자 추가: 다른 연산자의 상태를 이어받지 않는 한 연산자 상태가 기본 상태로 초기화됩니다.
  • 상태 저장 연산자 제거: 다른 연산자가 이어받지 않으면 제거된 연산자의 상태가 손실됩니다. 업그레이드된 애플리케이션을 시작할 때 상태를 버리겠다고 명시적으로 동의해야 합니다.
  • 연산자 입력/출력 타입 변경: 내부 상태를 가진 연산자 앞이나 뒤에 새 연산자를 추가할 때, 내부 연산자 상태의 데이터 타입을 보존하려면 상태 저장 연산자의 입력/출력 타입을 수정하지 않아야 합니다.
  • 연산자 체이닝 변경: 성능 향상을 위해 연산자를 체이닝할 수 있습니다. 1.3.x 이후 만들어진 savepoint에서 복원할 때 상태 일관성을 유지하면서 체인을 수정할 수 있습니다. 상태 저장 연산자가 체인 밖으로 이동하도록 체인을 끊는 것이 가능합니다. 또한 새거나 기존 상태 저장 연산자를 체인에 추가/주입하거나 체인 내 연산자 순서를 수정할 수도 있습니다. 하지만 savepoint를 1.3.x로 업그레이드할 때는 체이닝 측면에서 토폴로지가 변경되지 않았는지가 중요합니다. 체인의 일부인 모든 연산자에는 위의 연산자 상태 매칭 섹션에서 설명한 대로 ID가 할당되어야 합니다.

Table API & SQL

Table API & SQL 프로그램의 선언적 특성 때문에 기본 연산자 토폴로지와 상태 표현은 대부분 테이블 플래너에 의해 결정되고 최적화됩니다.

쿼리와 Flink 버전 모두에 대한 어떤 변경이든 상태 비호환을 초래할 수 있음을 인지하세요. 새로운 각 major-minor Flink 버전(예: 1.12에서 1.13)은 실행 계획을 바꾸는 새로운 옵티마이저 규칙이나 더 특화된 런타임 연산자를 도입할 수 있습니다. 하지만 커뮤니티는 patch 버전을 상태 호환(예: 1.13.1에서 1.13.2)으로 유지하려고 노력합니다.

자세한 내용은 테이블 상태 관리 섹션을 참고하세요.

이 섹션은 Flink를 버전 간 업그레이드하고 버전 간에 작업을 마이그레이션하는 일반적인 방법을 설명합니다.

간단히 말하면 이 절차는 두 가지 기본 단계로 구성됩니다:

  • 마이그레이션할 작업에 대해 이전(구) Flink 버전에서 savepoint를 만듭니다.
  • 새 Flink 버전에서 이전에 만든 savepoint로 작업을 재개합니다.

이 두 기본 단계 외에도 Flink 버전을 변경하는 방식에 따라 추가 단계가 필요할 수 있습니다. 이 가이드에서는 Flink 버전 간 업그레이드의 두 가지 접근 방식을 구분합니다: 제자리(in-place) 업그레이드와 섀도우 복사(shadow copy) 업그레이드입니다.

제자리 업데이트의 경우, savepoint를 만든 후 다음을 수행해야 합니다:

  • 실행 중인 모든 작업을 중지/취소합니다.
  • 구 Flink 버전을 실행하는 클러스터를 종료합니다.
  • 클러스터에서 Flink를 새 버전으로 업그레이드합니다.
  • 새 버전으로 클러스터를 재시작합니다.

섀도우 복사의 경우 다음을 수행해야 합니다:

  • savepoint에서 재개하기 전에 기존 Flink 설치 옆에 새 Flink 버전의 새 설치를 준비합니다.
  • 새 Flink 설치로 savepoint에서 재개합니다.
  • 모든 것이 잘 실행되면 구 Flink 클러스터를 중지하고 종료합니다.

다음에서는 먼저 성공적인 작업 마이그레이션의 전제 조건을 제시한 다음 앞서 개요한 단계에 대해 더 자세히 다룹니다.

전제 조건 (Preconditions)

마이그레이션을 시작하기 전에 마이그레이션하려는 작업이 savepoint에 대한 모범 사례를 따르고 있는지 확인하세요. 특히 작업의 연산자에 명시적 uid가 설정되었는지 확인하는 것을 권장합니다.

이것은 소프트 전제 조건이며, uid 할당을 잊었어도 복원 여전히 작동해야 합니다. 작동하지 않는 경우가 있으면 setUidHash(String hash) 호출을 사용해 이전 Flink 버전에서 생성된 레거시 vertex id를 작업에 수동으로 추가할 수 있습니다. 각 연산자(연산자 체인에서는 head 연산자만)에 대해 웹 UI나 로그에서 볼 수 있는 연산자 해시를 나타내는 32자 16진수 문자열을 할당해야 합니다.

연산자 uid 외에도 현재 작업 마이그레이션을 실패하게 만드는 하드 전제 조건이 두 가지 있습니다:

  • semi-asynchronous 모드로 checkpoint된 RocksDB의 상태에 대한 마이그레이션은 지원하지 않습니다. 구 버전 작업이 이 모드를 사용했다면, 마이그레이션의 기반으로 사용하는 savepoint를 만들기 전에 여전히 작업을 fully-asynchronous 모드로 변경할 수 있습니다.

  • 또 다른 중요한 전제 조건은 모든 savepoint 데이터가 새 설치에서 동일한(절대) 경로로 접근 가능해야 한다는 것입니다. 여기에는 savepoint 파일 내부에서 참조되는 추가 파일(상태 백엔드 스냅샷의 출력)에 대한 접근도 포함되며, State Processor API 수정으로 인한 추가 참조 savepoint도 포함하되 이에 국한되지 않습니다.

1단계: savepoint로 기존 작업 중지 (STEP 1: Stop the existing job with a savepoint)

버전 마이그레이션의 첫 번째 주요 단계는 savepoint를 만들고 구 Flink 버전에서 실행 중인 작업을 중지하는 것입니다. 다음과 같이 할 수 있습니다:

$ bin/flink stop [--savepointPath :savepointPath] :jobId

detached 모드에서 savepoint를 트리거하려면 명령에 -detached 옵션을 추가하세요. 자세한 내용은 savepoint 문서를 읽어보세요.

2단계: 클러스터를 새 Flink 버전으로 업데이트 (STEP 2: Update your cluster to the new Flink version)

이 단계에서는 클러스터의 프레임워크 버전을 업데이트합니다. 기본적으로 Flink 설치 내용을 새 버전으로 교체한다는 뜻입니다. 이 단계는 클러스터에서 Flink를 실행하는 방식(예: standalone 등)에 따라 달라질 수 있습니다. Flink를 클러스터에 설치하는 방법을 잘 모르면 배포 및 클러스터 설정 문서를 읽어보세요.

3단계: 새 Flink 버전에서 savepoint로 작업 재개 (STEP 3: Resume the job under the new Flink version from savepoint)

작업 마이그레이션의 마지막 단계로, 업데이트된 클러스터에서 위에서 만든 savepoint로부터 재개합니다. 다음과 같이 할 수 있습니다:

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

자세한 내용은 savepoint 문서를 확인하세요.

호환성 표 (Compatibility Table)

Savepoint는 아래 표에 표시된 대로 Flink 버전 간 호환됩니다.

참고: 여기서 "호환"은 구체적으로 savepoint 내부 데이터 포맷의 호환성을 가리킵니다. SQL 연산자나 다른 상위 수준 변경의 호환성은 다루지 않습니다. 실제로는 Flink SQL 의미론에 변화가 없다면 이 상태 포맷 호환성은 보통 작업 수준 호환성도 보장합니다.

생성 \ 복원 1.17.x 1.18.x 1.19.x 1.20.x 제한 사항
1.8.x O O O O
1.9.x O O O O
1.10.x O O O O
1.11.x O O O O
1.12.x O O O O
1.13.x O O O O 정렬되지 않은 checkpoint로 1.12.x에서 1.13.x로 업그레이드하지 마세요. 마이그레이션에는 savepoint를 사용하세요.
1.14.x O O O O
1.15.x O O O O Table API: 1.15.0과 1.15.1은 연산자에 대해 비결정적 UID를 생성해 상태 복원이나 다음 patch 버전으로의 업그레이드를 어렵게/불가능하게 합니다. 새로운 table.exec.uid.generation 설정 옵션(올바른 기본 동작)은 컴파일되지 않은 계획의 새 파이프라인에 UID 설정을 비활성화합니다. 기존 파이프라인은 안정적인 환경으로 인해 1.15.0/1 동작이 허용 가능했다면 table.exec.uid.generation=ALWAYS로 설정할 수 있습니다. 자세한 내용은 FLINK-28861 참고.
1.16.x O O O O
1.17.x O O O O
1.18.x O O O
1.19.x O O
1.20.x O

이전 Flink 버전의 savepoint 호환성 정보는 마지막 Compatibility Table을 참고하세요.

더 알아보기 (Learn more)