아파치 Flink의 효과적인 상태 관리와 일관성 유지 전략

스트림 처리 프레임워크에서 상태를 효율적으로 관리하는 것은 복잡한 데이터 처리 시나리오를 구현하는 데 필수적입니다. 대부분의 고급 스트림 처리 작업은 실시간으로 유입되는 데이터에 기반하여 지속적으로 업데이트되는 상태 정보를 필요로 합니다. 예를 들어, 데이터 스트림 내 중복을 제거하기 위해서는 이미 처리된 데이터를 상태로 기록하여 새로운 데이터가 유입될 때 이를 참조해야 합니다. 특정 데이터 패턴을 감지하기 위해 이전 요소들을 상태로 저장해야 할 수도 있으며, 시간 윈도우 기반의 집계 분석(예: 지난 한 시간 동안의 특정 지표 75 또는 99 퍼센타일 값 계산) 역시 상태 기능을 요구합니다.

일반적인 상태 관리 과정은 다음과 같습니다. 연산자 서브태스크는 입력 스트림을 수신하고, 관련 상태를 조회합니다. 새로운 계산 결과에 따라 상태를 갱신하고, 이 업데이트된 상태를 저장합니다. 예를 들어, 시간 윈도우 내의 특정 정수 필드를 합산하는 경우, 새로운 요소가 들어올 때 서브태스크는 저장된 현재 합계 값을 가져와 새로운 값을 더하고, 그 결과를 상태로 업데이트합니다.

상태 유형

Flink는 두 가지 기본적인 상태 유형을 제공합니다: 관리형 상태(Managed State)와 원시 상태(Raw State). 이 둘의 주요 차이점은 Flink가 상태를 관리하는 방식에 있습니다. 관리형 상태는 Flink 런타임에 의해 저장, 복구 및 최적화되지만, 원시 상태는 개발자가 직접 직렬화하고 관리해야 합니다.

  • 관리형 상태(Managed State): Flink 런타임이 관리하며, 저장 및 복구가 자동으로 이루어집니다. Flink는 저장소 관리 및 영구화에 대한 최적화를 제공하며, 애플리케이션의 병렬도가 변경될 때 상태가 자동으로 재분배됩니다. 관리형 상태는 ValueState, ListState, MapState 등과 같은 다양한 데이터 구조를 지원합니다. 대부분의 연산자는 Rich 함수 클래스를 상속하거나 제공된 인터페이스를 통해 관리형 상태를 사용할 수 있습니다.
  • 원시 상태(Raw State): 사용자 정의 상태로, 매우 낮은 수준의 바이트 배열 형태로 저장됩니다. 사용자는 데이터를 직접 직렬화해야 하며, Flink는 이 상태의 내부 구조를 알지 못합니다. 이는 기존 연산자와 관리형 상태로 충분하지 않을 때 사용자 정의 연산자를 구현하는 데 주로 사용됩니다.

관리형 상태는 다시 두 가지 세부 유형으로 나뉩니다: 키 기반 상태(Keyed State)와 연산자 상태(Operator State).

사용자 정의 Flink 연산자를 만들기 위해 RichFlatMapFunction과 같은 Rich Function 인터페이스를 재정의할 수 있습니다. 키 기반 상태는 Rich Function 인터페이스 내에서 생성 및 접근할 수 있습니다. 연산자 상태의 경우, 추가적으로 CheckpointedFunction 인터페이스를 구현해야 합니다.

키 기반 상태 (Keyed State)

Flink는 각 키에 대해 별도의 상태 인스턴스를 유지합니다. 동일한 키를 가진 모든 데이터는 동일한 연산자 태스크로 파티셔닝되며, 이 태스크는 해당 키에 대한 상태를 관리하고 처리합니다. 태스크가 데이터를 처리할 때, 자동으로 현재 데이터의 키에 따라 상태 접근 범위를 한정합니다. 따라서 동일한 키를 가진 모든 데이터는 동일한 상태에 접근합니다.

키 기반 상태는 KeyedStream에서만 사용할 수 있으며, 이는 스트림에 `keyBy(...)` 연산을 적용하여 얻을 수 있습니다.

Flink는 키 기반 상태를 관리하고 저장하기 위해 다음과 같은 데이터 구조를 제공합니다:

  • ValueState<T>: 단일 값을 저장합니다. update(T)로 갱신하고 value()로 조회합니다.
  • ListState<T>: 리스트 형태의 상태를 저장합니다. add(T) 또는 addAll(List)로 요소를 추가하고, update(T)로 갱신하며, get()으로 전체 리스트를 얻습니다.
  • ReducingState<T>: ReduceFunction의 결과로 계산된 값을 저장하며, add(T)로 요소를 추가합니다.
  • AggregatingState<IN, OUT>: AggregatingState의 결과로 계산된 값을 저장하며, add(IN)으로 요소를 추가합니다.
  • MapState<UK, UV>: 맵 형태의 상태를 유지합니다. get(UK)으로 조회하고, put(UK, UV)로 갱신하며, contains(UK)로 포함 여부를 확인하고, remove(UK)로 요소를 제거합니다.
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;

import java.util.ArrayList;
import java.util.List;

public class KeyedListStateExample extends RichFlatMapFunction<Tuple2<String, Long>, List<Tuple2<String, Long>>> {

