State V2로 작업하기

State V2로 작업하기 (Working with State V2, New APIs)

이 섹션에서는 상태 있는(stateful) 프로그램을 작성하기 위해 Flink가 제공하는 새로운 API에 대해 배웁니다. 상태 있는 스트림 처리의 배후 개념을 배우려면 Stateful Stream Processing을 참조하세요.

새 상태 API는 이전 API보다 더 유연하게 설계되었습니다. 사용자는 비동기 상태 연산을 수행할 수 있어 더 강력하고 더 효율적입니다. 비동기 상태 접근은 상태 백엔드가 큰 상태 크기를 처리하고 필요할 때 원격 파일시스템으로 spill할 수 있게 하는 데 필수적입니다. 이를 '분리된 상태 관리(disaggregated state management)'라고 합니다. 자세한 내용은 Disaggregated State Management를 참조하세요.

출처: 문서

본문

이 섹션에서는 상태 있는 프로그램을 작성하기 위해 Flink가 제공하는 새로운 API에 대해 배웁니다. 개념에 대해서는 Stateful Stream Processing을 참조하세요.

새 상태 API는 이전 API보다 더 유연하게 설계되었습니다. 사용자는 비동기 상태 연산을 수행할 수 있어 더 강력하고 더 효율적입니다. 비동기 상태 접근은 상태 백엔드가 큰 상태 크기를 처리하고 필요할 때 원격 파일시스템으로 spill할 수 있게 하는 데 필수적입니다. 이를 '분리된 상태 관리'라고 합니다. 자세한 내용은 Disaggregated State Management를 참조하세요.

Keyed DataStream

keyed state를 사용하려면 먼저 상태(그리고 스트림 자체의 레코드)를 파티셔닝하는 데 사용할 키를 DataStream에 지정해야 합니다. Java API에서 DataStreamkeyBy(KeySelector)를 사용해 키를 지정할 수 있습니다. 이는 keyed state를 사용하는 연산을 허용하는 KeyedStream을 생성합니다. 비동기 상태 연산을 활성화하려면 KeyedStream에서 enableAsyncState()를 수행해야 합니다.

키 셀렉터 함수는 단일 레코드를 입력으로 받아 해당 레코드의 키를 반환합니다. 키는 어떤 타입도 될 수 있으며 반드시 결정적 계산에서 파생되어야 합니다.

Flink의 데이터 모델은 키-값 쌍에 기반하지 않습니다. 따라서 데이터셋 타입을 키와 값으로 물리적으로 패킹할 필요가 없습니다. 키는 "가상(virtual)"입니다. 즉, 그룹화 오퍼레이터를 안내하기 위해 실제 데이터에 대한 함수로 정의됩니다.

다음 예시는 객체의 필드를 단순히 반환하는 키 셀렉터 함수를 보여줍니다.

Java

// some ordinary POJO
public class WordCount {
  public String word;
  public int count;

  public String getWord() { return word; }
}
DataStream<WordCount> words = // [...]
KeyedStream<WordCount> keyed = words
  .keyBy(WordCount::getWord).enableAsyncState();

Keyed State V2 사용 (Using Keyed State V2)

이전 상태 API와 달리 새 상태 API는 비동기 상태 접근을 위해 설계되었습니다. 각 상태 유형은 동기와 비동기 두 가지 버전의 API를 제공합니다. 동기 API는 상태 접근이 완료될 때까지 기다리는 차단 API입니다. 비동기 API는 비차단이며, 상태 접근이 완료될 때 완료되는 StateFuture를 반환합니다. 그 후에 콜백이나 후속 로직이 (있다면) 호출됩니다. 비동기 API가 더 효율적이므로 가능할 때마다 사용해야 합니다. 같은 사용자 함수에서 동기와 비동기 상태 접근을 혼합하는 것은 권장하지 않습니다.

keyed state 인터페이스는 모두 현재 입력 요소의 키에 범위가 한정된 다양한 유형의 상태에 대한 접근을 제공합니다. 이는 이런 유형의 상태가 Java에서 stream.keyBy(…)로 생성할 수 있는 KeyedStream에서만 사용할 수 있음을 의미합니다. 그리고 가장 중요한 것은 keyed stream이 enableAsyncState()를 호출해 비동기 상태 접근을 위해 활성화되어야 한다는 것입니다. 새 API 집합은 enableAsyncState()가 호출된 KeyedStream에서만 사용할 수 있습니다.

이제 사용 가능한 다양한 상태 유형을 살펴보고, 프로그램에서 어떻게 사용되는지 보겠습니다. 동기 API는 원래 API와 동일하므로 여기서는 비동기 API에만 초점을 맞춥니다.

반환 값 (The Return Values)

