분리형 상태 관리

분리형 상태 관리 (Disaggregated State Management)

개요 (Overview)

Flink의 첫 10년 동안 상태 관리는 TaskManager의 메모리나 로컬 디스크에 기반했습니다. 이 접근 방식은 대부분의 사용 사례에서 잘 동작하지만 몇 가지 한계가 있습니다:

  • 로컬 디스크 제약 (Local Disk Constraints): 상태 크기가 TaskManager의 메모리나 디스크 크기에 제한됩니다.
  • 성긴 리소스 사용 (Spiky Resource Usage): 로컬 상태 모델은 체크포인팅이나 SST 파일 압축 중에 주기적인 CPU 및 네트워크 I/O 버스트를 유발합니다.
  • 무거운 복구 (Heavy Recovery): 복구 중에 상태를 다운로드해야 합니다. 복구 시간은 상태 크기에 비례하며, 상태가 크면 느릴 수 있습니다.

Flink 2.0에서 분리형 상태 관리를 도입했습니다. 이 기능은 사용자가 S3, HDFS 같은 외부 저장 시스템에 상태를 저장할 수 있게 해줍니다. 이는 상태 크기가 매우 클 때 유용합니다. 더 비용 효율적인 방식으로 상태를 저장하거나, 더 가벼운 방식으로 상태를 유지·복구하는 데 사용할 수 있습니다. 분리형 상태 관리의 이점은 다음과 같습니다:

  • 무제한 상태 크기 (Unlimited State Size): 상태 크기는 외부 저장 시스템에 의해서만 제한됩니다.
  • 안정적인 리소스 사용 (Stable Resource Usage): 상태가 외부 저장소에 저장되므로 체크포인트가 매우 가벼울 수 있습니다. 그리고 SST 파일 압축은 원격으로 수행될 수 있습니다(TODO).
  • 빠른 복구 (Fast Recovery): 복구 중 상태를 다운로드할 필요가 없습니다. 복구 시간은 상태 크기와 무관합니다.
  • 유연성 (Flexible): 사용자는 하드웨어를 바꾸지 않고도 서로 다른 외부 저장 시스템이나 I/O 성능 수준을 쉽게 선택하거나, 요구 사항에 따라 저장소를 확장할 수 있습니다.
  • 비용 효율 (Cost-effective): 외부 저장소는 보통 로컬 디스크보다 저렴합니다. 사용자는 병목이 있으면 컴퓨팅 리소스와 저장 리소스를 독립적으로 유연하게 조정할 수 있습니다.

분리형 상태 관리는 세 부분으로 구성됩니다:

  • ForSt State Backend: 상태를 외부 저장 시스템에 저장하는 상태 백엔드입니다. 캐싱과 버퍼링을 위해 로컬 디스크를 활용할 수도 있습니다. 상태 읽기·쓰기에 비동기 I/O 모델이 사용됩니다. 자세한 내용은 ForSt State Backend를 참고하세요.
  • 새 상태 API (New State APIs): 분리형 상태 접근 시 높은 네트워크 지연을 극복하는 데 필수적인 비동기 상태 읽기·쓰기를 수행하기 위해 새 상태 API(State V2)가 도입되었습니다. 자세한 내용은 New State APIs를 참고하세요.
  • SQL 지원 (SQL Support): 많은 SQL 연산자가 분리형 상태 관리와 비동기 상태 접근을 지원하도록 재작성되었습니다. 사용자는 구성을 설정해 쉽게 활성화할 수 있습니다.

분리형 상태와 비동기 상태 접근은 큰 상태에 권장됩니다. 그러나 상태 크기가 작을 때는 동기 상태 접근을 사용한 로컬 상태 관리가 더 나은 선택입니다.

분리형 상태 관리는 아직 실험 단계입니다. 이 기능의 성능과 안정성을 개선하기 위해 노력 중입니다. API와 구성은 향후 릴리스에서 변경될 수 있습니다.

출처: 문서

본문

빠른 시작 (Quick Start)

SQL 작업 (For SQL Jobs)

SQL 작업에서 분리형 상태 관리를 활성화하려면 다음 구성을 설정할 수 있습니다:

state.backend.type: forst
table.exec.async-state.enabled: true

