상위: Flink
요약
Flink의 기본 변환 연산자들을 다룹니다. map, flatMap, filter와 같은 1:1 변환부터 keyBy, reduce, process와 같은 그룹화 및 집계 함수까지, 스트림 데이터를 처리하는 핵심 연산자들의 사용법을 학습합니다.
데이터 변환 및 연산
- 데이터 변환이란?
- 원시 데이터를 의미 있는 정보로 재구성
- 집계, 조인, 윈도우 연산 등 후속 처리를 위한 전처리 작업 수행
- 병렬 처리에 최적화, 높은 확장성과 내결함성 확보
기본 변환 함수
| 함수 | 역할 | 용도/설명 |
|---|---|---|
| map | 1:1 변환 | 데이터 포맷 변경, 필드 추출, 단순 계산 |
| flatMap | 1:N 변환 (0개 이상 출력) | 분해, 복수 결과 생성 |
| filter | 조건 기반 필터링 | 노이즈 제거, 유효 데이터 선별 |
map
input_data = [1, 2, 3, 4]
ds = env.from_collection(collection=input_data, type_info=Types.INT())
mapped_stream = ds.map(lambda x: x * 2, output_type=Types.INT())
🟢 결과: [2, 4, 6, 8]
flatMap
input_data = ["hello", "hi"]
ds = env.from_collection(collection=input_data, type_info=Types.STRING())
def split_string(s):
for ch in s:
yield ch
flat_mapped_stream = ds.flat_map(split_string, output_type=Types.STRING())
🟢 결과: ['h', 'e', 'l', 'l', 'o', 'h', 'i']
filter
input_data = [1, 2, 3, 4, 5, 6]
ds = env.from_collection(collection=input_data, type_info=Types.INT())
filtered_stream = ds.filter(lambda x: x % 2 == 0)
🟢 결과: [2, 4, 6]
그룹화 및 집계 함수
| 함수 | 역할 | 용도 |
|---|---|---|
| keyBy | 특정 키 기준으로 스트림 분리 | 그룹별 집계, 상태 저장 등 |
| reduce | 누적하여 단일 결과 도출 | 합계, 최대/최소 등 |
| process | 복잡한 로직 구현 | 타이밍 제어, 상태 기반 처리 등 |
keyBy
input_data = [("A", 1), ("B", 2), ("A", 3), ("B", 4)]
ds = env.from_collection(input_data, type_info=Types.TUPLE([Types.STRING(), Types.INT()]))
keyed_stream = ds.key_by(lambda x: x[0])
reduce
summed_stream = keyed_stream.reduce(lambda a, b: (a[0], a[1] + b[1]))
summed_stream.print()
env.execute("KeyBy Visualization")
🟢 Keyby + reduce결과 (중간 출력):
(A, 1)
(B, 2)
(A, 4)
(B, 6)
process
class MyProcessFunction(ProcessFunction):
def process_element(self, value, ctx):
if value % 2 == 0:
yield value * 2 # 짝수면 2배
else:
yield value # 홀수는 그대로
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)
data = [1, 2, 3, 4, 5, 6]
ds = env.from_collection(collection=data, type_info=Types.INT())
processed_stream = ds.process(MyProcessFunction(), output_type=Types.INT())
env.execute("ProcessFunction Example")
result = list(processed_stream.execute_and_collect())
print("Process 결과:", result)
🟢 결과: [1, 4, 3, 8, 5, 12]
관련 개념
- DataStream API - DataStream API 기본 구조
- Advanced Transformations - 고급 변환 연산자들
- Basic Example - 변환 연산자를 활용한 기본 예제