먼저 이런 비동기 상태 접근 메서드의 반환 값에 익숙해져야 합니다.

StateFuture<T>는 상태 접근의 결과로 완료될 future입니다. 반환 타입은 T입니다. 결과를 처리하는 여러 메서드를 제공하며, 다음과 같습니다.

  • StateFuture<Void> thenAccept(Consumer<T>): 이 메서드는 상태 접근이 완료될 때 결과와 함께 호출될 Consumer를 받습니다. Consumer가 끝나면 완료될 StateFuture<Void>를 반환합니다.
  • StateFuture<R> thenApply(Function<T, R>): 이 메서드는 상태 접근이 완료될 때 결과와 함께 호출될 Function을 받습니다. 함수의 반환 값은 Function이 끝나면 완료될 다음 StateFuture의 결과가 됩니다.
  • StateFuture<R> thenCompose(Function<T, StateFuture<R>>): 이 메서드는 상태 접근이 완료될 때 결과와 함께 호출될 Function을 받습니다. 함수의 반환 값은 thenCompose의 호출자에게 반환 값으로 노출되는 StateFuture<R>이어야 합니다. StateFuture<R>Function의 내부 StateFuture<R>이 끝나면 완료됩니다.
  • StateFuture<R> thenCombine(StateFuture<U>, BiFunction<T, U, R>): 이 메서드는 다른 StateFuture와, 둘 다 완료될 때 두 StateFuture의 결과로 호출될 BiFunction을 받습니다. BiFunction의 반환 값은 BiFunction이 끝나면 완료될 다음 StateFuture의 결과가 됩니다.

이 메서드들은 CompletableFuture의 해당 메서드들과 유사합니다. 이 메서드들 외에도 StateFuture는 현재 스레드를 상태 접근이 끝날 때까지 차단하는 get() 메서드를 제공하지 않는다는 점을 기억하세요. 이는 현재 스레드를 차단하면 재귀적 차단이 발생할 수 있기 때문입니다. StateFuture<T>는 또한 thenAccept, thenApply, thenCompose, thenCombine의 조건부 버전도 제공하는데, 이는 상태 접근이 완료되고 후속 로직이 상태 접근의 결과에 따라 두 갈래로 나뉘는 경우를 위한 것입니다. 이 메서드들의 조건부 버전은 thenConditionallyAccept, thenConditionallyApply, thenConditionallyCompose, thenConditionallyCombine입니다.

StateIterator<T>는 상태의 요소를 반복하는 데 사용할 수 있는 iterator입니다. 다음 메서드를 제공합니다.

  • boolean isEmpty(): iterator에 요소가 없으면 true를, 그렇지 않으면 false를 반환하는 동기 메서드입니다.
  • StateFuture<Void> onNext(Consumer<T>): 이 메서드는 상태 접근이 완료될 때 다음 요소와 함께 호출될 Consumer를 받습니다. Consumer가 끝나면 완료될 StateFuture<Void>를 반환합니다. 또한 onNext의 함수 버전인 StateFuture<Collection<R>> onNext(Function<T, R>)도 제공됩니다. 이 메서드는 상태 접근이 완료될 때 다음 요소와 함께 호출될 Function을 받습니다. 함수의 반환 값은 Function이 끝나면 완료될 다음 StateFuture의 컬렉션으로 수집되어 반환됩니다.

또한 StateFuture들을 처리하는 몇 가지 유틸리티 메서드를 포함하는 StateFutureUtils 클래스도 제공합니다. 메서드들은 다음과 같습니다.

  • StateFuture<T> completedFuture(T): 이 메서드는 주어진 값으로 완료된 StateFuture를 반환합니다. thenCompose 메서드에서 후속 처리를 위해 상수 값을 반환하고 싶을 때 유용합니다.
  • StateFuture<Void> completedVoidFuture(): 이 메서드는 null 값으로 완료된 StateFuture를 반환합니다. completedFuture의 void 값 버전입니다.
  • StateFuture<Collection<T>> combineAll(Collection<StateFuture<T>>): 이 메서드는 StateFuture 컬렉션을 받아 모든 입력 StateFuture가 완료될 때 완료되는 StateFuture를 반환합니다. 반환된 StateFuture의 결과는 입력 StateFuture들의 결과 컬렉션입니다. 여러 StateFuture의 결과를 결합하려 할 때 유용합니다.
  • StateFuture<Iterable<T>> toIterable(StateFuture<StateIterator<T>>): 이 메서드는 StateIteratorStateFuture를 받아 IterableStateFuture를 반환합니다. 반환된 StateFuture의 결과는 StateIterator의 모든 요소를 포함하는 Iterable입니다. StateIteratorIterable로 변환하려 할 때 유용합니다. 지연 로딩(lazy loading) 기능을 비활성화할 수 있으므로 그럴 이유는 없지만, 후속 계산이 iterator의 전체 데이터에 의존할 때만 유용합니다.

