본문으로 건너뛰기
김신건의 로그

[AWS] Managed Service for Apache Flink (구 Kinesis Data Analytics)

· 수정 · 📖 약 3분 · 1,182자/단어 #aws #cloud #kinesis #flink #streaming #real-time #analytics
AWS Kinesis Data Analytics, Kinesis Data Analytics, KDA, Amazon Managed Service for Apache Flink, Managed Flink, MSF, AWS Managed Flink, Flink Studio, Kinesis Data Analytics for SQL

정의

Amazon Managed Service for Apache Flink (구 Amazon Kinesis Data Analytics) 는 스트리밍 데이터를 실시간으로 처리 (aggregate, join, window, ML 추론) 하는 완전 관리형 Apache Flink 서비스. Java / Python / Scala / SQL 로 스트림 애플리케이션 작성.

이름 변경 & SQL 종료 (2025-2026)

IMPORTANT

Kinesis Data Analytics for SQL Applications 는 2026-01-27 종료. 신규 생성 불가, 기존 앱 모두 삭제. SQL 애플리케이션 사용자는 Managed Flink Studio (Zeppelin + Flink SQL) 로 마이그레이션.

시기별 정리:

시기이름 / 상태
~2018Kinesis Data Analytics for SQL Applications (SQL)
2018-2023Kinesis Data Analytics for Apache Flink (Java/Python/Scala) 추가
2023-08Amazon Managed Service for Apache Flink 로 이름 변경
2025-09-01KDA for SQL - 제한적 지원 시작
2025-10-15KDA for SQL - 신규 생성 차단
2026-01-27KDA for SQL 모든 앱 삭제, 서비스 종료

지금은 Managed Service for Apache Flink (MSF) 하나만.

스트리밍 처리는 배치 처리보다 어려운 문제들이 있음:

  • Event time vs Processing time - 이벤트가 언제 발생했는지 vs 언제 처리됐는지
  • Late-arriving events - 늦게 도착한 데이터 처리
  • Watermark - “이 시점 이전 이벤트는 다 왔다” 신호
  • Exactly-once semantics - 정확히 한 번 처리
  • Stateful processing - 세션, 집계 상태 관리
  • Windowing - Tumbling / Sliding / Session window

Flink 는 이 모든 것을 first-class 지원. Kafka Streams, Spark Streaming 대비 저지연 + stateful 강함.

지원 언어

언어특징
Java원류, 가장 완전. Flink DataStream API
PythonPyFlink. Java API wrapper
Scala사용 감소 (Flink 2.x 부터 core 제거)
SQLFlink SQL. Studio (Zeppelin notebook) 에서

지원 버전 (2026 시점)

  • Flink 1.20.0 (최신, 2026 기준)
  • Flink 1.19.1
  • Flink 1.18.1

LTS 정책: 각 버전 GA 시점부터 24개월 지원.

아키텍처

flowchart LR
    subgraph Sources
        KDS["Kinesis Data Streams"]
        MSK[Amazon MSK]
        Kafka[Self-managed Kafka]
        S3In[("S3 파일")]
    end

    Sources --> MSF["Managed Flink Application<br/>(Flink Job)"]

    MSF --> Sinks

    subgraph Sinks
        KDSOut["KDS"]
        FH["Data Firehose"]
        S3Out[("S3")]
        DDB[("DynamoDB")]
        OS[OpenSearch]
        RDS[("RDS/Aurora")]
        Custom[Custom HTTP]
    end

    State["Application State<br/>(RocksDB, 자동 체크포인트)"] <-.-> MSF
    CW["CloudWatch<br/>Logs/Metrics"] <-.-> MSF
  • Application: Flink job (JAR 또는 Python)
  • KPU (Kinesis Processing Unit): 1 vCPU + 4 GB. 병렬성 단위
  • State: RocksDB 기반, S3 로 자동 체크포인트 / 스냅샷
  • Auto-scaling: 스트림 트래픽에 따라 KPU 자동 조정

실행 유형

1. Application (프로덕션)

배포된 JAR / Python 을 24/7 실행.

Application configuration:
  runtime: FLINK-1_20
  code_location: s3://my-bucket/flink-app.jar
  parallelism: 4
  parallelismPerKPU: 1
  autoScalingEnabled: true

2. Studio (Zeppelin Notebook)

인터랙티브 개발/탐색. Flink SQL, PyFlink, Scala 지원.

-- Flink SQL 예: 5분 tumbling window
CREATE TABLE clicks (
  user_id BIGINT,
  page STRING,
  event_time TIMESTAMP(3),
  WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND
) WITH (
  'connector' = 'kinesis',
  'stream' = 'clicks',
  'format' = 'json'
);

SELECT
  page,
  TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start,
  COUNT(*) AS click_count
FROM clicks
GROUP BY page, TUMBLE(event_time, INTERVAL '5' MINUTE);

Studio 는 SQL 사용자의 마이그레이션 답. KDA for SQL 을 대체.

커넥터 (Connectors)

40+ 개 소스/싱크:

카테고리목록
AWSKinesis Data Streams, MSK, Firehose, S3, DynamoDB, RDS (JDBC), OpenSearch, CloudWatch
KafkaSelf-managed Kafka, Confluent
Message QueueRabbitMQ, ActiveMQ
DBJDBC (PostgreSQL, MySQL), Iceberg, Hudi, Delta
HTTPAsync HTTP sink

윈도우 (Windowing) 예제

Tumbling Window

