전체 그래프
Flink

DataStream API

data-engineeringflinkdatastreamapistreaming

상위: Flink

요약

DataStream API는 Flink의 핵심 스트림 처리 API입니다. 무한한 데이터 스트림을 처리하기 위한 표준 인터페이스를 제공하며, Source → Transformations → Sink 구조로 실시간 데이터 처리 파이프라인을 구성할 수 있습니다.

Stream Processing

동전 분류기, 줄서기

  • 모든 구성 요소는 직렬로 연결
  • 지속적으로 시스템에 입력 → 다양한 대기열로 출력(분류됨)
  • 스트림 처리 시스템은 무한 데이터셋 처리를 위해 데이터 기반 처리 방법 사용

실시간 분석 (Realtime Analytics)에 사용

  • 금융 거래, IoT 센서 데이터, 소셜 미디어 피드 등 거의 실시간 분석 가능
  • 빠른 반응 및 시스템 효율성
    • 낮은 지연 시간
    • 동시에 여러 이벤트 처리 가능 → 분산 환경에서 높은 처리량 유지
  • 확장성과 내결함성
    • 수평적 확장 가능 (클러스터 환경)
    • 체크포인트 등으로 상태 복구 가능

vs. Batch

flink-15.png

항목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 (출력/저장) 📤

flink-16.png

Sink 종류설명
PrintSinkprint()로 콘솔 출력
FileSink텍스트/CSV/JSON 파일로 저장
KafkaSinkKafka에 데이터 저장
DatabaseSinkJDBC 통해 DB(MySQL, PostgreSQL 등)에 저장
ElasticsearchSinkElasticsearch에 저장

Lazy Evaluation (게으른 평가)

연산은 실행 시점까지 수행되지 않고, 내부 실행 계획으로만 기록됨

  • 연산 계획 수립
    • 여러 연산(map, filter, flatMap 등)이 바로 실행되지 않고 그래프(DAG)로 기록
  • 최적화 기회 제공
    • 불필요한 연산 제거, 병합 등 전체 연산 최적화 가능
  • 실행 트리거
    • **sink** 또는 **execute()** 호출 시 DAG 기반으로 전체 연산 한 번에 실행됨
장점설명
효율성불필요한 계산 방지, 자원 절약
성능 향상연산 병합 및 중간 저장 최소화
유연성실행 전까지 실행 계획 수정/변경 가능

관련 개념