Learn Flink: Hands-On Training
Learn Flink: Hands-On Training
이 교육 과정은 확장 가능한 스트리밍 ETL, 분석, 이벤트 기반 애플리케이션을 작성하기 시작하는 데 필요한 만큼의 Apache Flink 소개를 제공합니다. 상태와 시간을 다루는 Flink API에 대한 직관적인 소개에 집중합니다.
출처: 문서
본문
이 교육의 목표와 범위
이 교육은 확장 가능한 스트리밍 ETL, 분석, 이벤트 기반 애플리케이션을 작성하기 시작하는 데 필요한 만큼(궁극적으로 중요한 많은 세부사항은 생략하면서) Apache Flink에 대한 소개를 제공합니다. 상태와 시간을 관리하는 Flink API에 대한 직관적인 소개를 제공하는 데 초점을 맞추며, 이러한 기초를 마스터하면 더 자세한 참조 문서에서 나머지 필요한 내용을 익히는 데 훨씬 더 잘 준비될 것입니다. 각 섹션 끝의 링크는 더 배울 수 있는 곳으로 안내합니다.
구체적으로 다음을 학습하게 됩니다:
- 스트리밍 데이터 처리 파이프라인을 구현하는 방법
- Flink가 상태를 관리하는 방법과 이유
- 이벤트 시간(event time)을 사용해 정확한 분석을 일관되게 계산하는 방법
- 연속 스트림 위에서 이벤트 기반 애플리케이션을 구축하는 방법
- Flink가 exactly-once 의미론으로 내결함성 있는 상태 기반 스트림 처리를 제공할 수 있는 방법
이 교육은 스트리밍 데이터의 연속 처리, 이벤트 시간, 상태 기반 스트림 처리, 상태 스냅샷이라는 네 가지 핵심 개념에 초점을 맞춥니다. 이 페이지는 이러한 개념을 소개합니다.
참고: 이 교육과 함께 제시되는 개념을 다루는 방법을 배우도록 안내하는 실습 연습 세트가 제공됩니다. 관련 연습에 대한 링크는 각 섹션 끝에 제공됩니다.
스트림 처리 (Stream Processing)
스트림은 데이터의 자연스러운 서식지입니다. 웹 서버의 이벤트, 주식 거래소의 거래, 공장 바닥의 기계에서 나오는 센서 판독값 등 데이터는 스트림의 일부로 생성됩니다. 하지만 데이터를 분석할 때 처리를 bounded(유한) 또는 unbounded(무한) 스트림을 기준으로 구성할 수 있으며, 이 패러다임 중 어떤 것을 선택하느냐는 심대한 결과를 가져옵니다.
배치 처리(Batch processing) 는 유한(bounded) 데이터 스트림을 처리할 때 작동하는 패러다임입니다. 이 작동 모드에서는 결과를 생성하기 전에 전체 데이터셋을 수집하도록 선택할 수 있습니다. 즉 데이터를 정렬하고, 전역 통계를 계산하고, 모든 입력을 요약하는 최종 보고서를 생성하는 것이 예를 들어 가능합니다.
반면 스트림 처리(Stream processing) 는 무한(unbounded) 데이터 스트림을 다룹니다. 개념적으로는 입력이 끝나지 않을 수 있으므로, 데이터가 도착할 때 계속해서 처리해야 합니다.
Flink에서 애플리케이션은 사용자 정의 operator 에 의해 변환될 수 있는 스트리밍 데이터플로우(streaming dataflows) 로 구성됩니다. 이러한 데이터플로우는 하나 이상의 source 로 시작해 하나 이상의 sink 로 끝나는 유향 그래프(directed graphs)를 형성합니다.
프로그램의 변환과 데이터플로우의 operator 사이에는 종종 일대일 대응이 있습니다. 그러나 때로는 하나의 변환이 여러 operator로 구성될 수 있습니다.
애플리케이션은 Apache Kafka나 Kinesis 같은 메시지 큐 또는 분산 로그와 같은 스트리밍 소스에서 실시간 데이터를 소비할 수 있습니다. 또한 Flink는 다양한 데이터 소스에서 유한(bounded)한 과거 데이터를 소비할 수도 있습니다. 마찬가지로 Flink 애플리케이션이 생성하는 결과 스트림은 sink 로 연결될 수 있는 다양한 시스템으로 전송될 수 있습니다.
병렬 데이터플로우 (Parallel Dataflows)
Flink의 프로그램은 본질적으로 병렬적이고 분산되어 있습니다. 실행 중 스트림은 하나 이상의 스트림 파티션(stream partitions) 을 가지며, 각 operator 는 하나 이상의 operator subtask 를 가집니다. operator subtask 는 서로 독립적이며 서로 다른 스레드에서, 그리고 다른 머신이나 컨테이너에서 실행될 수 있습니다.
operator subtask 의 수는 해당 operator 의 병렬도(parallelism) 입니다. 같은 프로그램의 서로 다른 operator 는 서로 다른 병렬도 수준을 가질 수 있습니다.
스트림은 두 operator 사이에서 데이터를 one-to-one(또는 forwarding) 패턴 또는 redistributing 패턴으로 전송할 수 있습니다:
- One-to-one 스트림(예: 위 그림의 Source 와 map() operator 사이)은 요소의 파티셔닝과 순서를 보존합니다. 즉 map() operator 의 subtask[1] 는 Source operator 의 subtask[1] 이 생성한 것과 같은 순서로 같은 요소를 보게 됩니다.
- Redistributing 스트림(위 그림에서 map() 과 keyBy/window 사이, 그리고 keyBy/window 와 Sink 사이)은 스트림의 파티셔닝을 변경합니다. 각 operator subtask 는 선택된 변환에 따라 서로 다른 대상 subtask 에 데이터를 보냅니다. 예로는 keyBy()(키를 해싱하여 재파티셔닝), broadcast(), rebalance()(무작위로 재파티셔닝) 등이 있습니다. 재분배 교환에서 요소 간의 순서는 보내는 subtask 와 받는 subtask 각 쌍 내에서만 보존됩니다(예: map() 의 subtask[1] 과 keyBy/window 의 subtask[2]). 따라서 위에 표시된 keyBy/window 와 Sink operator 사이의 재분배는 서로 다른 키에 대한 집계 결과가 Sink 에 도착하는 순서에 대해 비결정성을 도입합니다.
적시의 스트림 처리 (Timely Stream Processing)
대부분의 스트리밍 애플리케이션에서 과거 데이터를 처리할 때도 라이브 데이터 처리에 사용하는 것과 같은 코드로 재처리할 수 있고, 그와 관계없이 결정적이고 일관된 결과를 생성하는 것은 매우 가치 있습니다.
또한 이벤트가 처리되기 위해 전달된 순서가 아니라 발생한 순서에 주의를 기울이고, 일련의 이벤트가 언제 (또는 언제 완료되어야 하는지) 추론할 수 있는 것도 중요할 수 있습니다. 예를 들어 전자상거래 거래나 금융 거래에 관련된 이벤트 집합을 생각해 보세요.
적시의 스트림 처리를 위한 이러한 요구사항은 데이터를 처리하는 머신의 시계를 사용하는 대신 데이터 스트림에 기록된 이벤트 시간 타임스탬프를 사용함으로써 충족할 수 있습니다.
상태 기반 스트림 처리 (Stateful Stream Processing)
Flink의 연산은 상태를 가질 수 있습니다. 즉 하나의 이벤트가 처리되는 방식이 그 이전에 온 모든 이벤트의 누적된 효과에 의존할 수 있습니다. 상태는 대시보드에 표시하는 분당 이벤트 수를 세는 것 같은 단순한 것부터, 사기 탐지 모델의 특성을 계산하는 것 같은 더 복잡한 것까지 다양한 용도로 사용될 수 있습니다.
Flink 애플리케이션은 분산 클러스터에서 병렬로 실행됩니다. 주어진 operator 의 다양한 병렬 인스턴스는 별도의 스레드에서 독립적으로 실행되며 일반적으로 서로 다른 머신에서 실행됩니다.
상태 기반 operator 의 병렬 인스턴스 집합은 사실상 샤딩된(sharded) 키-값 저장소입니다. 각 병렬 인스턴스는 특정 키 그룹에 대한 이벤트를 처리할 책임이 있으며 해당 키에 대한 상태는 로컬에 유지됩니다.
아래 다이어그램은 작업 그래프에서 처음 세 operator 에 대해 병렬도 2로 실행되고 병렬도 1의 sink 로 끝나는 작업을 보여줍니다. 세 번째 operator 는 상태 기반이며 두 번째와 세 번째 operator 사이에 완전 연결 네트워크 셔플이 발생하는 것을 볼 수 있습니다. 이는 함께 처리되어야 하는 모든 이벤트가 함께 처리되도록 스트림을 일부 키로 파티셔닝하기 위해 수행됩니다.
상태는 항상 로컬에서 접근되며, 이는 Flink 애플리케이션이 높은 처리량과 낮은 지연 시간을 달성하는 데 도움이 됩니다. 상태를 JVM 힙에 유지하거나, 너무 크다면 효율적으로 구성된 디스크 기반 데이터 구조에 유지하도록 선택할 수 있습니다.
상태 스냅샷을 통한 내결함성 (Fault Tolerance via State Snapshots)
Flink는 상태 스냅샷과 스트림 재생(replay)의 결합을 통해 내결함성을 갖춘 exactly-once 의미론을 제공할 수 있습니다. 이러한 스냅샷은 분산 파이프라인의 전체 상태를 캡처하며, 입력 큐로의 오프셋과 그 시점까지 데이터를 수집한 결과 발생한 작업 그래프 전체의 상태를 기록합니다. 실패가 발생하면 source 가 되감기고, 상태가 복원되며, 처리가 재개됩니다. 위에서 설명한 대로 이러한 상태 스냅샷은 진행 중인 처리를 방해하지 않고 비동기적으로 캡처됩니다.