사용자 정의 함수

사용자 정의 함수 (Working with State)

상태 기반 프로그램을 작성하기 위한 Flink의 API를 배우는 절이에요. 상태 기반 스트림 처리의 개념이 궁금하다면 Stateful Stream Processing 문서를 먼저 보는 게 좋아요. 여기서는 구체적으로 키 상태(Keyed State) API부터, 연산자 상태(Operator State), 상태 TTL까지 실제 코드로 어떻게 쓰는지 다뤄요.

출처: Apache Flink 공식 문서 - Working with State

키 데이터 스트림 (Keyed DataStream)

키 상태를 쓰려면 먼저 DataStream에 상태(그리고 스트림의 레코드)를 파티셔닝할 키를 지정해야 해요. Java API에서는 keyBy(KeySelector), Python API에서는 key_by(KeySelector)DataStream에 호출하면 돼요. 그러면 KeyedStream이 나오고, 이걸로 키 상태를 쓰는 연산을 할 수 있어요.

키 셀렉터 함수는 단일 레코드를 입력으로 받아 그 레코드의 키를 반환해요. 키는 어떤 타입이든 될 수 있고 결정적( deterministic) 계산에서 나와야 해요. Flink의 데이터 모델은 키-값 쌍 기반이 아니어서, 데이터를 물리적으로 키와 값으로 묶을 필요는 없어요. 키는 "가상"이며 데이터 위에 정의되는 함수로 그룹핑 연산을 안내하지요.

// 평범한 POJO
public class WC {
    public String word;
    public int count;
    public String getWord() { return word; }
}

DataStream<WC> words = // [...]
KeyedStream<WC> keyed = words.keyBy(WC::getWord);
words = [...] # type: DataStream[Row]
keyed = words.key_by(lambda row: row[0])

Java API에는 키를 정의하는 대안으로 튜플 키(tuple key)와 표현식 키(expression key)도 있지만, 오늘날에는 권장하지 않아요. KeySelector 함수가 더 우월한데, Java 람다로 쓰기 쉽고 런타임 오버헤드도 잠재적으로 적어요.

키 상태 사용 (Using Keyed State)

키 상태 인터페이스는 현재 입력 요소의 키에 범위가 한정되는 여러 상태 타입에 접근하게 해줘요. 이 상태 타입은 KeyedStream(즉 stream.keyBy(...)/stream.key_by(...)로 만든 스트림)에서만 쓸 수 있어요. 사용 가능한 상태 원시 타입(primitive)은 다음과 같아요.

  • ValueState<T>: 갱신·조회할 수 있는 값 하나를 보관해요. update(T)로 설정하고 T value()로 조회해요.
  • ListState<T>: 요소들의 리스트를 보관해요. add(T)/addAll(List<T>)로 추가하고 Iterable<T> get()으로 조회하며, update(List<T>)로 기존 리스트를 덮어쓸 수 있어요.
  • ReducingState<T>: 상태에 추가된 모든 값의 집계를 나타내는 값 하나를 보관해요. ListState와 인터페이스가 비슷하지만 add(T)로 추가된 요소들이 지정된 ReduceFunction으로 축소(reduce)돼요.
  • AggregatingState<IN, OUT>: 상태에 추가된 모든 값의 집계를 나타내는 값 하나를 보관해요. ReducingState와 달리 집계 타입이 상태에 추가되는 요소의 타입과 달라도 돼요. add(IN)으로 추가된 요소들이 지정된 AggregateFunction으로 집계돼요.
  • MapState<UK, UV>: 매핑들의 리스트를 보관해요. put(UK, UV)/putAll(Map<UK, UV>)로 키-값 쌍을 넣고 get(UK)로 조회해요. entries(), keys(), values()로 순회 뷰를 얻고 isEmpty()로 비어 있는지 확인해요.