상태 프리미티브 (State Primitives)

사용 가능한 상태 프리미티브는 다음과 같습니다.

  • ValueState<T>: 갱신하고 검색할 수 있는 값을 유지합니다(위에서 언급한 것처럼 입력 요소의 키에 범위가 한정되므로, 연산이 보는 각 키에 대해 값이 하나씩 있을 수 있습니다). 값은 asyncUpdate(T)로 설정하고 StateFuture<T> asyncValue()로 검색할 수 있습니다.
  • ListState<T>: 요소 목록을 유지합니다. 요소를 추가하고 현재 저장된 모든 요소에 대한 StateIterator를 검색할 수 있습니다. 요소는 asyncAdd(T)asyncAddAll(List<T>)로 추가되며, Iterable은 StateFuture<StateIterator<T>> asyncGet()으로 검색할 수 있습니다. asyncUpdate(List<T>)로 기존 목록을 덮어쓸 수도 있습니다.
  • ReducingState<T>: 상태에 추가된 모든 값의 집계를 나타내는 단일 값을 유지합니다. 인터페이스는 ListState와 유사하지만 asyncAdd(T)로 추가된 요소는 지정된 ReduceFunction으로 집계로 줄어듭니다.
  • AggregatingState<IN, OUT>: 상태에 추가된 모든 값의 집계를 나타내는 단일 값을 유지합니다. ReducingState와 달리 집계 타입은 상태에 추가되는 요소의 타입과 다를 수 있습니다. 인터페이스는 ListState와 동일하지만 asyncAdd(IN)으로 추가된 요소는 지정된 AggregateFunction으로 집계됩니다.
  • MapState<UK, UV>: 매핑 목록을 유지합니다. 상태에 키-값 쌍을 넣고 현재 저장된 모든 매핑에 대한 StateIterator를 검색할 수 있습니다. 매핑은 asyncPut(UK, UV)asyncPutAll(Map<UK, UV>)로 추가됩니다. 사용자 키와 연관된 값은 asyncGet(UK)로 검색할 수 있습니다. 매핑, 키, 값에 대한 iterable 뷰는 각각 asyncEntries(), asyncKeys(), asyncValues()로 검색할 수 있습니다. 또한 asyncIsEmpty()로 이 맵이 키-값 매핑을 포함하는지 확인할 수 있습니다.

모든 상태 유형에는 현재 활성 키, 즉 입력 요소의 키에 대한 상태를 지우는 asyncClear() 메서드도 있습니다.

이 상태 객체들은 상태와의 인터페이스에만 사용된다는 점을 기억하는 것이 중요합니다. 상태는 반드시 내부에 저장되지는 않으며 디스크나 다른 곳에 있을 수 있습니다. 두 번째로 기억할 점은 상태에서 얻는 값이 입력 요소의 키에 의존한다는 것입니다. 따라서 사용자 함수의 한 호출에서 얻는 값은 관련된 키가 다르면 다른 호출의 값과 다를 수 있습니다.

상태 핸들을 얻으려면 StateDescriptor를 만들어야 합니다. 이는 상태의 이름(나중에 보겠지만 여러 상태를 만들 수 있고 참조할 수 있도록 고유한 이름을 가져야 합니다), 상태가 보유하는 값의 타입, 그리고 가능하면 ReduceFunction 같은 사용자 지정 함수를 보유합니다. 검색하려는 상태 타입에 따라 ValueStateDescriptor, ListStateDescriptor, AggregatingStateDescriptor, ReducingStateDescriptor, MapStateDescriptor 중 하나를 만듭니다. 이전 상태 API와 구분하기 위해 org.apache.flink.api.common.state.v2 패키지( v2 에 주목) 아래의 StateDescriptor를 사용해야 합니다.

상태는 RuntimeContext를 사용해 접근하므로 rich functions에서만 가능합니다. 자세한 내용은 여기를 참조하세요. RichFunction에서 사용 가능한 RuntimeContext에는 상태에 접근하는 다음 메서드가 있습니다.

  • ValueState<T> getState(ValueStateDescriptor<T>)
  • ReducingState<T> getReducingState(ReducingStateDescriptor<T>)
  • ListState<T> getListState(ListStateDescriptor<T>)
  • AggregatingState<IN, OUT> getAggregatingState(AggregatingStateDescriptor<IN, ACC, OUT>)
  • MapState<UK, UV> getMapState(MapStateDescriptor<UK, UV>)

이것은 모든 부분이 어떻게 맞물리는지 보여주는 FlatMapFunction 예시입니다.

Java

