상위: Flink
요약
Flink의 상태 관리 시스템을 다룹니다. Keyed State, Operator State, Broadcast State의 개념과 사용법을 학습하고, 상태를 활용한 복잡한 스트림 처리 로직을 구현하는 방법을 알아봅니다.
상태(state)와 데이터 저장
- 중간 결과나 정보를 상태(state)에 저장
- 데이터 스트림 처리 중 중간 결과 저장
- 예: 숫자의 합계 계산 시, 지금까지의 합이 상태에 해당됨

상태의 유형
- 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
상태 최적화
메모리 관리
- 상태 크기 모니터링
- 불필요한 상태 정리
- 압축 및 직렬화 최적화
성능 고려사항
- 상태 접근 빈도 최소화
- 배치 업데이트 활용
- 적절한 상태 백엔드 선택
관련 개념
- Checkpointing - 상태 저장 및 복구 메커니즘
- Savepoint - 수동 상태 저장
- DataStream API - 상태를 활용한 스트림 처리