전체 그래프
Flink

Evictors

data-engineeringflinkevictorsmemory-managementwindow-optimization

상위: Flink

요약

⚠️ 공식 문서 기준, PyFlink DataStream API에서는 Evictor가 아직 지원되지 않는다. 아래 코드는 Java/Scala DataStream API의 개념 설명이거나 윈도우에 실제 등록하지 않은 시뮬레이션 예시다.

Flink의 Evictor 메커니즘을 다룹니다. 윈도우에서 일부 요소를 제거하여 메모리 사용량을 최적화하고 성능을 향상시키는 방법을 학습합니다. CountEvictor, TimeEvictor, DeltaEvictor의 사용법과 커스텀 Evictor 구현 방법을 알아봅니다.

Evictor

  • 윈도우에서 일부 요소를 제거(evict) 하는 역할

  • 윈도우 연산이 수행되기 전에 적용되며, 남겨둘 요소와 버릴 요소 결정

  • 사용 이유 (예: 길이 1시간인 윈도우에서 최대값 구하기

    • 오래된 데이터를 제거해 메모리 사용 최적화
    • 최근 데이터 기반 분석 가능 → 실시간성 유지
    • 예: 1시간 윈도우에서 트리거는 그대로지만 최근 10분만 사용해 계산

CountEvictor(N)

  • 각 윈도우마다 N개의 요소만 유지하고 나머지는 제거
  • 예: 윈도우에 100개가 쌓였는데 CountEvictor(50)을 쓰면 50개만 남기고 50개 evict

TimeEvictor(t)

  • 현재 윈도우의 가장 늦은 타임스탬프 - t 보다 이전의 모든 요소를 제거
  • 예: TimeEvictor(10초) → 윈도우 내에서 최신 이벤트 시간으로부터 10초보다 더 오래된 이벤트는 제외하고 연산

DeltaEvictor(Δ)

  • 델타 함수를 이용한 사용자 정의 기준 (요소 간의 차이, 속성 변화량 등)
  • Δ 기준보다 크면 제거하는 방식
  • 사용자가 DeltaFunction을 직접 정의해야 함

TimeEvictor 예제

  • 오래된 데이터 제거
    • 윈도우는 60분치 데이터를 모으지만, 실제 계산에는 최근 10분 이내의 데이터만 사용
# 데이터 생성 (타임스탬프 포함)
data = [(i, time.time() - (i * 60)) for i in range(1, 61)]  # 1시간 데이터 ( 단위)

# Evictor: 최근 10분 데이터 유지
class TimeEvictor:
    def __init__(self, max_time_seconds):
        self.max_time_seconds = max_time_seconds

    def evict_before(self, elements):
        current_time = time.time()
        filtered = [e for e in elements if current_time - e[1] <= self.max_time_seconds]
        removed = [e for e in elements if e not in filtered]

        print("**Evictor 적용 결과**")
        print(f"총 입력 데이터 개수: {len(elements)}")
        print(f"유지된 데이터 개수 (최근 10분): {len(filtered)}")
        print(f"제거된 데이터 개수: {len(removed)}")
        print("**유지된 데이터 (Evictor 적용 후)**")
        print(filtered)
        print("**제거된 데이터**")
        print(removed)

        return filtered

# 윈도우 설정 (1시간) + Evictor 적용 (최근 10분 유지)
def time_evictor_example():
    print("**Evictor 적용 전 전체 데이터:**")
    print(data)  # 원본 데이터 출력

    evictor = TimeEvictor(max_time_seconds=600)  # 최근 10분 데이터 유지
    result = evictor.evict_before(data)

time_evictor_example()

출력 :

**Evictor 적용 결과**
 입력 데이터 개수: 60
유지된 데이터 개수 (최근 10분): 9
제거된 데이터 개수: 51

Evictor 구현 패턴

1. 기본 Evictor 구조

class CustomEvictor(Evictor):
    def __init__(self, max_elements):
        self.max_elements = max_elements

    def evict_before(self, elements):
        # 정렬 (최신 데이터 우선)
        sorted_elements = sorted(elements, key=lambda x: x[1], reverse=True)
        
        # 최대 개수만 유지
        return sorted_elements[:self.max_elements]

    def evict_after(self, elements):
        # 윈도우 연산  추가 정리 (필요시)
        return elements

2. 시간 기반 Evictor

class TimeBasedEvictor(Evictor):
    def __init__(self, time_window_seconds):
        self.time_window_seconds = time_window_seconds

    def evict_before(self, elements):
        if not elements:
            return elements
            
        # 최신 타임스탬프 기준
        latest_timestamp = max(elements, key=lambda x: x[1])[1]
        cutoff_time = latest_timestamp - self.time_window_seconds
        
        # 시간 기준 필터링
        filtered = [e for e in elements if e[1] >= cutoff_time]
        
        print(f"TimeEvictor: {len(elements)} -> {len(filtered)} elements")
        return filtered

3. 값 기반 Evictor

class ValueBasedEvictor(Evictor):
    def __init__(self, threshold):
        self.threshold = threshold

    def evict_before(self, elements):
        # 임계값 이상의 값만 유지
        filtered = [e for e in elements if e[1] >= self.threshold]
        
        print(f"ValueEvictor: {len(elements)} -> {len(filtered)} elements")
        return filtered

Evictor 최적화 기법

1. 메모리 효율성

class MemoryEfficientEvictor(Evictor):
    def __init__(self, max_memory_elements):
        self.max_memory_elements = max_memory_elements

    def evict_before(self, elements):
        if len(elements) <= self.max_memory_elements:
            return elements
            
        # 메모리 사용량이 임계값을 초과하면 오래된 데이터 제거
        sorted_elements = sorted(elements, key=lambda x: x[1])
        return sorted_elements[-self.max_memory_elements:]

2. 성능 최적화

class OptimizedEvictor(Evictor):
    def __init__(self, batch_size=1000):
        self.batch_size = batch_size

    def evict_before(self, elements):
        # 배치 단위로 처리하여 성능 향상
        if len(elements) <= self.batch_size:
            return self.process_batch(elements)
        
        #  데이터셋은 배치로 나누어 처리
        result = []
        for i in range(0, len(elements), self.batch_size):
            batch = elements[i:i + self.batch_size]
            result.extend(self.process_batch(batch))
        
        return result

    def process_batch(self, batch):
        # 배치 처리 로직
        return batch

3. 적응형 Evictor

class AdaptiveEvictor(Evictor):
    def __init__(self, initial_threshold, adaptation_rate=0.1):
        self.threshold = initial_threshold
        self.adaptation_rate = adaptation_rate
        self.performance_history = []

    def evict_before(self, elements):
        # 성능 기반 임계값 조정
        current_performance = self.measure_performance(elements)
        self.performance_history.append(current_performance)
        
        # 성능이 저하되면 임계값 조정
        if len(self.performance_history) > 5:
            avg_performance = sum(self.performance_history[-5:]) / 5
            if avg_performance < 0.8:  # 성능 저하 감지
                self.threshold = max(1, self.threshold - 1)
            elif avg_performance > 0.95:  # 성능 양호
                self.threshold = min(len(elements), self.threshold + 1)
        
        # 조정된 임계값으로 필터링
        return elements[:self.threshold]

    def measure_performance(self, elements):
        # 성능 측정 로직 (: 처리 시간, 메모리 사용량 )
        return 0.9  # 예시 

Evictor 사용 시나리오

1. 실시간 추천 시스템

# 사용자별 최근 행동만 유지
class UserBehaviorEvictor(Evictor):
    def evict_before(self, elements):
        # 사용자별로 최근 100개 행동만 유지
        user_actions = {}
        for element in elements:
            user_id = element[0]
            if user_id not in user_actions:
                user_actions[user_id] = []
            user_actions[user_id].append(element)
        
        #  사용자별로 최근 100개만 유지
        result = []
        for user_id, actions in user_actions.items():
            sorted_actions = sorted(actions, key=lambda x: x[2], reverse=True)
            result.extend(sorted_actions[:100])
        
        return result

2. 금융 거래 모니터링

# 고액 거래만 유지
class HighValueTransactionEvictor(Evictor):
    def __init__(self, min_amount):
        self.min_amount = min_amount

    def evict_before(self, elements):
        # 고액 거래만 유지
        high_value = [e for e in elements if e[1] >= self.min_amount]
        
        # 금액 순으로 정렬하여 상위 거래만 유지
        sorted_by_amount = sorted(high_value, key=lambda x: x[1], reverse=True)
        return sorted_by_amount[:1000]  # 상위 1000개만 유지

관련 개념