public class CountWindowAverage extends RichFlatMapFunction<Tuple2<Long, Long>, Tuple2<Long, Long>> {

    /**
     * The ValueState handle. The first field is the count, the second field a running sum.
     */
    private transient ValueState<Tuple2<Long, Long>> sum;

    @Override
    public void flatMap(Tuple2<Long, Long> input, Collector<Tuple2<Long, Long>> out) throws Exception {

        // access the state value
        sum.asyncValue().thenApply(currentSum -> {
            // if it hasn't been used before, it will be null
            Tuple2<Long, Long> current = currentSum == null ? Tuple2.of(0L, 0L) : currentSum;

            // update the count
            current.f0 += 1;

            // add the second field of the input value
            current.f1 += input.f1;

            return current;
        }).thenAccept(r -> {
            // if the count reaches 2, emit the average and clear the state
            if (r.f0 >= 2) {
                out.collect(Tuple2.of(input.f0, r.f1 / r.f0));
                sum.asyncClear();
            } else {
                sum.asyncUpdate(r);
            }
        });
    }

    @Override
    public void open(OpenContext ctx) {
        ValueStateDescriptor<Tuple2<Long, Long>> descriptor =
                new ValueStateDescriptor<>(
                        "average", // the state name
                        TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {})); // type information
        sum = getRuntimeContext().getState(descriptor);
    }
}

// this can be used in a streaming program like this (assuming we have a StreamExecutionEnvironment env)
env.fromElements(Tuple2.of(1L, 3L), Tuple2.of(1L, 5L), Tuple2.of(1L, 7L), Tuple2.of(1L, 4L), Tuple2.of(1L, 2L))
        .keyBy(value -> value.f0)
        .enableAsyncState()
        .flatMap(new CountWindowAverage())
        .print();

// the printed output will be (1,4) and (1,5)

이 예시는 poor man's counting window를 구현합니다. 튜플을 첫 번째 필드로 키 지정합니다(예시에서는 모두 같은 키 1을 가짐). 함수는 ValueState에 개수와 진행 중인 합계를 저장합니다. 개수가 2에 도달하면 평균을 내보내고 상태를 지워 0에서 다시 시작합니다. 첫 번째 필드에 다른 값이 있는 튜플이 있었다면 각각 다른 입력 키에 대해 다른 상태 값이 유지된다는 점에 유의하세요.

실행 순서 (Execution Order)

상태 접근 메서드는 비동기로 실행됩니다. 이는 상태 접근 메서드가 현재 스레드를 차단하지 않음을 의미합니다. 동기 API에서는 상태 접근 메서드가 호출된 순서대로 실행됩니다. 하지만 비동기 API에서는 상태 접근 메서드가 순서 없이 실행되며, 특히 다른 들어오는 요소에 대한 호출에서 그렇습니다. 위 예시에서 flatMap 함수가 두 개의 다른 들어오는 요소 A와 B에 대해 호출되면, A와 B에 대한 상태 접근 메서드가 실행됩니다. 먼저 A에 대해 asyncGet가 실행되고, 그 다음 B에 대해 asyncGet가 실행될 수 있습니다. 두 asyncGet의 완료 순서는 보장되지 않습니다. 따라서 두 StateFuture의 연속 실행 순서도 보장되지 않습니다. 따라서 A와 B에 대한 asyncClearasyncUpdate의 호출은 결정되지 않습니다.

상태 접근 메서드가 순서 없이 실행되지만, 이는 모든 사용자 코드가 병렬로 실행된다는 것을 의미하지는 않습니다. 상태 접근 메서드 다음의 processElement, flatMap 또는 thenXXxx 메서드의 사용자 코드는 단일 스레드(작업 스레드)에서 실행됩니다. 따라서 사용자 코드에는 동시성 문제가 없습니다.

일반적으로 상태 접근 메서드의 실행 순서에 대해 걱정할 필요는 없지만, Flink가 보장하는 몇 가지 규칙이 여전히 있습니다.

  • 같은 키의 요소에 대한 사용자 코드 진입 flatMap의 실행 순서는 요소 도착 순서대로 엄격히 호출됩니다.
  • thenXXxx 메서드에 전달된 소비자나 함수는 체인된 순서대로 실행됩니다. 체인되지 않았거나 여러 체인이 있으면 순서가 보장되지 않습니다.

비동기 API의 모범 사례 (Best practice of asynchronous APIs)