    private transient ListState<Tuple2<String, Long>> keyHistoryListState;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        ListStateDescriptor<Tuple2<String, Long>> descriptor =
                new ListStateDescriptor<>(
                        "keyHistory",
                        TypeInformation.of(new TypeHint<Tuple2<String, Long>>() {})
                );
        keyHistoryListState = getRuntimeContext().getListState(descriptor);
    }

    @Override
    public void flatMap(Tuple2<String, Long> inputTuple, Collector<List<Tuple2<String, Long>>> out) throws Exception {
        // 현재 상태에서 모든 요소를 가져와 새 리스트에 추가
        List<Tuple2<String, Long>> currentElements = new ArrayList<>();
        for (Tuple2<String, Long> element : keyHistoryListState.get()) {
            currentElements.add(element);
        }
        currentElements.add(inputTuple);

        // 상태를 업데이트하고 컬렉터로 내보내기
        keyHistoryListState.update(currentElements);
        out.collect(currentElements);
    }

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2); // 병렬 처리 설정

        DataStream<Tuple2<String, Long>> sourceStream = env.fromElements(
                Tuple2.of("userA", 10L), Tuple2.of("userA", 20L), Tuple2.of("userA", 30L),
                Tuple2.of("userB", 15L), Tuple2.of("userB", 25L), Tuple2.of("userB", 35L),
                Tuple2.of("userC", 5L), Tuple2.of("userC", 15L), Tuple2.of("userC", 25L)
        );

        sourceStream
                .keyBy(value -> value.f0) // 첫 번째 필드(String)로 키 설정
                .flatMap(new KeyedListStateExample())
                .print("Keyed State Output");

        env.execute("Keyed ListState Example");
    }
}

연산자 상태 (Operator State)

연산자 상태는 모든 연산자에서 사용할 수 있으며, 각 연산자 서브태스크 또는 연산자 인스턴스가 하나의 상태를 공유합니다. 이 서브태스크로 유입되는 모든 데이터는 해당 상태에 접근하고 업데이트할 수 있습니다. 연산자 상태는 동일하거나 다른 연산자의 다른 인스턴스에 의해 접근될 수 없습니다.

Flink는 연산자 상태를 위해 세 가지 기본적인 데이터 구조를 제공합니다:

  • ListState<T>: 리스트 형태의 상태를 저장합니다.
  • UnionListState<T>: 리스트 형태의 상태를 저장하지만, 병렬도가 변경될 때 ListState가 해당 연산자의 모든 병렬 상태 인스턴스를 통합하여 새로운 태스크에 균등하게 분배하는 것과 달리, UnionListState는 단순히 모든 병렬 상태 인스턴스를 모으고 구체적인 분할 동작은 사용자가 정의합니다.
  • BroadcastState<K, V>: 브로드캐스트용 연산자 상태입니다. 연산자에 여러 태스크가 있고 각 태스크의 상태가 모두 동일한 경우에 가장 적합합니다.

예를 들어, 모니터링 데이터 유형을 구분할 필요 없이, 특정 임계값을 초과하는 모니터링 데이터가 지정된 횟수만큼 발생하면 경보를 울리는 시나리오를 생각해 볼 수 있습니다.

import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.runtime.state.FunctionInitializationContext;
import org.apache.flink.runtime.state.FunctionSnapshotContext;
import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;

import java.util.ArrayList;
import java.util.List;

public class OperatorBufferingStateExample extends RichFlatMapFunction<Tuple2<String, Long>, List<Tuple2<String, Long>>>
        implements CheckpointedFunction {

    private final int bufferLimit;
    private transient ListState<Tuple2<String, Long>> managedBufferState;
    private List<Tuple2<String, Long>> currentBufferedElements;

    public OperatorBufferingStateExample(int limit) {
        this.bufferLimit = limit;
        this.currentBufferedElements = new ArrayList<>();
    }

    @Override
    public void flatMap(Tuple2<String, Long> value, Collector<List<Tuple2<String, Long>>> out) throws Exception {
        currentBufferedElements.add(value);
        if (currentBufferedElements.size() >= bufferLimit) {
            out.collect(new ArrayList<>(currentBufferedElements)); // 버퍼된 요소들을 내보내기
            currentBufferedElements.clear(); // 버퍼 비우기
        }
    }

    /**
     * 체크포인트 스냅샷을 생성할 때 호출됩니다.
     * @param context 스냅샷 컨텍스트
     * @throws Exception 예외 발생 시
     */
    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception {
        managedBufferState.clear(); // 이전 상태를 지우기
        for (Tuple2<String, Long> element : currentBufferedElements) {
            managedBufferState.add(element); // 현재 버퍼된 요소들을 상태에 추가
        }
    }

    /**
     * 상태를 초기화하거나 복구할 때 호출됩니다.
     * @param context 초기화 컨텍스트
     * @throws Exception 예외 발생 시
     */
    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {
        ListStateDescriptor<Tuple2<String, Long>> descriptor =
                new ListStateDescriptor<>(
                        "operatorBuffer",
                        TypeInformation.of(new TypeHint<Tuple2<String, Long>>() {})
                );
        managedBufferState = context.getOperatorStateStore().getListState(descriptor);

        // 장애로부터 복구되는 경우
        if (context.isRestored()) {
            for (Tuple2<String, Long> element : managedBufferState.get()) {
                currentBufferedElements.add(element); // 복구된 요소들을 현재 버퍼에 추가
            }
            // 복구 후에는 상태를 지우는 것이 일반적이지만, 필요에 따라 다를 수 있음.
            // 여기서는 snapshotState에서 clear를 하므로 유지.
        }
    }

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(500); // 500ms마다 체크포인트 활성화
        env.setParallelism(1); // 예시를 위해 병렬도 1로 설정

        DataStream<Tuple2<String, Long>> inputStream = env.fromElements(
                Tuple2.of("event1", 1L), Tuple2.of("event2", 2L), Tuple2.of("event3", 3L),
                Tuple2.of("event4", 4L), Tuple2.of("event5", 5L), Tuple2.of("event6", 6L),
                Tuple2.of("event7", 7L), Tuple2.of("event8", 8L), Tuple2.of("event9", 9L)
        );

        inputStream
                .flatMap(new OperatorBufferingStateExample(3)) // 3개 요소 버퍼링 후 출력
                .print("Operator State Output");

        env.execute("Operator State Buffering Example");
    }
}

