상위: 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)
- 알고리즘을 반복 수행하여 결과를 점진적으로 개선
- 드라이버에서 반복문을 통해 피드백 루프 구성
- 초기 설정 및 함수 정의
- 데이터 스트림 생성, 변환 연산 적용
- Job 실행 및 결과 수집
- 드라이버에서 반복 로직 수행
while반복문 내부에서 결과 수집 및 조건 검토
- 피드백 및 종료
- 조건 만족 여부에 따라 반복 종료
- 반복 종료 시 외부로 결과 출력
반복 처리 주의사항
- 종료 조건:
- 무한 루프 방지를 위해 명확한 조건 설정
- 조건이 부적절하면 시스템 자원 낭비 우려
- 상태 관리:
- 중간 상태가 반복적으로 업데이트되므로 체크포인트 필요
- 성능 최적화:
- 드라이버 반복 방식은 오버헤드 발생
- 반복 데이터 피드백 시 병목/지연 방지 설계 필요
반복 처리 예제
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]
관련 개념
- Advanced Transformations - 고급 변환 연산자들
- Transformations - 기본 변환 연산자들
- DataStream API - DataStream API 기본 구조