# enable checkpoints, checkpoint directory is required
execution.checkpointing.incremental: true
execution.checkpointing.dir: s3://your-bucket/flink-checkpoints

# We don't support the mini-batch and two-phase aggregation in asynchronous state access yet.
table.exec.mini-batch.enabled: false
table.optimizer.agg-phase-strategy: ONE_PHASE

이렇게 하면 SQL 작업에서 분리형 상태 관리와 비동기 상태 접근을 활용할 수 있습니다. SQL에서 비동기 상태 접근의 완전한 지원은 아직 구현되지 않았습니다. 사용 중인 SQL 연산자가 지원되지 않으면 연산자가 자동으로 동기 상태 구현으로 폴백(fallback)합니다. 이 경우 성능이 최적이 아닐 수 있습니다. 지원되는 상태 저장 연산자는 다음과 같습니다:

  • Rank (Top1, Append TopN)
  • Row Time Deduplicate
  • Aggregate (distinct 없음)
  • Join
  • Window Join
  • Tumble / Hop / Cumulative Window Aggregate

DataStream 작업 (For DataStream Jobs)

DataStream 작업에서 분리형 상태 관리를 활성화하려면 먼저 ForStStateBackend를 사용해야 합니다. per-job 모드에서 코드로 구성합니다:

Configuration config = new Configuration();
config.set(StateBackendOptions.STATE_BACKEND, "forst");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "s3://your-bucket/flink-checkpoints");
config.set(CheckpointingOptions.INCREMENTAL_CHECKPOINTS, true);
env.configure(config);

또는 config.yaml로 구성합니다:

state.backend.type: forst

# enable checkpoints, checkpoint directory is required
execution.checkpointing.incremental: true
execution.checkpointing.dir: s3://your-bucket/flink-checkpoints

그런 다음 새 상태 API로 datastream 작업을 작성해야 합니다. 자세한 내용은 State V2를 참고하세요.

고급 튜닝 옵션 (Advanced Tuning Options)

ForSt State Backend 튜닝 (Tuning ForSt State Backend)

ForStStateBackend는 성능을 튜닝할 수 있는 많은 구성을 가집니다. ForSt의 설계는 RocksDB와 매우 유사하고 구성 가능한 옵션도 거의 같으므로 large state tuning을 참조해 ForSt 상태 백엔드를 튜닝할 수 있습니다.

그 외에 다음 섹션은 ForSt에만 고유한 몇 가지 구성을 소개합니다.

ForSt 기본 저장 위치 (ForSt Primary Storage Location)

기본적으로 ForSt는 체크포인트 디렉터리에 상태를 저장합니다. 이 경우 ForSt는 가벼운 체크포인트와 빠른 복구를 수행할 수 있습니다. 그러나 사용자는 다른 위치(예: S3의 다른 버킷)에 상태를 저장하고 싶을 수 있습니다. 다음 구성을 설정해 기본 저장 위치를 지정할 수 있습니다:

state.backend.forst.primary-dir: s3://your-bucket/forst-state

참고: 이 구성을 설정하면 체크포인팅과 복구 중 ForSt가 기본 저장 위치와 체크포인트 디렉터리 사이에서 파일 복사를 수행하므로 가벼운 체크포인트와 빠른 복구를 활용하지 못할 수 있습니다.

ForSt 로컬 저장 위치 (ForSt Local Storage Location)

기본적으로 ForSt는 비동기 API(State V2)가 사용될 때만 상태를 분리(disaggregate)합니다. DataStream 및 SQL 작업에서 동기 상태 API를 사용할 때 ForSt는 로컬 상태 저장소로만 작동합니다. 작업은 혼합 API 사용으로 여러 ForSt 인스턴스를 포함할 수 있으므로, 동기 로컬 상태 접근과 비동기 원격 상태 접근을 함께 사용하면 전반적으로 더 나은 처리량을 얻을 수 있습니다. 동기 상태 API를 가진 연산자가 상태를 원격에 저장하길 원한다면 다음 구성이 도움이 됩니다:

state.backend.forst.sync.enforce-local: false

그리고 다음으로 로컬 저장 위치를 지정할 수 있습니다:

state.backend.forst.local-dir: path-to-local-dir
ForSt 파일 캐시 (ForSt File Cache)