겹치지 않는 고정 크기.

[--5분--][--5분--][--5분--]

Sliding Window

겹치는 크기.

[--5분--]
   [--5분--]
      [--5분--]

Session Window

이벤트 간 gap 기반.

[event][event] .......gap....... [event][event][event]
        ↑↑ 하나의 세션              ↑↑↑ 다른 세션

Global Window

전체 스트림 하나 (트리거 커스텀).

DataStream<Event> input = ...;
input
  .keyBy(Event::getUserId)
  .window(TumblingEventTimeWindows.of(Time.minutes(5)))
  .aggregate(new SumAggregate())
  .print();

Exactly-Once

Flink 체크포인트 + 커넥터 트랜잭션 지원.

  • 소스 (KDS, Kafka) offset 저장
  • 싱크 (S3, Kafka) transactional write
  • 실패 시 마지막 체크포인트에서 재실행

요금

  • KPU-hour: ~$0.11 / KPU-hour (Application)
  • Studio: ~$0.14 / KPU-hour
  • Running Application Storage: ~$0.10 / GB-월 (state)
  • Durable Application Backups: ~$0.023 / GB-월 (스냅샷)

계산 예:

Application 4 KPU × 730 시간 = 2920 KPU-hour
2920 × $0.11 = $321.20 /월

실무 시나리오

시나리오
실시간 이상 감지 (fraud detection)Managed Flink DataStream (Java)
5분 tumbling aggregationManaged Flink SQL (Studio)
CDC 로 KDS -> Iceberg 실시간Managed Flink + Iceberg connector
세션 window 분석Managed Flink Session Windows
배치 aggregationAthena / EMR (Flink 아님)
단순 적재 (변환 없음)Data Firehose (Flink 필요 X)
엔진특징언제
Flink저지연, exactly-once, stateful 강함복잡한 스트리밍 처리
Spark Structured Streaming배치와 통합, 익숙이미 Spark 스택
Kafka StreamsKafka 안에서만, 경량Kafka 전용
KSQL / ksqlDBSQL only, 단순Confluent 환경
Firehose변환 없이 적재코딩 회피

함정

WARNING

KDA for SQL 은 2026-01-27 종료. 신규 개발 금지. 기존 SQL 앱은 Managed Flink Studio (Zeppelin) 로 마이그레이션 계획 수립.

CAUTION

Flink 학습 곡선 급함. Event time / watermark / windowing 개념 이해 없이 도입하면 실패. 배치 사고 방식 X.

WARNING

KPU 과잉 프로비저닝 = 비용 폭탄. Auto-scaling + 실측 후 조정.

IMPORTANT

State 크기 증가 = 체크포인트 지연 + 비용. Session window state TTL 설정, keyed state 크기 모니터링.

CAUTION

Late event handling 미설정 = 데이터 조용히 drop. Watermark 전략 + allowedLateness 정의.

WARNING

Studio 는 개발/탐색용. 24/7 프로덕션은 Application 배포로. Studio 는 노트북 종료 시 앱도 중단.

IMPORTANT

Flink 2.x 는 Scala 제거. 새 프로젝트는 Java / Python.

관련 위키

이 글의 용어 (8개)
[AWS] Data Firehose: 서버리스 스트림 적재cloud
정의 Amazon Data Firehose (구 Kinesis Data Firehose, 2024-02 이름 변경) 는 스트리밍 데이터를 받아 S3/Redshift/OpenSea…
[AWS] EMR: 관리형 빅데이터 클러스터cloud
정의 Amazon EMR (Elastic MapReduce) 은 Hadoop / Spark 등 오픈소스 빅데이터 프레임워크를 관리형 클러스터로 제공 하는 AWS 서비스. AWS …
[AWS] Glue: 서버리스 ETL + Data Catalogcloud
정의 AWS Glue 는 서버리스 ETL + 통합 메타데이터 카탈로그 플랫폼. Apache Spark 로 데이터 변환을 실행하고, Hive 호환 Data Catalog 로 스키마…
[AWS] Kinesis Data Streams: 실시간 스트리밍 수집cloud
정의 Amazon Kinesis Data Streams (KDS) 는 실시간 스트리밍 데이터를 대규모로 수집/저장하고, 여러 소비자가 각자 실시간 처리/재처리 할 수 있는 관리형…
[AWS] Lambda: 서버리스 함수, 트리거, 동시성cloud
정의 AWS Lambda = 서버리스 함수 실행. 이벤트 트리거 → 함수 실행 → 결과 / 비동기 처리. 서버 관리 0. 사용 상황 | 상황 | Lambda 적합성 | |---|…
[AWS] S3: object storage, storage classes, lifecyclecloud
정의 S3 = AWS 의 object storage. bucket + key + object. 11 9's durability (99.999999999%), 무한 확장. 2026…
Apache Parquetdata-engineering
정의 Apache Parquet 은 분석 쿼리에 최적화된 오픈소스 컬럼형 이진 파일 포맷. 2013년 Twitter + Cloudera 가 Google Dremel 논문 (201…
ETL / ELTdata-engineering
정의 ETL = Extract (추출) + Transform (변환) + Load (적재). 여러 소스에서 데이터를 뽑아 정제한 뒤 데이터 웨어하우스 에 저장하는 파이프라인. E…

💬 댓글

사이트 검색 / 명령어

검색

스크롤 = 확대/축소 · 드래그 = 이동 · 0 = 원래 크기 · ESC = 닫기