관리 상태를 위한 사용자 정의 직렬화
관리 상태를 위한 사용자 정의 직렬화 (Custom Serialization for Managed State)
이 페이지는 상태에 사용자 정의 직렬화를 사용해야 하는 사용자를 위한 지침이에요. 사용자 정의 상태 직렬화기를 제공하는 방법과 상태 스키마 진화(schema evolution)를 허용하는 직렬화기 구현을 위한 지침·모범 사례를 다뤄요. 단순히 Flink 자체 직렬화기를 사용한다면 이 페이지는 관련이 없고 무시해도 돼요.
출처: 문서
본문
사용자 정의 상태 직렬화기 사용 (Using custom state serializers)
관리(managed) 연산자 또는 키드 상태를 등록할 때 상태의 이름과 상태 타입 정보를 지정하는 StateDescriptor가 필요해요. 타입 정보는 Flink의 타입 직렬화 프레임워크가 상태에 적절한 직렬화기를 만들 때 사용해요.
자신만의 TypeSerializer 구현으로 StateDescriptor를 직접 인스턴스화하면 이 과정을 완전히 우회하고 Flink가 관리 상태를 직렬화하는 데 사용자 정의 직렬화기를 사용하게 할 수도 있어요:
public class CustomTypeSerializer extends TypeSerializer<Tuple2<String, Integer>> {...};
ListStateDescriptor<Tuple2<String, Integer>> descriptor =
new ListStateDescriptor<>(
"state-name",
new CustomTypeSerializer());
checkpointedState = getRuntimeContext().getListState(descriptor);
상태 직렬화기와 스키마 진화 (State serializers and schema evolution)
이 섹션은 상태 직렬화와 스키마 진화와 관련된 사용자 노출 추상화, 그리고 Flink가 이 추상화와 상호작용하는 데 필요한 내부 세부사항을 설명해요.
savepoint에서 복원할 때 Flink는 이전에 등록된 상태를 읽고 쓰는 데 사용되는 직렬화기를 바꾸는 것을 허용해요. 그래서 사용자는 어떤 특정 직렬화 스키마에 얽매이지 않아요. 상태가 복원될 때 상태에 대해 새 직렬화기가 등록돼요(즉, 복원된 작업에서 상태에 접근하는 데 사용되는 StateDescriptor에 딸린 직렬화기). 이 새 직렬화기는 이전 직렬화기와 다른 스키마를 가질 수 있어요. 따라서 상태 직렬화기를 구현할 때 데이터 읽기/쓰기의 기본 로직 외에 또 하나 중요한 것은 직렬화 스키마가 미래에 어떻게 바뀔 수 있는가예요.
이 맥락에서 스키마라는 용어는 상태 타입의 데이터 모델과 상태 타입의 직렬화된 이진 포맷을 가리키는 것으로 서로 바꿔 쓸 수 있어요. 스키마는 일반적으로 몇 가지 경우에 바뀔 수 있어요:
- 상태 타입의 데이터 스키마가 진화했을 때, 즉 상태로 사용되는 POJO에서 필드를 추가하거나 제거
- 일반적으로 데이터 스키마가 변경된 후에는 직렬화기의 직렬화 포맷을 업그레이드해야 함
- 직렬화기의 구성이 바뀌었을 때
새 실행이 상태의 쓰여진 스키마에 대한 정보를 갖고 스키마가 바뀌었는지 감지하려면, 연산자 상태의 savepoint를 만들 때 상태 직렬화기의 스냅샷을 상태 바이트와 함께 써야 해요. 이는 TypeSerializerSnapshot으로 추상화되며 다음 하위 섹션에서 설명해요.
TypeSerializerSnapshot 추상화 (The TypeSerializerSnapshot abstraction)
public interface TypeSerializerSnapshot<T> {
int getCurrentVersion();
void writeSnapshot(DataOuputView out) throws IOException;
void readSnapshot(int readVersion, DataInputView in, ClassLoader userCodeClassLoader) throws IOException;
TypeSerializerSchemaCompatibility<T> resolveSchemaCompatibility(TypeSerializerSnapshot<T> oldSerializerSnapshot);
TypeSerializer<T> restoreSerializer();
}
public abstract class TypeSerializer<T> {
// ...
public abstract TypeSerializerSnapshot<T> snapshotConfiguration();
}
직렬화기의 TypeSerializerSnapshot은 상태 직렬화기의 쓰기 스키마와 주어진 시점과 동일한 직렬화기를 복원하는 데 필요한 추가 정보에 대한 유일한 진실의 원천(single source of truth) 역할을 하는 시점 정보(point-in-time information)예요. 직렬화기 스냅샷으로 복원 시 무엇을 쓰고 읽을지에 대한 로직은 writeSnapshot과 readSnapshot 메서드에 정의돼요.
스냅샷 자체의 쓰기 스키마도 시간이 지나며 바뀔 수 있다는 점에 주의해요(예: 직렬화기에 대한 정보를 스냅샷에 더 추가하고 싶을 때). 이를 용이하게 하기 위해 스냅샷은 버전이 매겨지며, 현재 버전 번호는 getCurrentVersion 메서드에 정의돼요. 복원 시 직렬화기 스냅샷이 savepoint에서 읽힐 때, 스냅샷이 쓰여진 스키마의 버전이 readSnapshot 메서드에 제공되므로 읽기 구현이 다른 버전을 처리할 수 있어요.
복원 시 새 직렬화기 스키마가 바뀌었는지 감지하는 로직은 resolveSchemaCompatibility 메서드에 구현해야 해요. 연산자의 복원된 실행에서 이전에 등록된 상태가 새 직렬화기로 다시 등록될 때 이전 직렬화기 스냅샷이 이 메서드를 통해 새 직렬화기의 스냅샷에 제공돼요. 이 메서드는 호환성 해석 결과를 나타내는 TypeSerializerSchemaCompatibility를 반환하며, 다음 중 하나일 수 있어요:
- TypeSerializerSchemaCompatibility.compatibleAsIs(): 새 직렬화기가 호환됨을 신호해요. 즉 새 직렬화기가 이전 직렬화기와 동일한 스키마를 가짐을 뜻해요. 새 직렬화기가 호환되도록
resolveSchemaCompatibility메서드에서 재구성되었을 수 있어요. - TypeSerializerSchemaCompatibility.compatibleAfterMigration(): 새 직렬화기가 다른 직렬화 스키마를 가지며, 이전 직렬화기(이전 스키마를 인식)로 바이트를 상태 객체로 읽은 다음 새 직렬화기(새 스키마를 인식)로 객체를 다시 바이트로 쓰는 방식으로 이전 스키마에서 마이그레이션하는 것이 가능함을 신호해요.
- TypeSerializerSchemaCompatibility.incompatible(): 새 직렬화기가 다른 직렬화 스키마를 가지며 이전 스키마에서 마이그레이션하는 것이 불가능함을 신호해요.
마지막 세부사항은 마이그레이션이 필요할 때 이전 직렬화기를 얻는 방법이에요. 직렬화기의 TypeSerializerSnapshot의 또 다른 중요한 역할은 이전 직렬화기를 복원하는 팩토리 역할을 한다는 것이에요. 더 구체적으로 TypeSerializerSnapshot은 restoreSerializer 메서드를 구현해 이전 직렬화기의 스키마와 구성을 인식하고 따라서 이전 직렬화기가 쓴 데이터를 안전하게 읽을 수 있는 직렬화기 인스턴스를 만들어야 해요.
Flink가 TypeSerializer와 TypeSerializerSnapshot 추상화와 상호작용하는 방법 (How Flink interacts with the TypeSerializer and TypeSerializerSnapshot abstractions)
정리하자면 이 섹션은 Flink, 더 구체적으로 상태 백엔드가 이 추상화와 어떻게 상호작용하는지 결론지어요. 상태 백엔드에 따라 상호작용이 약간 다르지만 이는 상태 직렬화기와 그 직렬화기 스냅샷의 구현과 직교해요.
비힙 상태 백엔드 (Off-heap state backends, e.g. EmbeddedRocksDBStateBackend)
- 스키마 A를 가진 상태 직렬화기로 새 상태 등록: 등록된 TypeSerializer가 모든 상태 접근에서 상태를 읽고 씁니다. 상태는 스키마 A로 쓰여져요.
- savepoint 생성: 직렬화기 스냅샷이
TypeSerializer#snapshotConfiguration메서드로 추출돼요. 직렬화기 스냅샷과 이미 직렬화된 상태 바이트(스키마 A)가 savepoint에 쓰여져요. - 복원된 실행이 스키마 B를 가진 새 상태 직렬화기로 복원된 상태 바이트에 재접근: 이전 상태 직렬화기의 스냅샷이 복원돼요. 상태 바이트는 복원 시 역직렬화되지 않고 상태 백엔드에 다시 로드만 돼요(따라서 여전히 스키마 A).
- 상태 백엔드의 상태 바이트를 스키마 A에서 B로 마이그레이션: 호환성 해석이 스키마가 바뀌었고 마이그레이션이 가능함을 반영하면 스키마 마이그레이션이 수행돼요. 스키마 A를 인식하는 이전 상태 직렬화기가 스키마 B를 인식하는 새 직렬화기로 다시 쓰기 위해 상태 바이트를 객체로 역직렬화하는 데 사용돼요. 접근된 상태의 모든 항목은 처리가 계속되기 전에 모두 함께 마이그레이션돼요. 해석이 비호환성을 신호하면 상태 접근이 예외로 실패해요.
힙 상태 백엔드 (Heap state backends, e.g. HashMapStateBackend)
- 스키마 A를 가진 상태 직렬화기로 새 상태 등록: 등록된 TypeSerializer가 상태 백엔드에 의해 유지돼요.
- 모든 상태를 스키마 A로 직렬화하며 savepoint 생성: 직렬화기 스냅샷이
TypeSerializer#snapshotConfiguration메서드로 추출되어 savepoint에 쓰여져요. 상태 객체는 이제 savepoint에 직렬화되어 스키마 A로 쓰여져요. - 복원 시 상태를 힙의 객체로 역직렬화: 이전 상태 직렬화기의 스냅샷이 복원돼요. 스키마 A를 인식하는 이전 직렬화기가
TypeSerializerSnapshot#restoreSerializer()를 통해 얻어져 상태 바이트를 객체로 역직렬화하는 데 사용돼요. 이후로 모든 상태는 이미 역직렬화돼 있어요. - 복원된 실행이 스키마 B를 가진 새 상태 직렬화기로 이전 상태에 재접근: 새 직렬화기를 받으면 이전 직렬화기의 스냅샷이
TypeSerializer#resolveSchemaCompatibility를 통해 새 직렬화기의 스냅샷에 제공돼 스키마 호환성을 확인해요. 호환성 검사가 마이그레이션 필요를 신호하면 힙 백엔드의 경우 모든 상태가 이미 객체로 역직렬화되어 있으므로 이 경우 아무 일도 일어나지 않아요. 해석이 비호환성을 신호하면 상태 접근이 예외로 실패해요. - 모든 상태를 스키마 B로 직렬화하며 또 다른 savepoint 생성: 2단계와 같지만 이제 상태 바이트가 모두 스키마 B에 있어요.
사전 정의된 편리한 TypeSerializerSnapshot 클래스 (Predefined convenient TypeSerializerSnapshot classes)
Flink는 일반적인 시나리오에 사용할 수 있는 두 개의 추상 기본 TypeSerializerSnapshot 클래스를 제공해요: SimpleTypeSerializerSnapshot과 CompositeTypeSerializerSnapshot. 이 사전 정의된 스냅샷을 직렬화기 스냅샷으로 제공하는 직렬화기는 항상 자신만의 독립적인 하위 클래스 구현을 가져야 해요. 이는 서로 다른 직렬화기 간에 스냅샷 클래스를 공유하지 않는 모범 사례에 해당하며, 다음 섹션에서 더 자세히 설명돼요.
SimpleTypeSerializerSnapshot 구현 (Implementing a SimpleTypeSerializerSnapshot)
SimpleTypeSerializerSnapshot은 상태나 구성이 없는 직렬화기를 위한 것이에요. 본질적으로 직렬화기의 직렬화 스키마가 직렬화기 클래스에 의해서만 정의됨을 의미해요. SimpleTypeSerializerSnapshot을 직렬화기 스냅샷 클래스로 사용할 때 호환성 해석의 가능한 결과는 두 가지만 있어요:
TypeSerializerSchemaCompatibility.compatibleAsIs(): 새 직렬화기 클래스가 동일하게 유지되면TypeSerializerSchemaCompatibility.incompatible(): 새 직렬화기 클래스가 이전 것과 다르면
다음은 SimpleTypeSerializerSnapshot을 사용하는 예제로, Flink의 IntSerializer를 예로 들어요:
public class IntSerializerSnapshot extends SimpleTypeSerializerSnapshot<Integer> {
public IntSerializerSnapshot() {
super(() -> IntSerializer.INSTANCE);
}
}
IntSerializer는 상태나 구성이 없어요. 직렬화 포맷은 직렬화기 클래스 자체에 의해서만 정의되고, 다른 IntSerializer만 읽을 수 있어요. 따라서 SimpleTypeSerializerSnapshot의 사용 사례에 적합해요.
SimpleTypeSerializerSnapshot의 기본 상위 생성자는 스냅샷이 현재 복원 중이든 스냅샷 작성 중이든 관계없이 해당 직렬화기의 인스턴스 Supplier를 기대해요. 그 supplier는 복원 직렬화기를 만들고 새 직렬화기가 예상된 직렬화기 클래스와 같은지 검증하는 타입 검사에도 사용돼요.
CompositeTypeSerializerSnapshot 구현 (Implementing a CompositeTypeSerializerSnapshot)
CompositeTypeSerializerSnapshot은 직렬화에 여러 중첩 직렬화기에 의존하는 직렬화기를 위한 것이에요. 자세히 설명하기 전에, 여러 중첩 직렬화기에 의존하는 직렬화기를 이 맥락에서 "외부(outer)" 직렬화기라고 부를게요. 예로는 MapSerializer, ListSerializer, GenericArraySerializer 등이 있어요. MapSerializer를 예로 들면 키와 값 직렬화기가 중첩 직렬화기이고 MapSerializer 자체가 "외부" 직렬화기예요.
이 경우 외부 직렬화기의 스냅샷은 중첩 직렬화기의 호환성을 독립적으로 확인할 수 있도록 중첩 직렬화기의 스냅샷도 포함해야 해요. 외부 직렬화기의 호환성을 해석할 때 각 중첩 직렬화기의 호환성을 고려해야 해요.
CompositeTypeSerializerSnapshot은 이런 종류의 복합 직렬화기를 위한 스냅샷 구현을 돕기 위해 제공돼요. 중첩 직렬화기 스냅샷의 읽기·쓰기, 그리고 모든 중첩 직렬화기의 호환성을 고려한 최종 호환성 결과 해석을 처리해요.
다음은 CompositeTypeSerializerSnapshot을 사용하는 예제로, Flink의 MapSerializer를 예로 들어요:
public class MapSerializerSnapshot<K, V> extends CompositeTypeSerializerSnapshot<Map<K, V>, MapSerializer> {
private static final int CURRENT_VERSION = 1;
public MapSerializerSnapshot() {
super(MapSerializer.class);
}
public MapSerializerSnapshot(MapSerializer<K, V> mapSerializer) {
super(mapSerializer);
}
@Override
public int getCurrentOuterSnapshotVersion() {
return CURRENT_VERSION;
}
@Override
protected MapSerializer createOuterSerializerWithNestedSerializers(TypeSerializer<?>[] nestedSerializers) {
TypeSerializer<K> keySerializer = (TypeSerializer<K>) nestedSerializers[0];
TypeSerializer<V> valueSerializer = (TypeSerializer<V>) nestedSerializers[1];
return new MapSerializer<>(keySerializer, valueSerializer);
}
@Override
protected TypeSerializer<?>[] getNestedSerializers(MapSerializer outerSerializer) {
return new TypeSerializer<?>[] { outerSerializer.getKeySerializer(), outerSerializer.getValueSerializer() };
}
}
CompositeTypeSerializerSnapshot의 하위 클래스로 새 직렬화기 스냅샷을 구현할 때 다음 세 메서드를 구현해야 해요:
getCurrentOuterSnapshotVersion(): 현재 외부 직렬화기 스냅샷의 직렬화된 이진 포맷의 버전을 정의해요.getNestedSerializers(TypeSerializer): 외부 직렬화기가 주어지면 그 중첩 직렬화기를 반환해요.createOuterSerializerWithNestedSerializers(TypeSerializer[]): 중첩 직렬화기가 주어지면 외부 직렬화기의 인스턴스를 만들어요.
위 예제는 중첩 직렬화기의 스냅샷 외에 스냅샷할 추가 정보가 없는 CompositeTypeSerializerSnapshot이에요. 따라서 그 외부 스냅샷 버전은 올림이 필요하지 않을 것으로 기대할 수 있어요. 그러나 일부 다른 직렬화기는 중첩 구성 요소 직렬화기와 함께 영속화해야 하는 추가 정적 구성이 있어요. 그 예는 배열 요소 타입의 클래스를 구성으로 포함하는 Flink의 GenericArraySerializer예요.
이런 경우 CompositeTypeSerializerSnapshot에 추가로 세 메서드를 구현해야 해요:
writeOuterSnapshot(DataOutputView): 외부 스냅샷 정보를 어떻게 쓰는지 정의해요.readOuterSnapshot(int, DataInputView, ClassLoader): 외부 스냅샷 정보를 어떻게 읽는지 정의해요.resolveOuterSchemaCompatibility(TypeSerializerSnapshot): 외부 스냅샷 정보를 기반으로 호환성을 확인해요.
기본적으로 CompositeTypeSerializerSnapshot은 읽기/쓸 외부 스냅샷 정보가 없다고 가정하므로 위 메서드에 대한 빈 기본 구현을 가져요. 하위 클래스에 외부 스냅샷 정보가 있으면 세 메서드를 모두 구현해야 해요.
다음은 외부 스냅샷 정보가 있는 복합 직렬화기 스냅샷에 CompositeTypeSerializerSnapshot을 사용하는 예제로, Flink의 GenericArraySerializer를 예로 들어요:
public final class GenericArraySerializerSnapshot<C> extends CompositeTypeSerializerSnapshot<C[], GenericArraySerializer> {
private static final int CURRENT_VERSION = 1;
private Class<C> componentClass;
public GenericArraySerializerSnapshot() {
super(GenericArraySerializer.class);
}
public GenericArraySerializerSnapshot(GenericArraySerializer<C> genericArraySerializer) {
super(genericArraySerializer);
this.componentClass = genericArraySerializer.getComponentClass();
}
@Override
protected int getCurrentOuterSnapshotVersion() {
return CURRENT_VERSION;
}
@Override
protected void writeOuterSnapshot(DataOutputView out) throws IOException {
out.writeUTF(componentClass.getName());
}
@Override
protected void readOuterSnapshot(int readOuterSnapshotVersion, DataInputView in, ClassLoader userCodeClassLoader) throws IOException {
this.componentClass = InstantiationUtil.resolveClassByName(in, userCodeClassLoader);
}
@Override
protected OuterSchemaCompatibility resolveOuterSchemaCompatibility(
TypeSerializerSnapshot<C[]> oldSerializerSnapshot) {
GenericArraySerializerSnapshot<C[]> oldGenericArraySerializerSnapshot =
(GenericArraySerializerSnapshot<C[]>) oldSerializerSnapshot;
return (this.componentClass == oldGenericArraySerializerSnapshot.componentClass)
? OuterSchemaCompatibility.COMPATIBLE_AS_IS
: OuterSchemaCompatibility.INCOMPATIBLE;
}
@Override
protected GenericArraySerializer createOuterSerializerWithNestedSerializers(TypeSerializer<?>[] nestedSerializers) {
TypeSerializer<C> componentSerializer = (TypeSerializer<C>) nestedSerializers[0];
return new GenericArraySerializer<>(componentClass, componentSerializer);
}
@Override
protected TypeSerializer<?>[] getNestedSerializers(GenericArraySerializer outerSerializer) {
return new TypeSerializer<?>[] { outerSerializer.getComponentSerializer() };
}
}
위 코드 조각에서 두 가지 중요한 점이 있어요. 첫째, 이 CompositeTypeSerializerSnapshot 구현은 스냅샷의 일부로 쓰여지는 외부 스냅샷 정보를 가지므로, getCurrentOuterSnapshotVersion()이 정의하는 외부 스냅샷 버전은 외부 스냅샷 정보의 직렬화 포맷이 바뀔 때마다 올려야 해요.
둘째, 구성 요소 클래스를 쓸 때 Java 직렬화를 피하고 클래스 이름만 쓴 다음 스냅샷을 읽을 때 동적으로 로드하는 방법을 주목해요. 직렬화기 스냅샷의 내용을 쓸 때 Java 직렬화를 피하는 것은 일반적으로 따를 만한 좋은 관행이에요. 이에 대한 자세한 내용은 다음 섹션에서 다뤄요.
구현 참고와 모범 사례 (Implementation notes and best practices)
1. Flink는 직렬화기 스냅샷을 클래스 이름으로 인스턴스화해 복원함
직렬화기의 스냅샷은 등록된 상태가 어떻게 직렬화되었는지에 대한 유일한 진실의 원천으로, savepoint에서 상태를 읽기 위한 진입점 역할을 해요. 이전 상태를 복원하고 접근하려면 이전 상태 직렬화기의 스냅샷을 복원할 수 있어야 해요. Flink는 먼저 TypeSerializerSnapshot을 그 클래스 이름(스냅샷 바이트와 함께 쓰여짐)으로 인스턴스화해 직렬화기 스냅샷을 복원해요. 따라서 의도치 않은 클래스 이름 변경이나 인스턴스화 실패에 영향을 받지 않도록 TypeSerializerSnapshot 클래스는:
- 익명 클래스나 중첩 클래스로 구현되는 것을 피해야 하고,
- 인스턴스화를 위해 public 무인자 생성자를 가져야 해요.
2. 서로 다른 직렬화기 간에 같은 TypeSerializerSnapshot 클래스를 공유하지 않기
스키마 호환성 검사가 직렬화기 스냅샷을 거치므로 여러 직렬화기가 같은 TypeSerializerSnapshot 클래스를 스냅샷으로 반환하면 TypeSerializerSnapshot#resolveSchemaCompatibility와 TypeSerializerSnapshot#restoreSerializer() 메서드의 구현이 복잡해져요. 이는 또한 관심사 분리에 좋지 않아요. 단일 직렬화기의 직렬화 스키마, 구성, 그리고 복원 방법은 자체 전용 TypeSerializerSnapshot 클래스에 통합되어야 해요.
3. 직렬화기 스냅샷 내용에 Java 직렬화 사용을 피하기
영속화된 직렬화기 스냅샷의 내용을 쓸 때 Java 직렬화를 전혀 사용해서는 안 돼요. 예를 들어 스냅샷의 일부로 대상 타입의 클래스를 영속화해야 하는 직렬화기가 있다고 해봐요. 클래스 정보는 Java로 클래스를 직접 직렬화하는 대신 클래스 이름을 써서 영속화해야 해요. 스냅샷을 읽을 때 클래스 이름을 읽고 그 이름으로 클래스를 동적으로 로드해요. 이 관행은 직렬화기 스냅샷을 항상 안전하게 읽을 수 있게 보장해요. 위 예제에서 타입 클래스가 Java 직렬화로 영속화되었다면, 클래스 구현이 바뀌어 Java 직렬화 세부사항에 따라 더 이상 이진 호환이 되지 않으면 스냅샷을 더 이상 읽을 수 없게 될 수 있어요.
Flink 1.7 이전의 deprecated 직렬화기 스냅샷 API에서 마이그레이션 (Migrating from deprecated serializer snapshot APIs before Flink 1.7)
이 섹션은 Flink 1.7 이전에 존재했던 직렬화기와 직렬화기 스냅샷에서의 API 마이그레이션 가이드예요. Flink 1.7 이전에는 직렬화기 스냅샷이 TypeSerializerConfigSnapshot(현재 deprecated이며, 미래에 제거되어 새 TypeSerializerSnapshot 인터페이스로 완전히 대체될 것)으로 구현되었어요. 또한 직렬화기 스키마 호환성 검사의 책임은 TypeSerializer 내에 있었고, TypeSerializer#ensureCompatibility(TypeSerializerConfigSnapshot) 메서드로 구현되었어요.
새 추상화와 이전 추상화의 또 다른 주요 차이는 deprecated TypeSerializerConfigSnapshot이 이전 직렬화기를 인스턴스화할 능력이 없었다는 것이에요. 따라서 직렬화기가 여전히 TypeSerializerConfigSnapshot의 하위 클래스를 스냅샷으로 반환하면, 복원 시 이전 직렬화기를 사용할 수 있도록 직렬화기 인스턴스 자체가 항상 Java 직렬화로 savepoint에 쓰여져요. 이는 매우 바람직하지 않은데, 작업 복원의 성공은 이전 직렬화기 클래스의 가용성, 또는 일반적으로 직렬화기 인스턴스를 복원 시 Java 직렬화로 읽을 수 있는지에 달려 있기 때문이에요. 이는 상태에 같은 직렬화기로 제한될 수밖에 없게 하며, 직렬화기 클래스를 업그레이드하거나 스키마 마이그레이션을 수행하려 할 때 문제가 될 수 있어요.
미래 대비와 상태 직렬화기·스키마 마이그레이션의 유연성을 위해 이전 추상화에서 마이그레이션하는 것을 강력히 권장해요. 단계는 다음과 같아요:
- TypeSerializerSnapshot의 새 하위 클래스를 구현해요. 이것이 직렬화기의 새 스냅샷이 돼요.
TypeSerializer#snapshotConfiguration()메서드에서 직렬화기 스냅샷으로 새 TypeSerializerSnapshot을 반환해요.- Flink 1.7 이전에 존재했던 savepoint에서 작업을 복원한 다음 다시 savepoint를 만들어요. 이 단계에서 직렬화기의 이전 TypeSerializerConfigSnapshot이 여전히 클래스패스에 있어야 하고,
TypeSerializer#ensureCompatibility(TypeSerializerConfigSnapshot)메서드의 구현을 제거해서는 안 돼요. 이 과정의 목적은 이전 savepoint에 쓰여진 TypeSerializerConfigSnapshot을 직렬화기용으로 새로 구현한 TypeSerializerSnapshot으로 대체하는 것이에요. - Flink 1.7로 만든 savepoint가 생기면 savepoint에는 상태 직렬화기 스냅샷으로 TypeSerializerSnapshot이 포함되고 직렬화기 인스턴스는 더 이상 savepoint에 쓰여지지 않아요. 이 시점이 되면 이전 추상화의 모든 구현을 제거해도 안전해요(이전 TypeSerializerConfigSnapshot 구현과
TypeSerializer#ensureCompatibility(TypeSerializerConfigSnapshot)를 직렬화기에서 제거).
Flink 1.19 이전의 deprecated TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializer newSerializer)에서 마이그레이션 (Migrating from deprecated TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializer newSerializer) before Flink 1.19)
이 섹션은 Flink 1.19 이전에 존재했던 직렬화기 스냅샷에서의 메서드 마이그레이션 가이드예요. Flink 1.19 이전에는 사용자 정의 직렬화기로 데이터를 처리할 때 이전 직렬화기(아마 Flink 라이브러리)의 스키마 호환성이 미래 요구를 충족해야 했어요. 그렇지 않으면 이전 직렬화기의 TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializer newSerializer)를 수정해야 했어요. 새 직렬화기에서 이전 직렬화기와의 호환성을 지정할 방법이 없어 일부 시나리오에서 스키마 진화가 지원되지 않았어요.
따라서 Flink 1.19부터 스키마 호환성 해석의 방향이 반전됐어요. 이전 메서드 TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializer newSerializer)는 제거되었으며 TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializerSnapshot oldSerializerSnapshot)로 대체해야 해요. 전환을 위해 다음 단계를 따르세요:
- 로직이 원래
TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializer newSerializer)와 같아야 하는TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializerSnapshot oldSerializerSnapshot)을 구현해요. - 이전 메서드
TypeSerializerSnapshot#resolveSchemaCompatibility(TypeSerializer newSerializer)를 제거해요.