상태의 수평 확장

스트림 처리 애플리케이션의 수평 확장 문제는 주로 Flink 애플리케이션의 병렬도가 변경될 때 발생합니다. 즉, 각 연산자의 병렬 인스턴스 또는 연산자 서브태스크 수가 변경되어 일부 서브태스크를 중단하거나 새로 시작해야 할 때, 기존 서브태스크에 있던 상태 데이터가 새로운 서브태스크로 원활하게 마이그레이션되어야 합니다.

Flink의 체크포인트 메커니즘은 연산자 간에 상태 데이터를 마이그레이션하는 효과적인 방법입니다. 연산자의 로컬 상태는 스냅샷으로 생성되어 HDFS와 같은 분산 저장소에 저장됩니다. 병렬도가 확장된 후, 연산자 서브태스크 수가 변경되면, 서브태스크는 재시작되고 해당 상태는 분산 저장소에서 복원됩니다.

키 기반 상태와 연산자 상태는 수평 확장 메커니즘에서 약간의 차이가 있습니다. 키 기반 상태는 항상 특정 키와 연결되어 있으므로, 수평 확장 시 키는 항상 특정 연산자 서브태스크에 자동으로 재할당됩니다. 따라서 키 기반 상태는 여러 병렬 서브태스크 간에 자동으로 마이그레이션됩니다. 반면, 비(非) 키 기반 스트림의 경우, 연산자 서브태스크로 유입되는 데이터는 병렬도 변경에 따라 달라질 수 있습니다. 예를 들어, 애플리케이션의 병렬도가 2에서 3으로 변경되거나 1로 변경되면, 데이터 스트림의 분할 방식이 변경되고 이에 따라 상태 저장 방식도 영향을 받습니다. 연산자 상태의 수평 확장 문제에 대해서는 두 가지 상태 할당 방식이 있습니다: 하나는 균등 분배 방식이고, 다른 하나는 모든 상태를 병합한 후 각 인스턴스에 분배하는 방식입니다.

체크포인트 메커니즘

Flink 상태의 뛰어난 장애 허용성을 위해 체크포인트 메커니즘이 제공됩니다. 체크포인트 메커니즘을 통해 Flink는 데이터 스트림에 주기적으로 체크포인트 배리어(barrier)를 생성합니다. 특정 연산자가 배리어를 수신하면 현재 상태를 기반으로 스냅샷을 생성한 다음, 이 배리어를 하위 연산자로 전달합니다. 하위 연산자도 배리어를 수신하면 현재 상태를 기반으로 스냅샷을 생성하고, 이 과정은 최종 싱크 연산자까지 반복됩니다. 예외 발생 시, Flink는 가장 최근의 스냅샷 데이터를 사용하여 모든 연산자를 이전 상태로 복구할 수 있습니다.

체크포인트 활성화

기본적으로 체크포인트는 비활성화되어 있습니다. StreamExecutionEnvironmentenableCheckpointing(n) 메서드를 호출하여 체크포인트를 활성화할 수 있으며, 여기서 n은 체크포인트 생성 간격(밀리초)입니다.

체크포인트는 Flink의 장애 허용 메커니즘에서 가장 핵심적인 기능입니다. 이는 구성에 따라 스트림 내 각 연산자의 상태를 기반으로 주기적으로 스냅샷을 생성하여 이 상태 데이터를 영구 저장합니다. Flink 프로그램이 예기치 않게 중단될 경우, 프로그램을 다시 실행할 때 이러한 스냅샷 중 하나를 선택하여 복구함으로써 장애로 인한 데이터 상태 중단을 수정할 수 있습니다.

체크포인트가 특정 시간 간격으로 트리거되어야 할 때마다, Flink 런타임의 분산 스트림 소스에 배리어 마커가 삽입됩니다. 이 배리어는 스트림의 데이터 레코드와 함께 하위의 각 연산자로 흐릅니다. 연산자가 배리어를 수신하면, 스트림에서 새로 수신되는 데이터 레코드 처리를 일시 중지합니다. 연산자에는 여러 입력 스트림이 있을 수 있으며, 각 스트림에는 해당 배리어가 존재합니다. 이 연산자는 모든 입력 스트림의 배리어가 도달할 때까지 기다립니다. 모든 스트림의 배리어가 연산자에 도달하면 (이는 시간상 같은 시점임을 나타냄), 연산자는 배리어보다 일찍 도착한 데이터 레코드들을 (Outgoing Records) 하위 연산자의 입력으로 방출한 후, 배리어에 해당하는 스냅샷을 이번 체크포인트의 결과 데이터로 방출합니다.

