State Schema Evolution
State Schema Evolution (상태 스키마 진화)
Apache Flink 스트리밍 애플리케이션은 일반적으로 무기한 또는 오랜 기간 동안 실행되도록 설계됩니다. 모든 장기 실행 서비스와 마찬가지로, 바뀌는 요구 사항에 적응하기 위해 애플리케이션을 업데이트해야 합니다. 애플리케이션이 사용하는 데이터 스키마도 마찬가지이며, 애플리케이션과 함께 진화합니다.
출처: 문서
본문
Apache Flink 스트리밍 애플리케이션은 일반적으로 무기한 또는 오랜 기간(주, 월, 심지어 년) 동안 실행되도록 설계됩니다. 모든 장기 실행 서비스와 마찬가지로 애플리케이션은 바뀌는 요구 사항에 적응하기 위해 업데이트되어야 합니다. 애플리케이션이 사용하는 데이터 스키마도 마찬가지이며 애플리케이션과 함께 진화합니다.
이 페이지는 상태 유형의 데이터 스키마를 진화시키는 방법에 대한 개요를 제공합니다. 현재 제한 사항은 유형과 상태 구조(ValueState, ListState 등)에 따라 다릅니다.
이 페이지의 정보는 Flink 자체의 type serialization framework가 생성한 상태 직렬화기를 사용하는 경우에만 관련됩니다. 즉, 상태를 선언할 때 제공된 상태 디스크립터가 특정
TypeSerializer나TypeInformation을 사용하도록 구성되지 않은 경우이며, 이 경우 Flink가 상태 유형에 대한 정보를 추론합니다:
ListStateDescriptor<MyPojoType> descriptor =
new ListStateDescriptor<>(
"state-name",
MyPojoType.class);
checkpointedState = getRuntimeContext().getListState(descriptor);
내부적으로 상태 스키마를 진화시킬 수 있는지 여부는 영속된 상태 바이트를 읽고/쓰는 데 사용되는 직렬화기에 달려 있습니다. 간단히 말해, 등록된 상태의 스키마는 해당 직렬화기가 이를 제대로 지원할 때에만 진화할 수 있습니다. 이는 Flink의 타입 직렬화 프레임워크가 생성한 직렬화기에 의해 투명하게 처리됩니다(현재 지원 범위는 아래에 나열).
상태 유형에 대해 커스텀 TypeSerializer를 구현하려고 하고 상태 스키마 진화를 지원하도록 직렬화기를 구현하는 방법을 배우고 싶다면 Custom State Serialization을 참고하세요. 그 문서에는 상태 스키마 진화를 지원하기 위한 상태 직렬화기와 Flink의 상태 백엔드 간 상호작용에 대한 필요한 내부 세부 사항도 다룹니다.
상태 스키마 진화 (Evolving state schema)
주어진 상태 유형의 스키마를 진화시키려면 다음 단계를 따릅니다:
- Flink 스트리밍 job의 savepoint를 생성합니다.
- 애플리케이션에서 상태 유형을 업데이트합니다(예: Avro 유형 스키마 변경).
- savepoint에서 job을 복원합니다. 상태에 처음 접근할 때 Flink는 상태에 대해 스키마가 변경되었는지 평가하고, 필요한 경우 상태 스키마를 마이그레이션합니다.
변경된 스키마에 적응하기 위해 상태를 마이그레이션하는 과정은 자동으로, 각 상태에 대해 독립적으로 발생합니다. 이 과정은 Flink가 상태의 새 직렬화기가 이전 직렬화기와 다른 직렬화 스키마를 가지는지 먼저 확인하여 내부적으로 수행됩니다. 그렇다면 이전 직렬화기로 상태를 객체로 읽고, 새 직렬화기로 다시 바이트로 씁니다.
마이그레이션 과정에 대한 추가 세부 사항은 이 문서의 범위를 벗어나며 여기를 참고하세요.
스키마 진화가 지원되는 데이터 유형 (Supported data types for schema evolution)
현재 스키마 진화는 POJO와 Avro 유형에 대해서만 지원됩니다. 따라서 상태의 스키마 진화가 중요하다면, 상태 데이터 유형에 항상 POJO 또는 Avro를 사용하는 것이 현재 권장됩니다.
더 많은 복합 유형에 대한 지원을 확장할 계획이 있습니다. 자세한 내용은 FLINK-10896을 참고하세요.
POJO 유형 (POJO types)
Flink는 POJO 유형의 스키마 진화를 다음 규칙 집합을 기반으로 지원합니다:
- 필드를 제거할 수 있습니다. 제거되면 미래의 체크포인트와 savepoint에서 제거된 필드의 이전 값은 버려집니다.
- 새 필드를 추가할 수 있습니다. 새 필드는 Java가 정의한 유형의 기본값으로 초기화됩니다.
- 선언된 필드 유형은 변경할 수 없습니다.
- POJO 유형의 클래스 이름은 클래스의 네임스페이스를 포함하여 변경할 수 없습니다.
Flink가 Flink 1.19부터 POJO 유형으로 취급하는 Java records에도 동일한 규칙이 적용됩니다.
POJO 유형 상태의 스키마는 1.8.0보다 새로운 Flink 버전으로 이전 savepoint에서 복원할 때만 진화할 수 있습니다. 1.8.0보다 오래된 Flink 버전으로 복원할 때는 스키마를 변경할 수 없습니다.
Avro 유형 (Avro types)
Flink는 Avro의 스키마 해석 규칙에 의해 스키마 변경이 호환 가능한 것으로 간주되는 한, Avro 유형 상태의 스키마 진화를 완전히 지원합니다.
하나의 제한은 상태 유형으로 사용되는 Avro 생성 클래스는 job이 복원될 때 재배치(relocate)되거나 다른 네임스페이스를 가질 수 없다는 것입니다.
스키마 마이그레이션 제한 사항 (Schema Migration Limitations)
Flink의 스키마 마이그레이션에는 정확성을 보장하기 위해 필요한 몇 가지 제한 사항이 있습니다. 이러한 제한을 우회하고 특정 사용 사례에서 안전하다는 것을 이해해야 하는 사용자는 custom serializer 또는 state processor api 사용을 고려하세요.
키의 스키마 진화는 지원되지 않습니다. (Schema evolution of keys is not supported)
키의 구조는 비결정적 동작으로 이어질 수 있으므로 마이그레이션할 수 없습니다. 예를 들어 POJO가 키로 사용되고 한 필드가 제거되면, 갑자기 동일해지는 여러 개의 분리된 키가 생길 수 있습니다. Flink는 해당 값을 병합할 방법이 없습니다.
또한 RocksDB 상태 백엔드는 hashCode 메서드가 아닌 바이너리 객체 정체성에 의존합니다. 키의 객체 구조에 대한 어떤 변경도 비결정적 동작으로 이어질 수 있습니다.
Kryo는 스키마 진화에 사용할 수 없습니다. (Kryo cannot be used for schema evolution)
Kryo를 사용하면 프레임워크가 호환되지 않는 변경이 이루어졌는지 검증할 방법이 없습니다.
즉, 특정 유형을 포함하는 데이터 구조가 Kryo로 직렬화되면, 포함된 유형은 스키마 진화를 거칠 수 없습니다.
예를 들어 POJO가 List를 포함하면 List와 그 내용물은 Kryo로 직렬화되며 SomeOtherPojo에 대해 스키마 진화는 지원되지 않습니다.