상위: Flink
요약
DataStream API는 Flink의 핵심 스트림 처리 API입니다. 무한한 데이터 스트림을 처리하기 위한 표준 인터페이스를 제공하며, Source → Transformations → Sink 구조로 실시간 데이터 처리 파이프라인을 구성할 수 있습니다.
Stream Processing
동전 분류기, 줄서기
- 모든 구성 요소는 직렬로 연결
- 지속적으로 시스템에 입력 → 다양한 대기열로 출력(분류됨)
- 스트림 처리 시스템은 무한 데이터셋 처리를 위해 데이터 기반 처리 방법 사용
실시간 분석 (Realtime Analytics)에 사용
- 금융 거래, IoT 센서 데이터, 소셜 미디어 피드 등 거의 실시간 분석 가능
- 빠른 반응 및 시스템 효율성
- 낮은 지연 시간
- 동시에 여러 이벤트 처리 가능 → 분산 환경에서 높은 처리량 유지
- 확장성과 내결함성
- 수평적 확장 가능 (클러스터 환경)
- 체크포인트 등으로 상태 복구 가능
vs. Batch

| 항목 | Batch 처리 | Streaming 처리 |
|---|---|---|
| 처리 방식 | 일정 기간 단위로 수집 후 일괄 처리 | 연속된 데이터를 하나씩 처리 |
| 처리량 | 대규모 데이터 단위 | 소량의 레코드 단위 |
| 속도 | 수분~시간 지연 시간 | (준)실시간 |
| 사용 환경 | 복잡한 분석이 요구됨 / 처리량 많음 | 실시간 분석 필요 / 고급 메시징 환경 |
| 예시 | 급여 및 청구 시스템 | ATM, 부정행위 탐지, SNS 분석 시스템 |
DataStream API
-
스트림 데이터를 처리하기 위한 핵심 API
- 유한/무한 불변 데이터 집합을 변환(transformation)하여 처리
- Source → Transformations → Sink 구조
**execute()**호출 시 실행 계획(데이터플로우 그래프)으로 실행
-
이벤트 기반 처리 지원
- 실시간 데이터 스트림 처리 위한 표준 인터페이스 제공
-
효율적 파이프라인 구성
- 여러 변환을 조합하여 효율적인 데이터 처리
-
다양한 활용 가능성
- 금융, IoT, SNS 등 다양한 분야의 실시간 분석 및 처리
-
실시간 데이터 처리
-
상태 관리 기능 (이전 이벤트 정보 유지)
-
확장성과 유연성 제공
- 다양한 소스와 통합, 대규모 처리 가능
Directed Acyclic Graph (DAG)
-
논리 모델(DAG): 방향성 있는 비순환 그래프로, 노드(연산자)와 간선(데이터 흐름)으로 구성
- 사용자가 작성한 처리 흐름을 순서대로 표현
-
연산자 순차/병렬 처리 가능
-
DAG 형태 → 각 노드를 병렬 확장 가능
-
최적화 가능: 데이터 교환, 파티셔닝, 스케줄링
-
장애 발생 시 체크포인트 통해 복구 용이
데이터 소스 (Data Source)
-
파일 시스템
- 로컬 파일, HDFS 등에서 데이터를 읽어 들임
-
메시지 큐 / 스트리밍 플랫폼
- 예시: Apache Kafka, Amazon Kinesis
- 실시간 스트림 처리에 사용
-
소켓 스트림
- 네트워크 소켓을 통해 실시간 로그/이벤트 데이터 수신
-
Flink의 StreamExcutionEnvironment 사용
- 실행 환경 객체 생성 후
**read_text_file**,**from_collection**등으로 데이터 소스 정의
- 실행 환경 객체 생성 후
-
커넥터 API로 외부 시스템과 쉽게 연결 가능
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")
Sink
📥 데이터 소스 → 연산 (변환) → Sink (출력/저장) 📤

| Sink 종류 | 설명 |
|---|---|
| PrintSink | print()로 콘솔 출력 |
| FileSink | 텍스트/CSV/JSON 파일로 저장 |
| KafkaSink | Kafka에 데이터 저장 |
| DatabaseSink | JDBC 통해 DB(MySQL, PostgreSQL 등)에 저장 |
| ElasticsearchSink | Elasticsearch에 저장 |
Lazy Evaluation (게으른 평가)
연산은 실행 시점까지 수행되지 않고, 내부 실행 계획으로만 기록됨
- 연산 계획 수립
- 여러 연산(map, filter, flatMap 등)이 바로 실행되지 않고 그래프(DAG)로 기록
- 최적화 기회 제공
- 불필요한 연산 제거, 병합 등 전체 연산 최적화 가능
- 실행 트리거
**sink**또는**execute()**호출 시 DAG 기반으로 전체 연산 한 번에 실행됨
| 장점 | 설명 |
|---|---|
| 효율성 | 불필요한 계산 방지, 자원 절약 |
| 성능 향상 | 연산 병합 및 중간 저장 최소화 |
| 유연성 | 실행 전까지 실행 계획 수정/변경 가능 |
관련 개념
- Basic Example - DataStream API 기본 사용법
- Transformations - 기본 변환 연산자들
- Advanced Transformations - 고급 변환 연산자들