체크포인트의 다른 속성들은 다음과 같습니다:

  • 정확히 한 번(exactly-once) vs 최소 한 번(at-least-once): enableCheckpointing(long interval, CheckpointingMode mode) 메서드에 모드를 전달하여 두 가지 보장 수준 중 하나를 선택할 수 있습니다. 대부분의 애플리케이션에는 정확히 한 번이 더 좋은 선택입니다. 최소 한 번은 매우 낮은 지연 시간을 가진 애플리케이션(항상 몇 밀리초)에 더 적합할 수 있습니다.
  • 체크포인트 타임아웃: 체크포트 실행 시간이 구성된 임계값을 초과하면, 진행 중인 체크포인트 작업은 폐기됩니다.
  • 체크포인트 간 최소 시간: 이 속성은 체크포인트 사이에 필요한 최소 시간을 정의하여, 스트림 애플리케이션이 체크포인트 간에 충분한 진행을 보장하도록 합니다. 예를 들어, 값이 5000으로 설정되면, 체크포인트 지속 시간과 간격이 얼마든 이전 체크포인트가 완료된 후 최소 5초 후에 다음 체크포인트가 시작됩니다.
  • 동시 체크포인트 수: 기본적으로 이전 체크포인트가 완료되지 않은 (실패 또는 성공) 경우, 시스템은 다른 체크포인트를 트리거하지 않습니다. 이는 토폴로지가 체크포인트에 너무 많은 시간을 소비하여 정상적인 처리 흐름에 영향을 미치지 않도록 합니다. 그러나 여러 체크포인트를 병렬로 진행하는 것이 가능하며, 처리 지연이 일정하지만 (예: 시간이 오래 걸리는 외부 서비스 호출) 장애 후 재처리 파이프라인을 최소화하기 위해 빈번한 체크포인트를 원하는 경우에 유용합니다.
  • 외부화된 체크포인트 (Externalized checkpoints): 체크포인트를 외부 시스템에 주기적으로 저장하도록 구성할 수 있습니다. 외부화된 체크포인트는 메타데이터를 영구 저장소에 기록하며 작업이 실패할 때 자동으로 삭제되지 않습니다. 이런 방식으로, 작업이 실패하더라도 복구할 수 있는 기존 체크포인트가 있게 됩니다. 자세한 내용은 외부화된 체크포인트 배포 문서를 참조하세요.
  • 체크포인트 오류 시 태스크 실패 또는 계속 진행: 태스크 체크포인트 과정에서 오류가 발생할 때 태스크를 실패시킬지 아니면 계속 진행할지를 결정합니다. 기본 동작은 실패시키는 것입니다. 이 옵션을 비활성화하면, 태스크는 단순히 체크포인트 오류 정보를 체크포인트 코디네이터에 보고하고 계속 실행됩니다.
  • 체크포인트 복구 선호 (prefer checkpoint for recovery): 이 속성은 작업이 가장 최근의 세이브포인트가 사용 가능하더라도 가장 최근 체크포인트로 되돌릴지 여부를 결정합니다. 이는 잠재적으로 복구 시간을 줄일 수 있습니다 (체크포인트 복구가 세이브포인트 복구보다 빠릅니다).
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.environment.CheckpointConfig;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class FlinkCheckpointConfig {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 1000ms마다 체크포인트 시작
        env.enableCheckpointing(1000);

        // 체크포인트 고급 옵션 설정
        CheckpointConfig config = env.getCheckpointConfig();

        // 모드를 EXACTLY_ONCE (기본값)으로 설정
        config.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

        // 체크포인트 간 최소 500ms의 간격을 보장
        config.setMinPauseBetweenCheckpoints(500);

        // 체크포인트는 60초 내에 완료되어야 하며, 그렇지 않으면 폐기
        config.setCheckpointTimeout(60000);

        // 동시에 하나의 체크포인트만 진행 허용
        config.setMaxConcurrentCheckpoints(1);

        // 작업이 취소된 후에도 외부화된 체크포인트 유지 활성화
        config.enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

        // 복구 시 가장 최근 세이브포인트보다 체크포인트를 선호 (복구 시간 단축 가능)
        config.setPreferCheckpointForRecovery(true);

        // ... 스트림 처리 로직 ...
        env.execute("Flink Checkpoint Configuration Example");
    }
}

여러 체크포인트 보존

기본적으로 Flink는 체크포인트 옵션이 설정된 경우 가장 최근에 성공적으로 생성된 하나의 체크포인트만 보존합니다. Flink 프로그램이 실패하면 이 최근 체크포인트에서 복구할 수 있습니다. 그러나 여러 체크포인트를 보존하고 실제 필요에 따라 그중 하나를 선택하여 복구하는 것이 더 유연할 수 있습니다. 예를 들어, 최근 4시간 동안 데이터 처리에 문제가 발생하여 전체 상태를 4시간 이전으로 되돌리고 싶을 수 있습니다.

Flink는 여러 체크포인트를 보존할 수 있도록 지원합니다. Flink의 설정 파일 `conf/flink-conf.yaml`에 다음 설정을 추가하여 최대 보존할 체크포인트 수를 지정할 수 있습니다:

state.checkpoints.num-retained: 20

이 설정은 최근 20개의 체크포인트를 보존하도록 합니다. 특정 체크포인트 지점으로 되돌리고 싶다면, 해당 체크포인트의 경로를 지정하기만 하면 됩니다.

체크포인트로부터 복구

지정된 체크포인트에서 Flink 작업을 시작하려면, 일반적으로 현재 실행 중인 Flink 세션을 중지한 다음 다음 명령을 통해 시작할 수 있습니다. 예를 들어, `/flink/checkpoints/my_app_checkpoints/job_id_xyz/chk-12345`와 같은 경로에서 복구할 수 있습니다.

/path/to/flink/bin/flink run -s /flink/checkpoints/my_app_checkpoints/job_id_xyz/chk-12345/_metadata \
-c com.example.MyFlinkJob /path/to/my-flink-job.jar

이 명령을 스크립트에 포함시켜 매번 체크포인트 복구 스크립트를 직접 실행할 수 있습니다.

세이브포인트 (Savepoints)