비동기 API는 동기 API보다 더 효율적이고 더 강력하게 설계되었습니다. 비동기 API를 사용할 때 따르면 좋은 몇 가지 모범 사례가 있습니다.

  • 동기와 비동기 상태 접근을 혼합하지 마세요.
  • thenXXxx 메서드의 체이닝을 사용해 상태 접근의 결과를 처리하고 다른 상태 접근이나 결과를 제공하세요. 로직을 thenXXxx 메서드로 구분된 여러 단계로 나누세요.
  • 사용자 함수(RichFlatMapFunction)의 가변 멤버에 접근하는 것을 피하세요. 상태 접근 메서드가 순서 없이 실행되므로 가변 멤버가 예측할 수 없는 순서로 접근될 수 있습니다. 대신 상태 접근의 결과를 사용해 다른 단계 사이에서 데이터를 전달하세요. StateFutureUtils.completedFuturethenApply 메서드를 사용해 데이터를 전달할 수 있습니다. 또는 각 flatMap 호출에 대해 초기화되고 람다 사이에서 공유되는 캡처된 컨테이너(AtomicReference)를 사용할 수 있습니다.

상태 Time-To-Live (TTL)

time-to-live(TTL)은 어떤 타입의 keyed state에도 할당할 수 있습니다. TTL이 구성되고 상태 값이 만료되면 저장된 값은 아래에서 더 자세히 논의하는 best effort 방식으로 정리됩니다.

모든 상태 컬렉션 유형은 항목별 TTL을 지원합니다. 이는 목록 요소와 맵 항목이 독립적으로 만료된다는 뜻입니다.

상태 TTL을 사용하려면 먼저 StateTtlConfig 구성 객체를 만들어야 합니다. 그런 다음 구성을 전달해 어떤 상태 디스크립터에서든 TTL 기능을 활성화할 수 있습니다.

Java

import org.apache.flink.api.common.state.StateTtlConfig;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import java.time.Duration;

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Duration.ofSeconds(1))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .build();

ValueStateDescriptor<String> stateDescriptor = new ValueStateDescriptor<>("text state", String.class);
stateDescriptor.enableTimeToLive(ttlConfig);

구성에는 고려할 여러 옵션이 있습니다.

newBuilder 메서드의 첫 번째 파라미터는 필수이며 time-to-live 값입니다.

update type은 상태 TTL이 언제 새로고침되는지 구성합니다(기본적으로 OnCreateAndWrite).

  • StateTtlConfig.UpdateType.OnCreateAndWrite - 생성과 쓰기 접근 시에만
  • StateTtlConfig.UpdateType.OnReadAndWrite - 읽기 접근 시에도

(참고: 동시에 상태 표시성을 StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp으로 설정하면 상태 읽기 캐시가 비활성화되어 PyFlink에서 일부 성능 손실이 발생합니다.)

상태 표시성은 아직 정리되지 않은 경우 읽기 접근 시 만료된 값이 반환되는지 여부를 구성합니다(기본적으로 NeverReturnExpired).

  • StateTtlConfig.StateVisibility.NeverReturnExpired - 만료된 값은 절대 반환되지 않습니다.

(참고: 상태 읽기/쓰기 캐시가 비활성화되어 PyFlink에서 일부 성능 손실이 발생합니다.)

  • StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp - 여전히 사용 가능하면 반환됩니다.

NeverReturnExpired의 경우, 만료된 상태는 여전히 제거되어야 하더라도 더 이상 존재하지 않는 것처럼 동작합니다. 이 옵션은 데이터가 TTL 이후 읽기 접근에 엄격히 사용 불가능해야 하는 경우, 예를 들어 개인정보 보호에 민감한 데이터를 다루는 애플리케이션에 유용할 수 있습니다.

ReturnExpiredIfNotCleanedUp 옵션은 정리 전에 만료된 상태를 반환할 수 있게 합니다.

참고:

  • 상태 백엔드는 마지막 수정 타임스탬프를 사용자 값과 함께 저장하므로, 이 기능을 활성화하면 상태 저장소의 소비가 증가합니다. Heap 상태 백엔드는 사용자 상태 객체에 대한 참조와 원시 long 값이 있는 추가 Java 객체를 메모리에 저장합니다. RocksDB/ForSt 상태 백엔드는 저장된 값, 목록 항목 또는 맵 항목마다 8바이트를 추가합니다.
  • 현재는 processing time에 대한 TTL만 지원됩니다.
  • TTL 없이 이전에 구성된 상태를 TTL 활성화 디스크립터로 복원하거나 그 반대의 경우는 호환성 실패와 StateMigrationException을 초래합니다.
  • TTL 구성은 체크포인트나 savepoint의 일부가 아니라, 현재 실행 중인 작업에서 Flink가 상태를 다루는 방식입니다.
  • TTL을 짧은 값에서 긴 값으로 조정해 체크포인트 상태를 복원하는 것은 권장되지 않으며, 잠재적인 데이터 오류를 일으킬 수 있습니다.
  • TTL이 있는 맵 상태는 사용자 값 직렬 변환기가 null 값을 처리할 수 있는 경우에만 null 사용자 값을 지원합니다. 직렬 변환기가 null 값을 지원하지 않으면 직렬화된 형식에 추가 바이트라는 비용으로 NullableSerializer로 감쌀 수 있습니다.
  • TTL 활성화 구성에서는 이미 더 이상 사용되지 않는 StateDescriptordefaultValue가 더 이상 효과가 없습니다. 이는 의미를 더 명확하게 하고, 상태 내용이 null이거나 만료된 경우 사용자가 기본 값을 수동으로 관리하게 하기 위함입니다.