모든 상태 타입에는 현재 활성 키(입력 요소의 키)의 상태를 지우는 clear() 메서드가 있어요. 이 상태 객체들은 상태와 인터페이스하기 위한 핸들이라는 점을 기억해야 해요. 상태가 반드시 그 안에 저장되는 건 아니고 디스크나 다른 곳에 있을 수도 있어요. 또 상태에서 얻는 값은 입력 요소의 키에 달라지므로, 같은 함수의 다른 호출에서 얻는 값이 키에 따라 달라질 수 있어요.

상태 핸들을 얻으려면 StateDescriptor를 만들어야 해요. 이 디스크립터는 상태 이름(여러 상태를 만들면 중복 없이 참조할 수 있도록 고유해야 해요), 상태가 보관하는 값의 타입, 그리고 선택적으로 ReduceFunction 같은 사용자 지정 함수를 담아요. ValueStateDescriptor, ListStateDescriptor, AggregatingStateDescriptor, ReducingStateDescriptor, MapStateDescriptor 중 원하는 걸 만들면 돼요.

상태는 RuntimeContext로 접근하므로 rich 함수에서만 가능해요. RichFunctionRuntimeContextgetState, getReducingState, getListState, getAggregatingState, getMapState 메서드를 제공해요.

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

public class CountWindowAverage extends RichFlatMapFunction<Tuple2<Long, Long>, Tuple2<Long, Long>> {
    /** ValueState 핸들. 첫 필드는 개수, 둘째 필드는 누적 합. */
    private transient ValueState<Tuple2<Long, Long>> sum;

    @Override
    public void flatMap(Tuple2<Long, Long> input, Collector<Tuple2<Long, Long>> out) throws Exception {
        // 상태 값 접근
        Tuple2<Long, Long> currentSum = sum.value();
        // 개수 갱신
        currentSum.f0 += 1;
        // 입력 값의 둘째 필드 추가
        currentSum.f1 += input.f1;
        // 상태 갱신
        sum.update(currentSum);
        // 개수가 2에 도달하면 평균을 내보내고 상태를 비움
        if (currentSum.f0 >= 2) {
            out.collect(new Tuple2<>(input.f0, currentSum.f1 / currentSum.f0));
            sum.clear();
        }
    }

    @Override
    public void open(OpenContext ctx) {
        ValueStateDescriptor<Tuple2<Long, Long>> descriptor =
            new ValueStateDescriptor<>(
                "average", // 상태 이름
                TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {}), // 타입 정보
                Tuple2.of(0L, 0L)); // 상태 기본값
        sum = getRuntimeContext().getState(descriptor);
    }
}

이 예시는 "서민용 카운팅 윈도우"를 구현해요. 튜플을 첫 필드로 키잉하고(예시에서는 모두 키 1), 함수가 개수와 누적 합을 ValueState에 저장해요. 개수가 2에 도달하면 평균을 내보내고 상태를 비워 0부터 다시 시작해요. 만약 첫 필드가 다른 값을 가진 튜플이 있었다면 각 입력 키마다 다른 상태 값이 유지됐을 거예요.

상태 TTL (State Time-To-Live)

어떤 타입의 키 상태든 **time-to-live(TTL)**을 부여할 수 있어요. TTL이 설정되고 상태 값이 만료되면, 저장된 값은 최선을 다해(best effort) 정리돼요. 모든 상태 컬렉션 타입이 항목별 TTL을 지원해서 리스트 요소와 맵 항목이 독립적으로 만료돼요.

TTL을 쓰려면 먼저 StateTtlConfig 설정 객체를 만들고, 아무 상태 디스크립터에 그 설정을 전달해서 활성화하면 돼요.

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))   // TTL 값 (필수)
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .build();

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

UpdateType는 상태 TTL이 언제 갱신되는지 설정해요(기본 OnCreateAndWrite): 생성·쓰기 시에만, 또는 읽기도 포함. StateVisibility는 아직 정리되지 않은 만료 값을 읽을 때 반환할지 설정해요(기본 NeverReturnExpired): 만료 값을 절대 반환하지 않거나, 아직 남아 있으면 반환. NeverReturnExpired에서는 만료 상태가 마치 더 이상 존재하지 않는 것처럼 동작해서, TTL 이후 엄격히 읽기 접근이 불가해야 하는 개인정보 민감 데이터 등에 유용해요.