세이브포인트는 체크포인트 메커니즘의 특별한 구현으로, 수동으로 체크포인트를 트리거하고 그 결과를 지정된 경로에 영구적으로 저장할 수 있게 합니다. 주로 Flink 클러스터 재시작이나 업그레이드 시 상태 손실을 방지하는 데 사용됩니다. 사용 예시는 다음과 같습니다:

# 지정된 job ID의 세이브포인트를 트리거하고, 그 결과를 지정된 디렉터리에 저장
/path/to/flink/bin/flink savepoint :jobId [:targetDirectory]

수동 세이브포인트

다음과 같이 특정 작업 ID에 대해 세이브포인트를 생성할 수 있습니다. 예를 들어, `0409251eaff826ef2dd775b6a2d5e219`라는 작업 ID에 대해 `hdfs://bigdata/flink/savepoints` 경로에 세이브포인트를 생성하는 경우:

/app/local/flink-1.6.2/bin/flink savepoint 0409251eaff866ef2dd775b6a2d5e219 hdfs://bigdata/flink/savepoints

세이브포인트가 성공적으로 생성되면 일반적으로 `Savepoint completed. Path: hdfs://path...`와 같은 메시지가 표시됩니다.

수동 작업 취소

체크포인트가 예외적으로 중단되거나 수동으로 `Kill`되는 것과 달리, 세이브포인트는 보통 작업자가 수동으로 작업을 중지하고 코드를 업데이트할 때 사용됩니다. `flink cancel` 명령을 사용할 수 있습니다:

/app/local/flink-1.6.2/bin/flink cancel 0409251eaff866ef2dd775b6a2d5e219

지정된 세이브포인트에서 작업 시작

/path/to/flink/bin/flink run -p 8 -s hdfs:///flink/savepoints/my-savepoint-123456 \
-c com.example.analytics.CustomStreamingJob /path/to/my-job-artifacts.jar

상태 백엔드

Flink는 상태의 저장 방식과 위치를 지정하기 위해 다양한 상태 백엔드를 제공합니다. 상태는 Java의 힙 메모리 또는 오프-힙 메모리에 저장될 수 있습니다. 상태 백엔드에 따라 Flink는 애플리케이션의 상태를 직접 관리할 수도 있습니다. 애플리케이션이 매우 큰 상태를 유지할 수 있도록 Flink는 필요한 경우 메모리를 디스크로 오버플로우시키면서 자체적으로 메모리를 관리할 수 있습니다. 기본적으로 모든 Flink 작업은 `flink-conf.yaml` 파일에 지정된 상태 백엔드를 사용합니다. 그러나 작업 내에서 명시적으로 지정된 상태 백엔드는 설정 파일의 기본값을 덮어씁니다.

상태 관리자 유형

  • MemoryStateBackend: 기본 방식으로, JVM 힙 메모리를 기반으로 상태를 저장합니다. 주로 로컬 개발 및 디버깅에 적합합니다.
  • FsStateBackend: 파일 시스템(로컬 파일 시스템 또는 HDFS와 같은 분산 파일 시스템)을 기반으로 상태를 저장합니다. FsStateBackend를 사용하더라도 실제 처리 중인 데이터는 TaskManager의 메모리에 저장되며, 체크포인트 시에만 상태 스냅샷이 지정된 파일 시스템에 기록된다는 점을 주의해야 합니다.
  • RocksDBStateBackend: Flink에 내장된 타사 상태 관리자로, 임베디드 키-값 데이터베이스인 RocksDB를 사용하여 처리 중인 데이터를 저장합니다. 체크포인트 시에는 그 데이터를 지정된 파일 시스템에 영구적으로 저장하므로, RocksDBStateBackend를 사용할 때도 영구 저장용 파일 시스템을 구성해야 합니다. 이는 RocksDB가 임베디드 데이터베이스로서 안정성이 상대적으로 낮지만, 전체 파일 시스템 방식보다 읽기 속도가 빠르고, 전체 메모리 방식보다 저장 공간이 크기 때문에 균형 잡힌 솔루션으로 간주됩니다.

구성 방법

Flink는 두 가지 방식으로 백엔드 관리자를 구성할 수 있도록 지원합니다:

첫 번째 방식: 코드 기반 구성으로, 현재 작업에만 적용됩니다.

// FsStateBackend 구성
env.setStateBackend(new org.apache.flink.runtime.state.filesystem.FsStateBackend("hdfs://namenode:8020/flink/checkpoints"));

// RocksDBStateBackend 구성
env.setStateBackend(new org.apache.flink.contrib.streaming.state.RocksDBStateBackend("hdfs://namenode:8020/flink/checkpoints"));

RocksDBStateBackend를 구성할 때는 다음 종속성을 추가해야 합니다:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-statebackend-rocksdb_2.12</artifactId>
    <version>1.17.1</version> <!-- Flink 버전에 맞게 조정 -->
</dependency>

두 번째 방식: `flink-conf.yaml` 설정 파일 기반 구성으로, 해당 클러스터에 배포된 모든 작업에 적용됩니다.

state.backend: filesystem
state.checkpoints.dir: hdfs://namenode:8020/flink/checkpoints

상태 일관성

실제 애플리케이션에서 스트림 처리 애플리케이션은 스트림 프로세서 외에 데이터 소스(예: Kafka)와 영구 시스템으로의 출력(싱크)을 포함합니다. 종단 간(end-to-end) 일관성 보장은 결과의 정확성이 전체 스트림 처리 애플리케이션을 통해 유지됨을 의미합니다. 각 구성 요소가 자체 일관성을 보장하며, 전체 종단 간 일관성 수준은 모든 구성 요소 중 일관성 수준이 가장 낮은 구성 요소에 의해 결정됩니다. 구체적으로 다음과 같이 나눌 수 있습니다:

  • 내부 보장: 체크포인트에 의존합니다.
  • 소스 측: 외부 소스가 데이터 읽기 위치를 재설정할 수 있어야 합니다.
  • 싱크 측: 장애로부터 복구 시, 데이터가 외부 시스템에 중복으로 기록되지 않도록 보장해야 합니다.

