Flink DataStream API 프로그래밍 가이드
Flink DataStream API 프로그래밍 가이드
Flink의 DataStream 프로그램은 데이터 스트림에 대한 변환(예: 필터링, 상태 갱신, 윈도우 정의, 집계)을 구현하는 일반적인 프로그램이에요. 데이터 스트림은 처음에 다양한 소스(예: 메시지 큐, 소켓 스트림, 파일)에서 생성돼요. 결과는 싱크(sink)를 통해 반환되며, 싱크는 예를 들어 데이터를 파일에 쓰거나 표준 출력(예: 명령줄 터미널)으로 출력할 수 있어요. Flink 프로그램은 다양한 컨텍스트에서 실행되는데, standalone으로 또는 다른 프로그램에 임베드(embed)되어 실행될 수 있어요. 실행은 로컬 JVM에서 혹은 여러 머신으로 이루어진 클러스터에서 일어날 수 있어요.
출처: 문서
본문
참고: DataStream API V2는 기존 DataStream API를 점진적으로 대체하기 위한 새로운 API 집합이에요. 현재 실험 단계(experimental stage)에 있으며 생산(production) 환경에서 완전히 사용 가능한 상태는 아니에요.
DataStream이란 무엇인가? (What is a DataStream?)
DataStream API는 Flink 프로그램에서 데이터의 모음을 나타내는 데 사용되는 특별한 DataStream 클래스에서 이름을 얻었어요. 이를 중복을 포함할 수 있는 불변(immutable) 데이터 모음으로 생각할 수 있어요. 이 데이터는 유한하거나 무한할 수 있으며, 이 데이터를 다루는 데 사용하는 API는 동일해요.
DataStream은 사용 방식 측면에서 일반적인 Java Collection과 비슷하지만 몇 가지 중요한 점에서 상당히 달라요. DataStream은 불변(immutable)이라서 한 번 생성되면 요소를 추가하거나 제거할 수 없어요. 또한 내부의 요소를 단순히 검사할 수 없고, 변환(transformation)이라고도 불리는 DataStream API 연산을 통해서만 작업할 수 있어요.
Flink 프로그램에서 소스를 추가해 초기 DataStream을 만들 수 있어요. 그런 다음 process, connectAndProcess 같은 API 메서드를 사용해 이 초기 스트림에서 새 스트림을 파생해 결합할 수 있어요.
기본 프리미티브와 확장 (Fundamental Primitives and Extensions)
그 기능이 Flink에 의해 제공되어야 하는지에 따라 DataStream API의 관련 개념을 두 범주로 나눌 수 있어요: 기본 프리미티브(fundamental primitives)와 고수준 확장(high-level extensions).
기본 프리미티브 (Fundamental primitives)
기본 프리미티브는 상태를 가진 스트림 처리 애플리케이션(stateful stream processing application)을 정의하기 위해 Flink가 제공해야 하는 기본적이고 필수적인 의미론(semantics)이에요. 이것은 프레임워크가 제공하지 않으면 사용자가 달성할 수 없는 것들이에요. 여기에는 데이터 스트림, 파티셔닝, 프로세스 함수(process function), 상태(state), 처리 타이머 서비스(processing timer service), 워터마크(watermark)가 포함돼요.
자세한 내용은 다음에서 찾을 수 있어요:
- Building Blocks: DataStream API의 가장 기본적인 요소를 다뤄요.
- State Processing: 상태를 가진 애플리케이션을 개발하는 방법을 설명해요.
- Time Processing # Processing Timer Service: 프로세싱 타임을 처리하는 방법을 설명해요.
- Watermark: 사용자 정의 이벤트를 정의하고 처리하는 방법을 설명해요.
고수준 확장 (High-Level Extensions)
고수준 확장은 단축키/문법 설탕(sugar)과 같아요. 이것 없이도 사용자가 기본 API를 사용해 동일한 동작을 달성할 수 있겠지만, 내장(builtin) 지원 덕분에 훨씬 쉬워져요. 여기에는 이벤트 타이머 서비스(event timer service), 윈도우(window), 조인(join)이 포함돼요.
자세한 내용은 다음에서 찾을 수 있어요:
- Time Processing # Event Timer Service: 확장을 통해 이벤트 타임을 처리하는 방법을 설명해요.
- Builtin Functions: 확장을 통해 윈도우 집계와 조인을 수행하는 방법을 설명해요.
Flink DataStream 프로그램의 구조 (Anatomy of a Flink DataStream Program)
Flink 프로그램은 DataStream을 변환하는 일반적인 프로그램처럼 보여요. 각 프로그램은 동일한 기본 부분으로 구성돼요:
Execution Environment를 얻는다,- 초기 데이터를 로드/생성한다,
- 이 데이터에 대한 변환을 지정한다,
- 계산 결과를 어디에 둘지 지정한다,
- 프로그램 실행을 트리거한다
Execution Environment 얻기
ExecutionEnvironment는 모든 Flink 프로그램의 기반이에요. ExecutionEnvironment의 다음 정적 메서드를 사용해 얻을 수 있어요:
ExecutionEnvironment env = ExecutionEnvironment.getInstance();
IDE 안에서 또는 일반적인 Java 프로그램으로 프로그램을 실행하면, 로컬 머신에서 프로그램을 실행할 로컬 환경이 만들어져요. 프로그램에서 JAR 파일을 만들어 명령줄을 통해 실행하면, Flink 클러스터 매니저가 main 메서드를 실행하고 ExecutionEnvironment.getInstance()는 클러스터에서 프로그램을 실행하기 위한 실행 환경을 반환해요.
초기 데이터 로드/생성 (Load/create the Initial Data)
소스(source)는 프로그램이 입력을 읽는 곳이에요. ExecutionEnvironment.fromSource(source, sourceName)을 사용해 프로그램에 소스를 연결할 수 있어요. Flink에는 사전 구현된 소스가 여러 개 함께 제공돼요. FLIP-27 기반 소스는 DataStreamV2SourceUtils.wrapSource(source)로 사용하거나, 테스트/디버깅 목적으로 DataStreamV2SourceUtils.fromData(collection)를 사용할 수 있어요.
예를 들어, 사전 정의된 컬렉션에서 데이터만 읽으려면 다음을 사용할 수 있어요:
ExecutionEnvironment env = ExecutionEnvironment.getInstance();
NonKeyedPartitionStream<String> input =
env.fromSource(
DataStreamV2SourceUtils.fromData(Arrays.asList("1", "2", "3")),
"source"
);
소스의 데이터에는 명확한 파티셔닝이 없기 때문에, 이 코드는 NonKeyedPartitionStream을 반환해요. 이 스트림에 변환을 적용해 새로 파생된 DataStream을 만들 수 있어요. 다른 타입의 DataStream에 대해서는 Building Blocks # DataStream을 참고하세요.
이 데이터에 대한 변환 지정 (Specify Transformations on this Data)
ProcssFunction(편의를 위해 여기서는 람다 표현식 사용)을 가진 DataStream의 메서드를 호출해 변환을 적용해요. 예를 들어, map 변환은 다음과 같아요:
NonKeyedPartitionStream<String> input = ...;
NonKeyedPartitionStream<Integer> parsed = source.process(
(OneInputStreamProcessFunction<String, Integer>)
(record, output, ctx) -> {
output.collect(Integer.parseInt(record));
});
이렇게 하면 원래 컬렉션의 모든 String을 Integer로 변환해 새 DataStream을 만들어요. 프로세싱에 대한 자세한 내용은 Building Blocks # Process Function을 참고하세요.
계산 결과를 둘 위치 지정 (Specify Where to Put the Results of Your Computations)
최종 결과를 담은 DataStream을 얻으면, 싱크를 만들어 외부 시스템에 쓸 수 있어요.
데이터 싱크는 DataStream을 소비하고 파일, 소켓, 외부 시스템으로 전달하거나 출력(print)해요. Flink에는 다양한 내장 싱크 구현이 함께 제공돼요. SinkV2 기반 싱크는 DataStreamV2SinkUtils.wrapSink(sink)를 사용할 수 있어요.
다음은 결과를 출력하기 위한 싱크를 만드는 예시예요:
parsed.toSink(DataStreamV2SinkUtils.wrapSink(new PrintSink<>()));
프로그램 실행 트리거 (Trigger the Program Execution)
완전한 프로그램을 지정했다면 ExecutionEnvironment에서 execute()를 호출해 프로그램 실행을 트리거해야 해요.
ExecutionEnvironment의 타입에 따라 실행은 로컬 머신에서 트리거되거나 프로그램을 클러스터에서 실행하도록 제출돼요. execute() 메서드는 작업이 끝날 때까지 기다려요.
프로그램 실행에 관한 이 마지막 부분은 Flink 연산이 언제 그리고 어떻게 실행되는지 이해하는 데 중요해요. 모든 Flink 프로그램은 지연(lazily) 실행돼요: 프로그램의 main 메서드가 실행될 때 데이터 로딩과 변환이 직접 일어나지 않아요. 대신 각 연산이 생성되고 데이터플로우 그래프에 추가돼요. 연산은 실행 환경에서 execute() 호출로 실행이 명시적으로 트리거될 때 실제로 실행돼요. 프로그램이 로컬로 실행되는지 클러스터에서 실행되는지는 실행 환경의 타입에 따라 달라져요.
지연 평가(lazy evaluation) 덕분에 Flink가 하나의 전체적으로 계획된 단위로 실행하는 정교한 프로그램을 구성할 수 있어요.