만료된 상태 정리 (Cleanup of Expired State)

기본적으로 만료된 값은 ValueState#value 같은 읽기 시에 명시적으로 제거되고, 구성된 상태 백엔드가 지원하면 백그라운드에서 주기적으로 가비지 컬렉션됩니다. 백그라운드 정리는 StateTtlConfig에서 비활성화할 수 있습니다.

Java

import org.apache.flink.api.common.state.StateTtlConfig;
StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Duration.ofSeconds(1))
    .disableCleanupInBackground()
    .build();

현재 ForSt 상태 백엔드는 컴팩션 과정에서만 상태를 정리합니다. ForSt가 아닌 다른 상태 백엔드를 사용하고 다른 정리 전략에 대해 읽으려면 만료된 상태 정리에 대한 State V1 문서를 참조하세요.

컴팩션 중 정리 (Cleanup during compaction)

ForSt 상태 백엔드를 사용하면 Flink 특정 컴팩션 필터가 백그라운드 정리를 위해 호출됩니다. ForSt는 상태 업데이트를 병합하고 저장소를 줄이기 위해 주기적으로 비동기 컴팩션을 실행합니다. Flink 컴팩션 필터는 TTL이 있는 상태 항목의 만료 타임스탬프를 확인하고 만료된 값을 제외합니다.

이 기능은 StateTtlConfig에서 구성할 수 있습니다.

Java

import org.apache.flink.api.common.state.StateTtlConfig;

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Duration.ofSeconds(1))
    .cleanupInRocksdbCompactFilter(1000, Duration.ofHours(1))
    .build();

ForSt 컴팩션 필터는 특정 수의 상태 항목을 처리한 후 매번 Flink에서 현재 타임스탬프를 조회해 만료를 확인합니다. 이를 변경하고 StateTtlConfig.newBuilder(...).cleanupInRocksdbCompactFilter(long queryTimeAfterNumEntries) 메서드에 사용자 정의 값을 전달할 수 있습니다. 타임스탬프를 더 자주 업데이트하면 정리 속도가 향상될 수 있지만, 네이티브 코드에서 JNI 호출을 사용하므로 컴팩션 성능이 저하됩니다. ForSt 백엔드의 기본 백그라운드 정리는 1000개 항목이 처리될 때마다 현재 타임스탬프를 조회합니다.

주기적 컴팩션은 만료된 상태 항목 정리를 가속화할 수 있으며, 특히 거의 접근되지 않는 상태 항목에서 그렇습니다. 이 값보다 오래된 파일은 컴팩션을 위해 선택되고 이전과 같은 레벨로 다시 작성됩니다. 이는 파일이 주기적으로 컴팩션 필터를 거치도록 보장합니다. 이를 변경하고 StateTtlConfig.newBuilder(...).cleanupInRocksdbCompactFilter(long queryTimeAfterNumEntries, Duration periodicCompactionTime) 메서드에 사용자 정의 값을 전달할 수 있습니다. 주기적 컴팩션의 기본값은 30일입니다. 0으로 설정해 주기적 컴팩션을 끄거나 작은 값으로 설정해 만료된 상태 항목 정리를 가속화할 수 있지만, 더 많은 컴팩션이 트리거됩니다.

FlinkCompactionFilter에 대해 debug 레벨을 활성화해 ForSt 필터의 네이티브 코드에서 디버그 로그를 활성화할 수 있습니다.

log4j.logger.org.forstdb.FlinkCompactionFilter=DEBUG

참고:

  • 컴팩션 중 TTL 필터 호출은 컴팩션을 느리게 합니다. TTL 필터는 컴팩션되는 각 키에 대해 모든 저장된 상태 항목의 마지막 접근 타임스탬프를 파싱하고 만료를 확인해야 합니다. 컬렉션 상태 유형(목록 또는 맵)의 경우 저장된 요소마다 검사가 호출됩니다.
  • 이 기능을 고정 바이트 길이가 아닌 요소를 가진 목록 상태와 함께 사용하면, 네이티브 TTL 필터는 최소 첫 요소가 만료된 각 상태 항목에 대해 JNI로 요소의 Flink java 타입 직렬 변환기를 추가로 호출해 다음 만료되지 않은 요소의 오프셋을 결정해야 합니다.
  • 기존 작업의 경우 이 정리 전략은 StateTtlConfig에서 언제든지(예: savepoint에서 재시작한 후) 활성화하거나 비활성화할 수 있습니다.
  • 주기적 컴팩션은 TTL이 활성화된 경우에만 작동할 수 있습니다.