싱크 측의 경우, 두 가지 구체적인 구현 방식이 있습니다:

  • 멱등적(Idempotent) 쓰기: 멱등 연산이란 여러 번 반복 실행하더라도 결과가 한 번만 변경되는 연산을 의미합니다. 즉, 이후의 반복 실행은 효과가 없습니다.
  • 트랜잭션(Transactional) 쓰기: 외부 시스템에 쓰기 위해 트랜잭션을 구축해야 합니다. 구축된 트랜잭션은 체크포인트에 해당하며, 체크포인트가 실제로 완료될 때까지는 모든 해당 결과가 싱크 시스템에 기록되지 않습니다.

트랜잭션 쓰기에는 구체적으로 두 가지 구현 방식이 있습니다: 선행 기록 로그(WAL)와 2단계 커밋(2PC). Flink DataStream API는 `GenericWriteAheadSink` 템플릿 클래스와 `TwoPhaseCommitSinkFunction` 인터페이스를 제공하여 이 두 가지 방식의 트랜잭션 쓰기를 편리하게 구현할 수 있도록 합니다.

Flink + Kafka를 통한 종단 간 정확히 한 번(exactly-once) 의미론

종단 간 상태 일관성 구현은 각 구성 요소가 이를 구현해야 합니다. Flink + Kafka 데이터 파이프라인 시스템(Kafka 입력, Kafka 출력)의 경우, 각 구성 요소가 정확히 한 번 의미론을 어떻게 보장할까요?

  • 내부: 체크포인트 메커니즘을 활용하여 상태를 저장하고, 장애 발생 시 복구하여 내부 상태 일관성을 보장합니다.
  • 소스: Kafka 컨슈머는 소스로서 오프셋을 저장할 수 있습니다. 후속 태스크에 장애가 발생할 경우, 커넥터가 오프셋을 재설정하여 데이터를 다시 소비하고 일관성을 보장할 수 있습니다.
  • 싱크: Kafka 프로듀서는 싱크로서 2단계 커밋 싱크를 사용합니다. 이는 `TwoPhaseCommitSinkFunction` 인터페이스를 구현하고 내부 체크포인트 메커니즘에 의존합니다.

`EXACTLY_ONCE` 의미론(EOS)은 각 입력 메시지가 최종 결과에 한 번만 영향을 미친다는 것을 의미합니다. 여기서 중요한 것은 '한 번 처리'가 아니라 '한 번 영향'입니다. Flink는 항상 EOS를 지원한다고 주장하지만, 이는 주로 Flink 애플리케이션 내부에서 적용되며, 외부 시스템(종단 간)에는 강력한 제약이 따릅니다.

  • 외부 시스템 쓰기가 멱등성을 지원해야 합니다.
  • 외부 시스템이 트랜잭션 방식으로 쓰기를 지원해야 합니다.

Kafka는 0.11 버전 이전에는 `At-Least-Once` 또는 `At-Most-Once` 의미론만 보장할 수 있었지만, 0.11 버전부터 멱등성 전송과 트랜잭션을 도입하여 `EXACTLY_ONCE` 의미론을 보장하기 시작했습니다.

Flink 1.4.0 버전부터 `TwoPhaseCommitSinkFunction` 인터페이스를 도입하여 2단계 커밋 로직을 캡슐화하고, Kafka Sink 커넥터에서 이를 구현했습니다. 이는 Kafka 0.11+ 버전에 의존합니다.

import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;

import java.util.Properties;

public class FlinkKafkaExactlyOnceExample {

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 1000ms마다 체크포인트 활성화 및 EXACTLY_ONCE 모드 설정
        env.enableCheckpointing(1000, CheckpointingMode.EXACTLY_ONCE);
        env.getCheckpointConfig().setCheckpointTimeout(60000); // 60초 타임아웃
        env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); // 동시에 하나의 체크포인트만 허용
        env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500); // 체크포인트 간 최소 500ms 일시 중지

        // Kafka 소스 설정
        Properties consumerProps = new Properties();
        consumerProps.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        consumerProps.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "my-flink-consumer-group");
        consumerProps.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 처음부터 읽기

        DataStream<String> kafkaSourceStream = env.addSource(
                new FlinkKafkaConsumer<>(
                        "input-topic", // 입력 토픽 이름
                        new SimpleStringSchema(),
                        consumerProps
                )
        );

        // 간단한 변환 및 예외 테스트
        DataStream<String> processedStream = kafkaSourceStream.map(data -> {
            // 특정 조건에서 예외를 발생시켜 복구 테스트
            if (data.contains("error_trigger")) {
                System.out.println("Error trigger encountered: " + data);
                throw new RuntimeException("Simulated processing error for: " + data);
            }
            return "Processed: " + data;
        }).name("DataProcessor");

        // Kafka 싱크 설정
        Properties producerProps = new Properties();
        producerProps.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        producerProps.setProperty(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "flink-kafka-transactional-producer"); // Kafka 트랜잭션 ID

        FlinkKafkaProducer<String> kafkaSink = new FlinkKafkaProducer<>(
                "output-topic", // 출력 토픽 이름
                new SimpleStringSchema(),
                producerProps,
                FlinkKafkaProducer.Semantic.EXACTLY_ONCE // EXACTLY_ONCE 의미론 보장
        );

        processedStream.addSink(kafkaSink).name("KafkaSink");

        env.execute("Flink Kafka End-to-End Exactly-Once Example");
    }
}

