전체 그래프
Flink

Stream Iteration

data-engineeringflinkiterationstream-splittingfilter

상위: Flink

요약

Flink에서 스트림을 분할하고 반복 처리하는 방법을 다룹니다. filter 연산자를 활용한 스트림 분할과 드라이버 기반 반복 처리 방식을 통해 복잡한 데이터 처리 시나리오를 구현하는 방법을 학습합니다.

데이터 스트림의 분할 및 반복

데이터 스트림 분할 (Filter)

  • 하나의 스트림에서 서로 다른 조건이나 처리 목적에 따라 데이터를 여러 개의 하위 스트림으로 나누는 작업
    • 서로 다른 처리 로직이나 파이프라인을 개별 하위 스트림에 적용
    • 복잡한 데이터 처리 시 유연성 향상
    • 초기 Flink 버전: split(), select() 사용
    • 현재는 filter 변환 연산 주로 사용
택배 물류 데이터 스트림 
  - 무거운 택배 라인 (Heavy Stream)
  - 도서산간 라인 (Remote Stream)
  - 냉장 택배 라인 (Cold Stream)
  - 일반 택배 라인 (Main Stream)

기본 스트림 분할 예제

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.common.typeinfo import Types

env = StreamExecutionEnvironment.get_execution_environment()

# 예제 입력: 정수 리스트
data = env.from_collection([1, 2, 3, 4, 5, 6], type_info=Types.INT())

# filter 연산을 이용한 데이터 스트림 분할
even_stream = data.filter(lambda x: x % 2 == 0)  # 짝수 스트림
odd_stream = data.filter(lambda x: x % 2 != 0)   # 홀수 스트림

env.execute("Split Stream Example")

# 결과 수집
evens = list(even_stream.execute_and_collect())
odds = list(odd_stream.execute_and_collect())

print("Even stream:", evens)
print("Odd stream:", odds)

🟢 결과:

  • 짝수(메인): [2, 4, 6]
  • 홀수(Side Output): [1, 3, 5]

택배 분류 예제

data = [
  (25, "서울 강남구", False),
  (15, "강원도 산간 지역", False),
  (10, "부산 해운대", True),
  (5, "서울 관악구", False)
]

# 무거운 택배
heavy = ds.filter(lambda x: x[0] > 20)
# 산간 지역
remote = ds.filter(lambda x: "산간" in x[1])
# 냉장 택배
cold = ds.filter(lambda x: x[2])
# 일반 택배
normal = ds.filter(lambda x: not (x[0] > 20 or "산간" in x[1] or x[2]))

env.execute("택배 분리 처리")

🟢 결과:

  • 무거운 택배: (25, '서울 강남구', False)
  • 산간 택배: (15, '강원도 산간 지역', False)
  • 냉장 택배: (10, '부산 해운대', True)
  • 일반 택배: (5, '서울 관악구', False)

데이터 스트림 반복 (Iteration)

  • 알고리즘을 반복 수행하여 결과를 점진적으로 개선
    • 드라이버에서 반복문을 통해 피드백 루프 구성
  1. 초기 설정 및 함수 정의
    • 데이터 스트림 생성, 변환 연산 적용
    • Job 실행 및 결과 수집
  2. 드라이버에서 반복 로직 수행
    • while 반복문 내부에서 결과 수집 및 조건 검토
  3. 피드백 및 종료
    • 조건 만족 여부에 따라 반복 종료
    • 반복 종료 시 외부로 결과 출력

반복 처리 주의사항

  • 종료 조건:
    • 무한 루프 방지를 위해 명확한 조건 설정
    • 조건이 부적절하면 시스템 자원 낭비 우려
  • 상태 관리:
    • 중간 상태가 반복적으로 업데이트되므로 체크포인트 필요
  • 성능 최적화:
    • 드라이버 반복 방식은 오버헤드 발생
    • 반복 데이터 피드백 시 병목/지연 방지 설계 필요

반복 처리 예제

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.common.typeinfo import Types

# 데이터에 map 연산 (값을 2배)
def run_flink_job(input_data):
    env = StreamExecutionEnvironment.get_execution_environment()
    ds = env.from_collection(input_data, type_info=Types.INT())

    mapped_ds = ds.map(lambda x: x * 2, output_type=Types.INT())
    env.execute("Driver-based iteration job")

    result = list(mapped_ds.execute_and_collect())
    return result

# 초기 데이터
data = [10]

# 드라이버  반복
while True:
    data = run_flink_job(data)
    print("중간 결과:", data)
    if all(x >= 100 for x in data):
        break

print("최종 결과:", data)

🟢 결과:

  • 중간 결과: [20]
  • 중간 결과: [40]
  • 중간 결과: [80]
  • 중간 결과: [160]
  • 최종 결과: [160]

관련 개념