Operator State

Operator State(또는 non-keyed state)는 하나의 병렬 오퍼레이터 인스턴스에 바인딩된 상태입니다. Kafka Connector는 Flink에서 Operator State 사용에 대한 좋은 동기 부여 예시입니다. Kafka 소비자의 각 병렬 인스턴스는 Operator State로 토픽 파티션과 오프셋의 맵을 유지합니다.

Operator State 인터페이스는 병렬도가 변경될 때 병렬 오퍼레이터 인스턴스 사이에서 상태를 재분배하는 것을 지원합니다. 이 재분배를 수행하는 방법에는 여러 방식이 있습니다.

일반적인 상태 있는 Flink 애플리케이션에서는 operator state가 필요하지 않습니다. 그것은 대부분 소스/싱크 구현과 상태가 파티셔닝될 키가 없는 시나리오에서 사용되는 특수한 상태 유형입니다.

Broadcast State

Broadcast State는 특수한 Operator State 유형입니다. 하나의 스트림의 레코드가 모든 다운스트림 작업에 브로드캐스트되어 모든 서브태스크에서 같은 상태를 유지하는 데 사용되는 사용 사례를 지원하기 위해 도입되었습니다. 그런 다음 이 상태는 두 번째 스트림의 레코드를 처리하는 동안 접근할 수 있습니다. broadcast state가 자연스러운 예로, 다른 스트림에서 오는 모든 요소에 대해 평가하려는 규칙 집합을 포함하는 저처리량 스트림을 상상할 수 있습니다. 위 유형의 사용 사례를 염두에 두면, broadcast state는 다음과 같은 점에서 다른 operator states와 다릅니다.

  • 맵 형식을 가집니다.
  • 입력으로 broadcasted 스트림과 non-broadcasted 스트림을 가지는 특정 오퍼레이터에서만 사용할 수 있습니다.
  • 그러한 오퍼레이터는 이름이 다른 multiple broadcast states를 가질 수 있습니다.

Operator State 사용 (Using Operator State)

operator state를 사용하려면 상태 있는 함수가 CheckpointedFunction 인터페이스를 구현할 수 있습니다.

CheckpointedFunction

CheckpointedFunction 인터페이스는 다른 재분배 방식을 가진 non-keyed state에 대한 접근을 제공합니다. 두 메서드의 구현이 필요합니다.

void snapshotState(FunctionSnapshotContext context) throws Exception;

void initializeState(FunctionInitializationContext context) throws Exception;

체크포인트가 수행되어야 할 때마다 snapshotState()가 호출됩니다. 반대인 initializeState()는 사용자 정의 함수가 초기화될 때마다 호출됩니다. 함수가 처음 초기화될 때이든 이전 체크포인트에서 실제로 복구될 때이든 관계없이 호출됩니다. 이를 고려하면 initializeState()는 다양한 유형의 상태가 초기화되는 장소일 뿐만 아니라 상태 복구 로직이 포함되는 장소이기도 합니다.

현재 목록 스타일 operator state가 지원됩니다. 상태는 직렬화 가능한 객체의 List로 예상되며, 서로 독립적이므로 재스케일링 시 재분배에 적합합니다. 즉, 이 객체들은 non-keyed state가 재분배될 수 있는 가장 세밀한 단위입니다. 상태 접근 메서드에 따라 다음과 같은 재분배 방식이 정의됩니다.

  • 균등 분할 재분배(Even-split redistribution): 각 오퍼레이터는 상태 요소의 List를 반환합니다. 전체 상태는 논리적으로 모든 리스트의 연결(concatenation)입니다. 복원/재분배 시 목록은 병렬 오퍼레이터 수만큼의 하위 목록으로 균등하게 나뉩니다. 각 오퍼레이터는 비어 있을 수 있거나 하나 이상의 요소를 포함할 수 있는 하위 목록을 가집니다. 예를 들어 병렬도 1에서 오퍼레이터의 체크포인트 상태가 요소 element1element2를 포함한다면, 병렬도를 2로 늘릴 때 element1은 오퍼레이터 인스턴스 0에, element2는 오퍼레이터 인스턴스 1로 갈 수 있습니다.
  • Union 재분배: 각 오퍼레이터는 상태 요소의 List를 반환합니다. 전체 상태는 논리적으로 모든 리스트의 연결입니다. 복원/재분배 시 각 오퍼레이터는 완전한 상태 요소 목록을 가집니다. 목록의 카디널리티가 높을 수 있다면 이 기능을 사용하지 마세요. 체크포인트 메타데이터는 각 목록 항목에 대한 오프셋을 저장하며, 이는 RPC framesize나 메모리 부족 오류로 이어질 수 있습니다.