Flink는 JobManager가 각 TaskManager와 협력하여 체크포인트를 저장하도록 조정합니다. 체크포인트는 StateBackend에 저장되며, 기본 StateBackend는 메모리 기반이지만, 파일 기반으로 변경하여 영구적으로 저장할 수도 있습니다.

체크포인트가 시작되면, JobManager는 체크포인트 배리어(barrier)를 데이터 스트림에 주입하며, 이 배리어는 연산자 간에 전달됩니다. 각 연산자는 현재 상태의 스냅샷을 찍어 상태 백엔드에 저장합니다. 소스 태스크의 경우, 현재 오프셋을 상태로 저장합니다. 다음번 체크포인트 복구 시, 소스 태스크는 오프셋을 다시 제출하여 마지막으로 저장된 위치부터 데이터를 다시 소비할 수 있습니다.

내부의 각 변환(transform) 태스크는 배리어를 만나면 상태를 체크포인트에 저장합니다.

싱크 태스크는 먼저 데이터를 외부 Kafka에 기록합니다. 이 데이터들은 모두 사전 커밋된 트랜잭션에 속하며 (아직 소비할 수 없음), 배리어를 만나면 상태를 상태 백엔드에 저장하고, 다음 체크포인트 데이터를 커밋하기 위한 새로운 사전 커밋 트랜잭션을 시작합니다.

모든 연산자 태스크의 스냅샷이 완료되어 이번 체크포인트가 완료되면, JobManager는 모든 태스크에 알림을 보내 이번 체크포인트가 완료되었음을 확인합니다. 싱크 태스크가 확인 알림을 받으면, 이전 트랜잭션을 공식적으로 커밋하고, Kafka에 미확인 상태였던 데이터는 "확인됨"으로 변경되어 실제로 소비될 수 있게 됩니다.

따라서 실행 과정은 실제로 2단계 커밋이며, 각 연산자 실행이 완료되면 "사전 커밋"을 수행하고, 싱크 작업이 완료될 때까지 "확인 커밋"을 시작합니다. 실행이 실패하면 사전 커밋은 취소됩니다.

구체적인 2단계 커밋 단계는 다음과 같이 요약할 수 있습니다:

  • 첫 번째 데이터가 도착한 후, Kafka 트랜잭션(transaction)을 시작하고, Kafka 파티션 로그에 정상적으로 기록하지만 '미커밋'으로 표시합니다. 이것이 "사전 커밋"입니다.
  • JobManager가 체크포인트 작업을 트리거하면, 배리어가 소스에서 하위로 전달되기 시작합니다. 배리어를 만나는 연산자는 상태를 상태 백엔드에 저장하고 JobManager에 알립니다.
  • 싱크 커넥터는 배리어를 수신하면 현재 상태를 저장하고 체크포인트에 저장한 다음 JobManager에 알리고, 다음 체크포인트 데이터를 커밋하기 위한 다음 단계의 트랜잭션을 시작합니다.
  • JobManager는 모든 태스크로부터 알림을 받으면, 체크포인트가 완료되었음을 나타내는 확인 정보를 발행합니다.
  • 싱크 태스크는 JobManager의 확인 정보를 수신하면 이 기간의 데이터를 공식적으로 커밋합니다.
  • 외부 Kafka는 트랜잭션을 닫고, 커밋된 데이터는 정상적으로 소비될 수 있습니다.

따라서, 시스템이 다운되어 StateBackend를 통해 복구해야 할 경우, 모든 확인 커밋된 작업만 복구할 수 있음을 알 수 있습니다.

Kafka 멱등성 및 트랜잭션

앞서 언급했듯이, Kafka는 0.11 버전 이전에는 `At-Least-Once` 또는 `At-Most-Once` 의미론만 보장할 수 있었지만, 0.11 버전부터 멱등성 전송과 트랜잭션을 도입하여 `EXACTLY_ONCE` 의미론을 보장하기 시작했습니다.

멱등성

멱등성이 도입되기 전에는 Kafka 메시지 전송 및 재시도 과정에서 중복 메시지가 발생할 수 있었습니다. 이를 해결하기 위해 Kafka는 프로듀서 ID (PID)와 시퀀스 번호(Sequence Number)를 도입했습니다. 새로운 프로듀서가 초기화될 때마다 고유한 PID가 할당되며, 이는 사용자에게는 투명합니다.

프로듀서가 각 ``에 대해 메시지를 전송할 때 시퀀스 번호는 0부터 단조 증가합니다. 브로커는 각 ``에 대한 시퀀스 번호를 유지하며, 메시지를 커밋할 때마다 이 번호를 1씩 증가시킵니다. 수신된 각 메시지에 대해, 해당 시퀀스 번호가 브로커가 유지하는 시퀀스 번호(마지막으로 커밋된 메시지의 시퀀스 번호)보다 1 이상 크면 브로커는 이를 수락합니다. 그렇지 않으면 메시지를 폐기합니다:

  • 시퀀스 번호가 브로커가 유지하는 번호보다 1 이상 크다면, 순서가 뒤바뀐 경우를 나타낼 수 있습니다.
  • 시퀀스 번호가 브로커가 유지하는 번호보다 작다면, 이 메시지는 이미 저장되었고 중복 데이터임을 나타냅니다.

멱등성 덕분에, Kafka 메시지 전송 및 재시도 시 중복이 발생하지 않고 안전하게 처리될 수 있습니다.

