상위: Flink
요약
Flink의 기본 코드 구조와 실행 방법을 다룹니다. StreamExecutionEnvironment를 통한 실행 환경 설정부터 데이터 소스, 변환, 싱크까지의 전체 파이프라인을 구성하는 방법을 학습합니다.
Flink 코드 구조 (PyFlink)
1. 실행 환경 생성
from pyflink.datastream import StreamExecutionEnvironment
# 실행 환경 생성
env = StreamExecutionEnvironment.get_execution_environment()
- Flink 작업을 구성하고 실행하는 중심 객체
- 데이터 소스, 변환, 싱크 등을 생성하는 메서드 제공
2. 데이터 소스(Source) 정의
# 리스트 데이터를 스트림으로 변환
data_stream = env.from_collection([1, 2, 3, 4, 5])
- 외부 데이터를 PyFlink 내부로 불러오는 역할
- 예: 파일, 소켓, 컬렉션 등
3. 데이터 변환(Transformation) 적용
# 각 숫자에 * 2 연산 수행
transformed_stream = data_stream.map(lambda x: x * 2)
- 입력 데이터를 원하는 형태로 가공 (예: 매핑, 필터링, 집계 등)
4. 데이터 싱크(Sink) 설정
# 결과를 콘솔에 출력
transformed_stream.print()
- 처리된 데이터를 외부 시스템으로 출력
- 예: 콘솔, 파일, 데이터베이스 등
5. 작업 실행
# Flink 작업 실행
env.execute("Simple Flink Job")
- 지금까지 구성한 소스, 변환, 싱크를 기반으로 실행
완전한 예제
from pyflink.datastream import StreamExecutionEnvironment
# 1. 실행 환경 생성
env = StreamExecutionEnvironment.get_execution_environment()
# 2. 데이터 소스(Source) 정의 (리스트 데이터를 스트림으로 변환)
data_stream = env.from_collection([1, 2, 3, 4, 5])
# 3. 데이터 변환(Transformation) 적용 (각 숫자에 * 2 연산 수행)
transformed_stream = data_stream.map(lambda x: x * 2)
# 4. 데이터 싱크(Sink) 설정 (결과를 콘솔에 출력)
transformed_stream.print()
# 5. 작업 실행
env.execute("Simple Flink Job")
예상 결과
- 예상 결과: 2 → 4 → 6 → 8 → 10
- 실제 출력 순서는 병렬 처리이기 때문에 랜덤하게 나타날 수 있음
DAG 시각화 예제
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.common.typeinfo import Types
# 실행 환경 생성
env = StreamExecutionEnvironment.get_execution_environment()
# 병렬도 설정 (선택 사항)
env.set_parallelism(1)
# 예시 데이터 소스 생성
data = env.from_collection(
collection=[("apple", 1), ("banana", 1), ("apple", 1)],
type_info=Types.TUPLE([Types.STRING(), Types.INT()])
)
# transformation 적용: keyBy + sum
result = data.key_by(lambda x: x[0]) \
.sum(1)
# DAG 실행 계획 출력 (JSON 문자열 형식)
execution_plan = env.get_execution_plan()
print("DAG Execution Plan:\n", execution_plan)
# 출력: 콘솔에 출력
result.print()
# 실행
env.execute("Basic PyFlink Job")
실행 결과
- 노드 3개 구성:
- Node 1 : Source (데이터 생성)
- from_collection() 을 통해 생성된 데이터 스트림
- Node 2 : KeyBy 연산자
- key_by(lambda x : x[0]) 에 해당
- Node 1 에서 받아옴 ( ship_strategy: " FORWARD " → 로컬 전달 )
- Node 4 : Reduce 연산자 (sum)
- Reduce 연산자, sum(1) 연산에 해당
- Node 2 의 출력 ( ship_strategy: " HASH "→ 키 해시 기반 분산 )
- Node 1 : Source (데이터 생성)
- 출력:
- ('apple', 1)
- ('banana', 1)
- ('apple', 2)
파일 소스 예제
from pyflink.datastream import StreamExecutionEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
file_path = "../data/data.csv"
text_stream = env.read_text_file(file_path)
text_stream.print()
env.execute("File Source Example")
Java Encoder를 사용한 파일 출력
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors.file_system import FileSink
from pyflink.datastream.formats import Encoder
from pyflink.java_gateway import get_gateway
# Java Encoder 생성
gateway = get_gateway()
j_string_encoder = gateway.jvm.org.apache.flink.api.common.serialization.SimpleStringEncoder()
# Python Encoder 래핑
encoder = Encoder(j_string_encoder)
# FileSink 정의
file_sink = FileSink.for_row_format(
"../output/result", encoder
).build()
# 데이터 연결 및 실행
data_stream.sink_to(file_sink)
env.execute("File Sink Example")
관련 개념
- Installation - Flink 설치 및 환경 설정
- DataStream API - DataStream API 상세 사용법
- Transformations - 다양한 변환 연산자들