다음은 CheckpointedFunction을 사용해 외부 세계로 보내기 전에 요소를 버퍼링하는 상태 있는 SinkFunction의 예시입니다. 기본 균등 분할 재분배 목록 상태를 보여줍니다.

Java

public class BufferingSink
        implements SinkFunction<Tuple2<String, Integer>>,
                   CheckpointedFunction {

    private final int threshold;

    private transient ListState<Tuple2<String, Integer>> checkpointedState;

    private List<Tuple2<String, Integer>> bufferedElements;

    public BufferingSink(int threshold) {
        this.threshold = threshold;
        this.bufferedElements = new ArrayList<>();
    }

    @Override
    public void invoke(Tuple2<String, Integer> value, Context context) throws Exception {
        bufferedElements.add(value);
        if (bufferedElements.size() >= threshold) {
            for (Tuple2<String, Integer> element: bufferedElements) {
                // send it to the sink
            }
            bufferedElements.clear();
        }
    }

    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception {
        checkpointedState.update(bufferedElements);
    }

    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {
        ListStateDescriptor<Tuple2<String, Integer>> descriptor =
            new ListStateDescriptor<>(
                "buffered-elements",
                TypeInformation.of(new TypeHint<Tuple2<String, Integer>>() {}));

        checkpointedState = context.getOperatorStateStore().getListState(descriptor);

        if (context.isRestored()) {
            for (Tuple2<String, Integer> element : checkpointedState.get()) {
                bufferedElements.add(element);
            }
        }
    }
}

initializeState 메서드는 인자로 FunctionInitializationContext를 받습니다. 이는 non-keyed state "컨테이너"를 초기화하는 데 사용됩니다. 이들은 non-keyed state 객체가 체크포인팅 시에 저장될 ListState 타입의 컨테이너입니다.

상태가 keyed state와 유사하게 상태 이름과 상태가 보유한 값의 타입에 대한 정보를 포함하는 StateDescriptor로 초기화되는 방법을 주목하세요.

Java

ListStateDescriptor<Tuple2<String, Integer>> descriptor =
    new ListStateDescriptor<>(
        "buffered-elements",
        TypeInformation.of(new TypeHint<Tuple2<String, Integer>>() {}));

checkpointedState = context.getOperatorStateStore().getListState(descriptor);

상태 접근 메서드의 명명 규칙은 재분배 패턴 다음에 상태 구조를 포함합니다. 예를 들어 복원 시 union 재분배 방식으로 목록 상태를 사용하려면 getUnionListState(descriptor)로 상태에 접근합니다. 메서드 이름에 재분배 패턴이 포함되지 않으면(예: getListState(descriptor)) 기본 균등 분할 재분배 방식이 사용됨을 의미합니다.

컨테이너를 초기화한 후, 컨텍스트의 isRestored() 메서드를 사용해 실패 후 복구 중인지 확인합니다. true이면, 즉 복구 중이면 복원 로직이 적용됩니다.

수정된 BufferingSink 코드에서 보여주듯이 상태 초기화 중에 복구된 이 ListStatesnapshotState()에서 나중에 사용하기 위해 클래스 변수에 유지됩니다. 거기서 ListState는 이전 체크포인트가 포함한 모든 객체에서 지워진 다음 체크포인트하려는 새 객체들로 채워집니다.

참고로 keyed state는 initializeState() 메서드에서도 초기화할 수 있습니다. 이는 제공된 FunctionInitializationContext를 사용해 수행할 수 있습니다.

옛 상태 API에서 마이그레이션 (Migrate from the Old State API)

옛 상태 API에서 새 것으로 마이그레이션하는 것은 매우 쉽습니다. 다음 단계를 따르세요.

  • KeyedStream에서 enableAsyncState()를 호출해 새 상태 API를 활성화합니다.
  • StateDescriptorv2 패키지 아래의 새 것으로 바꿉니다. 또한 옛 상태 핸들을 v2 패키지 아래의 새 것으로 바꿉니다.
  • 옛 상태 접근 메서드를 새 비동기 메서드로 다시 작성합니다.
  • 새 상태 API에는 비동기 상태 접근을 수행할 수 있는 ForSt State Backend를 사용할 것을 권장합니다. 다른 상태 백엔드는 새 상태 API와 함께 사용될 수 있지만 상태 접근의 동기 실행만 지원합니다.

더 알아보기 (Learn more)