전체 그래프
Flink

Transformations

data-engineeringflinktransformationsmapfilterreduce

상위: Flink

요약

Flink의 기본 변환 연산자들을 다룹니다. map, flatMap, filter와 같은 1:1 변환부터 keyBy, reduce, process와 같은 그룹화 및 집계 함수까지, 스트림 데이터를 처리하는 핵심 연산자들의 사용법을 학습합니다.

데이터 변환 및 연산

  • 데이터 변환이란?
    • 원시 데이터를 의미 있는 정보로 재구성
    • 집계, 조인, 윈도우 연산 등 후속 처리를 위한 전처리 작업 수행
    • 병렬 처리에 최적화, 높은 확장성과 내결함성 확보

기본 변환 함수

함수역할용도/설명
map1:1 변환데이터 포맷 변경, 필드 추출, 단순 계산
flatMap1: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]

관련 개념