트랜잭션

트랜잭션은 모든 작업이 하나의 원자적 단위로 성공하거나 실패하며, 부분적인 성공이나 실패는 발생하지 않는 것을 의미합니다. 예를 들어, A가 B에게 1000원을 송금할 때, A의 계좌에서 1000원이 감소하고 B의 계좌에 1000원이 증가하는 두 작업은 반드시 하나의 트랜잭션으로 처리되어야 합니다. 그렇지 않으면 감소만 되거나 증가만 되는 문제가 발생할 수 있습니다. 따라서 두 작업 모두 실패하거나 모두 성공해야 합니다. 분산 환경에서 트랜잭션을 보장하기 위해 일반적으로 2단계 커밋 프로토콜이 사용됩니다.

세션 간 및 모든 파티션에 걸쳐 `EXACTLY-ONCE` 문제를 해결하기 위해 Kafka는 0.11 버전부터 트랜잭션을 도입했습니다.

트랜잭션을 지원하기 위해 Kafka는 전체 트랜잭션 진행을 조정하는 `Transaction Coordinator`를 도입했으며, 트랜잭션을 오프셋 및 그룹 저장과 유사하게 내부 토픽에 영구 저장할 수 있습니다.

사용자는 애플리케이션에 전역 `Transaction ID`를 제공하며, 애플리케이션이 재시작되더라도 이 ID는 변경되지 않습니다. 새로운 프로듀서가 시작된 후, 동일한 `Transaction ID`를 가진 이전 프로듀서가 무효화되도록 보장하기 위해, 프로듀서가 `Transaction ID`를 통해 `PID`를 얻을 때마다 단조 증가하는 `epoch`도 함께 얻습니다. 이전 프로듀서의 `epoch`는 새로운 프로듀서의 `epoch`보다 작으므로, Kafka는 해당 프로듀서가 오래된 프로듀서임을 쉽게 식별하고 해당 요청을 거부할 수 있습니다. `Transaction ID`를 통해 Kafka는 다음을 보장할 수 있습니다:

  • 세션 간 데이터 멱등성 전송: 동일한 `Transaction ID`를 가진 새로운 프로듀서 인스턴스가 생성되어 작동하면, 이전 프로듀서는 작동을 멈춥니다.
  • 세션 간 트랜잭션 복구: 특정 애플리케이션 인스턴스가 다운되면, 새로운 인스턴스는 미완료된 이전 트랜잭션이 커밋되거나 중단되도록 보장하여, 새 인스턴스가 정상 상태에서 작업을 시작할 수 있도록 합니다.

Kafka 트랜잭션의 주요 단계는 다음과 같습니다:

  • 프로듀서는 임의의 브로커에게 `FindCoordinatorRequest`를 보내 `Transaction Coordinator`의 주소를 얻습니다.
  • `Transaction Coordinator`를 찾은 후, 멱등성 특성을 가진 프로듀서는 `InitPidRequest`를 보내 `PID`를 획득해야 합니다.
  • `beginTransaction()` 메서드를 호출하여 트랜잭션을 시작합니다. 프로듀서 로컬에서는 트랜잭션이 시작되었음을 기록하지만, `Transaction Coordinator`는 프로듀서가 첫 번째 메시지를 보낸 후에야 트랜잭션이 시작되었다고 간주합니다.
  • 데이터를 소비하고 변환하여 생산하는 단계(Consume-Transform-Produce)는 전체 트랜잭션 데이터 처리 과정을 포함하며, 다양한 요청을 포함합니다.
  • 위의 데이터 쓰기 작업이 완료되면, 애플리케이션은 `KafkaProducer`의 `commitTransaction` 또는 `abortTransaction` 메서드를 호출하여 현재 트랜잭션을 종료해야 합니다.

2단계 커밋 프로토콜 (Two-Phase Commit Protocol)

2단계 커밋은 분산 트랜잭션을 구현하는 데 자주 사용되는 프로토콜로, 간단히 말해 사전 커밋 + 실제 커밋으로 이해할 수 있습니다. 일반적으로 코디네이터(Coordinator, 이하 C)와 여러 트랜잭션 참여자(Participant, 이하 P) 두 가지 역할로 나뉩니다.

1단계: 준비 (Prepare Phase)

  1. C는 먼저 준비 요청을 로컬 로그에 기록한 다음, 각 P에게 준비 요청을 보냅니다.
  2. P는 준비 요청을 받은 후 트랜잭션을 실행하기 시작합니다. 실행이 성공하면 C에게 Yes 또는 OK 상태를 반환하고, 그렇지 않으면 No를 반환하며, 이 상태를 로컬 로그에 저장합니다.

2단계: 커밋/취소 (Commit/Abort Phase)

  1. C는 P로부터 반환된 상태를 수신합니다. 만약 모든 P의 상태가 Yes라면, C는 트랜잭션 커밋 작업을 시작하고 각 P에게 커밋 요청을 보냅니다. P는 커밋 요청을 받은 후 각자 커밋 트랜잭션 작업을 실행합니다.
  2. 만약 하나 이상의 P의 상태가 No라면, C는 취소 작업을 실행하고 각 P에게 취소 요청을 보냅니다. P는 취소 요청을 받은 후 각자 취소 트랜잭션 작업을 실행합니다.

C 또는 P가 송수신 메시지를 먼저 로그에 기록하는 것은 주로 장애 발생 시 복구를 위한 것으로, WAL(Write-Ahead Log)과 유사한 개념입니다.

태그: Flink 상태관리 키기반상태 연산자상태 체크포인트

8월 14일 18:33에 게시됨