전체 그래프
Flink

State Management

data-engineeringflinkstatemanagementkeyed-state

상위: Flink

요약

Flink의 상태 관리 시스템을 다룹니다. Keyed State, Operator State, Broadcast State의 개념과 사용법을 학습하고, 상태를 활용한 복잡한 스트림 처리 로직을 구현하는 방법을 알아봅니다.

상태(state)와 데이터 저장

  • 중간 결과나 정보를 상태(state)에 저장
    • 데이터 스트림 처리 중 중간 결과 저장
    • 예: 숫자의 합계 계산 시, 지금까지의 합이 상태에 해당됨

flink-17.png

상태의 유형

  • Keyed State: 키 기반 상태 (예: 사용자 ID별 세션 정보 저장)
  • Operator State: 연산자 전체에 공유되는 상태 (예: 외부 시스템에서 읽은 마지막 위치)
  • Broadcast State: 모든 하위 작업에 동일한 상태를 브로드캐스트

상태 관리 중요성

  • 정확성 보장: 중간 결과 저장을 통해 처리 정확성 유지
  • 장애 복구: 체크포인트 & 세이브포인트로 장애 발생 시 복구 가능
  • 성능 최적화: 메모리 효율 향상, 대규모 데이터에도 안정적인 성능

상태 관리 예제

class KeyedSum(KeyedProcessFunction):
    def __init__(self):
        self.state = None  # Keyed State 저장 변수

    def open(self, runtime_context):
        self.state = runtime_context.get_state(
            ValueStateDescriptor("sum", Types.LONG())
        )

    def process_element(self, value, ctx):
        current_sum = self.state.value() or 0
        new_sum = current_sum + value[1]
        self.state.update(new_sum)
        print(f"현재 상태: ID={value[0]}, 누적 금액={new_sum}")
        return value[0], new_sum
def main():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_parallelism(1)
    env.enable_checkpointing(10000)  # 10초마다 체크포인트

    df = pd.read_csv("../data/data.csv")
    transactions = df[['transaction_id', 'amount']].dropna().values.tolist()

    transaction_stream = env.from_collection(transactions)
    keyed_stream = transaction_stream.key_by(lambda x: x[0]).process(KeyedSum())
    keyed_stream.print()

    env.execute("Keyed State Example")

Keyed State

개념

  • 키별로 독립적인 상태를 유지
  • 같은 키의 데이터만 같은 상태에 접근 가능
  • 파티셔닝과 함께 사용되어 확장성 보장

사용 예시

# 사용자별 구매 금액 누적
class UserPurchaseState(KeyedProcessFunction):
    def __init__(self):
        self.purchase_sum = None
        self.purchase_count = None

    def open(self, runtime_context):
        self.purchase_sum = runtime_context.get_state(
            ValueStateDescriptor("purchase_sum", Types.LONG())
        )
        self.purchase_count = runtime_context.get_state(
            ValueStateDescriptor("purchase_count", Types.INT())
        )

    def process_element(self, value, ctx):
        # 현재 상태 가져오기
        current_sum = self.purchase_sum.value() or 0
        current_count = self.purchase_count.value() or 0
        
        # 상태 업데이트
        new_sum = current_sum + value[1]
        new_count = current_count + 1
        
        self.purchase_sum.update(new_sum)
        self.purchase_count.update(new_count)
        
        return (value[0], new_sum, new_count)

Operator State

개념

  • 연산자 전체에 공유되는 상태
  • 키와 무관하게 연산자 레벨에서 관리
  • 주로 소스나 싱크에서 사용

사용 예시

class CustomSource(SourceFunction, CheckpointedFunction):
    def __init__(self):
        self.offset_state = None  # Operator State (ListState만 지원)

    def initialize_state(self, context):
        # Operator State는 CheckpointedFunction.initialize_state에서
        # OperatorStateStore를 통해 ListState로 획득한다 (ValueState 미지원)
        self.offset_state = context.get_operator_state_store().get_list_state(
            ListStateDescriptor("offset", Types.LONG())
        )

    def run(self, ctx):
        current_offset = self.offset.value() or 0
        # 데이터 읽기 로직
        # ...
        self.offset.update(new_offset)

Broadcast State

개념

  • 모든 하위 작업에 동일한 상태를 브로드캐스트
  • 설정 정보나 참조 데이터를 모든 태스크에 전파
  • 동적 설정 변경에 유용

사용 예시

class BroadcastProcessFunction(BroadcastProcessFunction):
    def __init__(self):
        self.broadcast_state = None

    def open(self, runtime_context):
        # Broadcast State 초기화
        self.broadcast_state = runtime_context.get_broadcast_state(
            MapStateDescriptor("config", Types.STRING(), Types.STRING())
        )

    def process_element(self, value, ctx):
        # Broadcast State에서 설정 읽기
        config = self.broadcast_state.get("threshold")
        if value > int(config):
            # 처리 로직
            pass

상태 최적화

메모리 관리

  • 상태 크기 모니터링
  • 불필요한 상태 정리
  • 압축 및 직렬화 최적화

성능 고려사항

  • 상태 접근 빈도 최소화
  • 배치 업데이트 활용
  • 적절한 상태 백엔드 선택

관련 개념