ForSt는 캐싱과 버퍼링에 로컬 디스크를 사용합니다. 캐시의 세분성은 전체 파일입니다. 이는 기본적으로 활성화되어 있으며, 기본 저장 위치가 로컬로 설정된 경우는 제외됩니다. 캐시에는 두 가지 용량 제한 정책이 있습니다:

  • 크기 기반 (Size-based): 캐시 크기가 제한을 초과하면 가장 오래된 파일을 퇴출합니다.
  • 예약 기반 (Reserved-based): 디스크(캐시 디렉터리가 있는 디스크)의 예약 공간이 충분하지 않으면 가장 오래된 파일을 퇴출합니다.

해당 구성을 예로 들면:

state.backend.forst.cache.size-based-limit: 1GB
state.backend.forst.cache.reserve-size: 10GB

이 둘은 함께 적용될 수 있습니다. 그런 경우 캐시 크기가 크기 기반 제한이나 예약 크기 제한 중 하나라도 초과하면 캐시가 가장 오래된 파일을 퇴출합니다.

캐시 디렉터리는 다음으로 지정할 수 있습니다:

state.backend.forst.cache.dir: /tmp/forst-cache
ForSt 비동기 스레드 (ForSt Asynchronous Threads)

ForSt는 상태를 읽고 쓰기 위해 비동기 I/O를 사용합니다. 세 가지 유형의 스레드가 있습니다:

  • 코디네이터 스레드 (Coordinator thread): 비동기 읽기·쓰기를 조정하는 스레드.
  • 읽기 스레드 (Read thread): 상태를 비동기로 읽는 스레드.
  • 쓰기 스레드 (Write thread): 상태를 비동기로 쓰는 스레드.

비동기 스레드 수는 구성 가능합니다. 일반적으로 기본값이 대부분의 경우 충분하므로 이러한 값을 조정할 필요가 없습니다. 특별한 요구가 있을 경우 다음 구성을 설정해 비동기 스레드 수를 지정할 수 있습니다:

  • state.backend.forst.executor.read-io-parallelism: 읽기를 위한 비동기 스레드 수. 기본값 3.
  • state.backend.forst.executor.write-io-parallelism: 쓰기를 위한 비동기 스레드 수. 기본값 1.
  • state.backend.forst.executor.inline-write: 코디네이터 스레드에서 쓰기 연산을 인라인할지 여부. 기본값 true. false로 설정하면 CPU 사용이 올라갑니다.
  • state.backend.forst.executor.inline-coordinator: 작업 스레드가 코디네이터 스레드가 되게 할지 여부. 기본값 true. false로 설정하면 CPU 사용이 올라갑니다.

ForStStateBackend는 기본 데이터베이스 코어로 ForSt를 활용합니다. ForSt의 현재 버전은 frocksdb에서 포크되었으며, Flink 로컬 상태 관리를 위해 특별히 설계된 임베디드 데이터베이스 코어로 아키텍처링되었습니다.

DB를 분리형 아키텍처로 전환하는 동안 기존 프레임워크 내에서 상당한 아키텍처·엔지니어링 제약에 직면했습니다. 이러한 과제를 해결하기 위해 커뮤니티 구성원들은 현재 Rust로 작성된 차세대 클라우드 네이티브 ForSt DB를 작업 중입니다.

새로운 ForSt DB의 주요 장점 (Key Advantages of the New ForSt DB):
  • 아키텍처 단순성 (Architectural Simplicity): 높은 확장성을 위해 설계된 간소화된 코드베이스.
  • 스트림 네이티브 설계 (Stream-Native Design): 대규모 스트림 처리의 고유한 요구에 맞게 최적화됨.
  • 클라우드 네이티브 (Cloud-Native): 분리(disaggregation)를 지원하도록 처음부터 구축됨.
로드맵 & 유지보수 (Roadmap & Maintenance):
  • 릴리스 일정 (Release Schedule): 첫 번째 안정적인 오픈소스 버전은 2026년 후반(낙관적으로 8월)으로 예상됩니다.
  • 폐기 공지 (Deprecation Notice): 새로운 Rust 기반 구현으로 초점을 옮기면서 frocksdb 기반 ForSt 버전은 더 이상 활발히 개발되지 않으며, 새 릴리스 이후 단계적으로 폐기(deprecated)될 예정입니다.

더 알아보기 (Learn more)