여기서 알아 둘 점이 몇 가지 있어요.

  • 상태 백엔드는 사용자 값과 함께 마지막 수정 시각을 저장해서, 이 기능을 켜면 상태 저장소 소비가 늘어나요(Heap 백엔드는 참조용 자바 객체 + long 값, RocksDB 백엔드는 저장 항목당 8바이트 추가).
  • 현재는 처리 시간 기준의 TTL만 지원돼요.
  • TTL 설정은 체크포인트나 세이브포인트의 일부가 아니라, 현재 실행 중인 잡에서 Flink가 상태를 다루는 방식이에요.
  • 짧은 TTL에서 긴 TTL로 조정하며 체크포인트를 복원하는 것은 권장하지 않아요(데이터 오류 가능성).

만료 상태 정리 (Cleanup of Expired State)

기본적으로 만료 값은 ValueState#value 같은 읽기 시점에 명시적으로 제거되고, 상태 백엔드가 지원하면 백그라운드에서 주기적으로 가비지 컬렉트돼요. 백그라운드 정리를 세부 조정할 수 있는 방법으로는 전체 스냅샷 시점 정리(full snapshot에서 정리 — 스냅샷 크기 감소), 증분 정리(incremental cleanup — 상태 접근마다/레코드 처리마다 일부 항목 확인, Heap 백엔드에서만 구현), RocksDB 컴팩션 시 정리(RocksDB 컴팩션 필터가 만료 항목 제외)가 있어요. 기존 잡에서도 StateTtlConfig에서 언제든 활성화/비활성화할 수 있어요.

연산자 상태 (Operator State)

키 상태와 더불어 **연산자 상태(operator state)**도 있어요. 이 상태는 키에 연결되는 대신 연산자의 단일 병렬 인스턴스(또는 서브태스크)에 연결돼요. 연산자 상태는 CheckpointedFunction 인터페이스를 구현해서 사용해요. 병렬 재분배 방식은 even-split redistribution(전체 상태를 병렬도만큼 균등 분할)과 union redistribution(각 연산자가 전체 리스트를 받음 — 높은 카디널리티에선 사용 금지)이 있어요. 예를 들어 체크포인트된 상태에 element1, element2가 있고 병렬도를 1에서 2로 늘리면, even-split으로 element1은 인스턴스 0, element2는 인스턴스 1로 갈 수 있어요.

다음은 CheckpointedFunction으로 요소를 버퍼링한 뒤 외부로 보내는 상태 기반 SinkFunction 예시예요.

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) {
                // 외부로 전송
            }
            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를 인자로 받아 비키 상태 "컨테이너"를 초기화해요. 컨테이너는 ListState 타입으로, 체크포인트 시 비키 상태 객체가 저장되는 곳이에요. 컨테이너를 초기화한 뒤 isRestored() 메서드로 실패 후 복구 중인지 확인하고, 복구 중이면 복구 로직을 적용해요. 상태 접근 메서드의 이름 규약은 재분배 패턴 + 상태 구조의 조합이에요. 예를 들어 union 재분배로 리스트 상태에 접근하려면 getUnionListState(descriptor)를, 기본 even-split 방식을 쓰려면 getListState(descriptor)를 써요.

상태 기반 소스 함수 (Stateful Source Functions)

상태 기반 소스는 다른 연산자보다 조금 더 신경 써야 해요. 상태 갱신과 출력 수집을 원자적으로 만들기 위해(exactly-once 시맨틱스에 필요) 소스 컨텍스트에서 잠금(lock)을 얻어야 해요. ctx.getCheckpointLock()으로 잠금을 얻고, 출력과 상태 갱신을 synchronized(lock) 블록 안에서 수행해요. 또한 어떤 연산자는 체크포인트가 완전히 확인(acknowledge)된 시점을 알아야 외부 시스템과 소통할 수 있는데, 그럴 때는 org.apache.flink.api.common.state.CheckpointListener 인터페이스를 사용해요.

더 알아보기 (Learn more)