Flink DataStream API 프로그래밍 가이드
Flink DataStream API 프로그래밍 가이드
Flink의 DataStream 프로그램은 데이터 스트림에 대한 변환(예: 필터링, 상태 갱신, 윈도우 정의, 집계)을 구현하는 일반적인 프로그램이에요. 데이터 스트림은 처음에 다양한 소스(예: 메시지 큐, 소켓 스트림, 파일)에서 생성돼요. 결과는 싱크(sink)를 통해 반환되며, 싱크는 예를 들어 데이터를 파일에 쓰거나 표준 출력(예: 명령줄 터미널)으로 출력할 수 있어요. Flink 프로그램은 standalone 또는 다른 프로그램에 임베드되어 다양한 컨텍스트에서 실행돼요. 실행은 로컬 JVM에서 또는 여러 머신으로 이루어진 클러스터에서 일어날 수 있어요.
출처: 문서
본문
자신만의 Flink DataStream 프로그램을 만들려면 Flink 프로그램의 구조부터 시작하고, 점차 자신만의 스트림 변환을 추가하는 것을 권장해요. 나머지 섹션들은 추가 연산과 고급 기능에 대한 참조 역할을 해요.
DataStream이란 무엇인가? (What is a DataStream?)
DataStream API는 Flink 프로그램에서 데이터의 모음을 나타내는 데 사용되는 특별한 DataStream 클래스에서 이름을 얻었어요. 이를 중복을 포함할 수 있는 불변(immutable) 데이터 모음으로 생각할 수 있어요. 이 데이터는 유한하거나 무한할 수 있으며, 이 데이터를 다루는 데 사용하는 API는 동일해요.
DataStream은 사용 방식 측면에서 일반적인 Java Collection과 비슷하지만 몇 가지 중요한 점에서 상당히 달라요. DataStream은 불변(immutable)이라서 한 번 생성되면 요소를 추가하거나 제거할 수 없어요. 또한 내부의 요소를 단순히 검사할 수 없고, 변환(transformation)이라고도 불리는 DataStream API 연산을 통해서만 작업할 수 있어요.
Flink 프로그램에서 소스를 추가해 초기 DataStream을 만들 수 있어요. 그런 다음 map, filter 같은 API 메서드를 사용해 이 초기 스트림에서 새 스트림을 파생해 결합할 수 있어요.
Flink 프로그램의 구조 (Anatomy of a Flink Program)
Flink 프로그램은 DataStream을 변환하는 일반적인 프로그램처럼 보여요. 각 프로그램은 동일한 기본 부분으로 구성돼요:
execution environment을 얻는다,- 초기 데이터를 로드/생성한다,
- 이 데이터에 대한 변환을 지정한다,
- 계산 결과를 어디에 둘지 지정한다,
- 프로그램 실행을 트리거한다
Java
이제 각 단계에 대한 개요를 설명할게요. 자세한 내용은 각 섹션을 참조하세요. Java DataStream API의 모든 핵심 클래스는 org.apache.flink.streaming.api에서 찾을 수 있다는 점에 유의하세요.
StreamExecutionEnvironment는 모든 Flink 프로그램의 기반이에요. StreamExecutionEnvironment의 다음 정적 메서드를 사용해 얻을 수 있어요:
getExecutionEnvironment();
createLocalEnvironment();
createRemoteEnvironment(String host, int port, String... jarFiles);
일반적으로 getExecutionEnvironment()만 사용하면 돼요. 이것은 컨텍스트에 따라 올바른 일을 하기 때문이에요: IDE 안에서 또는 일반적인 Java 프로그램으로 프로그램을 실행하면, 로컬 머신에서 프로그램을 실행할 로컬 환경이 만들어져요. 프로그램에서 JAR 파일을 만들어 명령줄로 실행하면, Flink 클러스터 매니저가 main 메서드를 실행하고 getExecutionEnvironment()는 클러스터에서 프로그램을 실행하기 위한 실행 환경을 반환해요.
데이터 소스를 지정하기 위해 실행 환경에는 다양한 방법으로 파일을 읽는 여러 메서드가 있어요: 라인별로 읽거나, CSV 파일로 읽거나, 제공된 다른 소스를 사용할 수 있어요. 텍스트 파일을 일련의 라인으로 읽으려면 다음을 사용할 수 있어요:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
FileSource<String> fileSource = FileSource.forRecordStreamFormat(
new TextLineInputFormat(), new Path("file:///path/to/file")
).build();
DataStream<String> text = env.fromSource(
fileSource,
WatermarkStrategy.noWatermarks(),
"file-input"
);
이것은 변환을 적용해 새로 파생된 DataStream을 만들 수 있는 DataStream을 제공해요.
변환 함수를 가진 DataStream의 메서드를 호출해 변환을 적용해요. 예를 들어, map 변환은 다음과 같아요:
DataStream<String> input = ...;
DataStream<Integer> parsed = input.map(new MapFunction<String, Integer>() {
@Override
public Integer map(String value) {
return Integer.parseInt(value);
}
});
이렇게 하면 원래 컬렉션의 모든 String을 Integer로 변환해 새 DataStream을 만들어요.
최종 결과를 담은 DataStream을 얻으면, 싱크를 만들어 외부 시스템에 쓸 수 있어요. 다음은 싱크를 만드는 몇 가지 예시 메서드예요:
stream.sinkTo(
FileSink.forRowFormat(
new Path("outputPath"),
new SimpleStringEncoder<>()
).build()
);
stream.print();
완전한 프로그램을 지정했다면 StreamExecutionEnvironment에서 execute()를 호출해 프로그램 실행을 트리거해야 해요.
ExecutionEnvironment의 타입에 따라 실행은 로컬 머신에서 트리거되거나 프로그램을 클러스터에서 실행하도록 제출돼요.
execute() 메서드는 작업이 끝날 때까지 기다린 후 JobExecutionResult를 반환해요. 이것은 실행 시간과 누산기 결과를 포함해요.
작업이 끝날 때까지 기다리지 않으려면 StreamExecutionEnvironment에서 executeAsync()를 호출해 비동기 작업 실행을 트리거할 수 있어요. 이것은 방금 제출한 작업과 통신할 수 있는 JobClient를 반환해요. 예를 들어, executeAsync()를 사용해 execute()의 의미론을 구현하는 방법은 다음과 같아요.
final JobClient jobClient = env.executeAsync();
final JobExecutionResult jobExecutionResult = jobClient.getJobExecutionResult().get();
프로그램 실행에 관한 이 마지막 부분은 Flink 연산이 언제 그리고 어떻게 실행되는지 이해하는 데 중요해요. 모든 Flink 프로그램은 지연(lazily) 실행돼요: 프로그램의 main 메서드가 실행될 때 데이터 로딩과 변환이 직접 일어나지 않아요. 대신 각 연산이 생성되고 데이터플로우 그래프에 추가돼요. 연산은 실행 환경에서 execute() 호출로 실행이 명시적으로 트리거될 때 실제로 실행돼요. 프로그램이 로컬로 실행되는지 클러스터에서 실행되는지는 실행 환경의 타입에 따라 달라져요.
지연 평가(lazy evaluation) 덕분에 Flink가 하나의 전체적으로 계획된 단위로 실행하는 정교한 프로그램을 구성할 수 있어요.
예시 프로그램 (Example Program)
다음 프로그램은 웹 소켓에서 오는 단어를 5초 윈도우에서 세는 스트리밍 윈도우 단어 개수 애플리케이션의 완전하고 동작하는 예시예요. 코드를 복사 & 붙여넣기해 로컬에서 실행할 수 있어요.
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.time.Duration;
import org.apache.flink.util.Collector;
public class WindowWordCount {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<Tuple2<String, Integer>> dataStream = env
.socketTextStream("localhost", 9999)
.flatMap(new Splitter())
.keyBy(value -> value.f0)
.window(TumblingProcessingTimeWindows.of(Duration.ofSeconds(5)))
.sum(1);
dataStream.print();
env.execute("Window WordCount");
}
public static class Splitter implements FlatMapFunction<String, Tuple2<String, Integer>> {
@Override
public void flatMap(String sentence, Collector<Tuple2<String, Integer>> out) throws Exception {
for (String word: sentence.split(" ")) {
out.collect(new Tuple2<String, Integer>(word, 1));
}
}
}
}
예시 프로그램을 실행하려면 먼저 터미널에서 netcat으로 입력 스트림을 시작해요:
nc -lk 9999
몇 단어를 입력하고 새 단어를 위해 Enter를 눌러요. 이들이 단어 개수 프로그램의 입력이 돼요. 1보다 큰 카운트를 보려면, 5초 안에 같은 단어를 계속 입력하면 돼요(타이핑이 빠르지 않으면 윈도우 크기를 5초에서 늘리세요 ☺).
데이터 소스 (Data Sources)
소스는 프로그램이 입력을 읽는 곳이에요. StreamExecutionEnvironment.addSource(sourceFunction)을 사용해 프로그램에 소스를 연결할 수 있어요. Flink에는 사전 구현된 소스 함수가 여러 개 함께 제공되지만, 비병렬 소스의 경우 SourceFunction을 구현하거나, 병렬 소스의 경우 ParallelSourceFunction 인터페이스를 구현하거나 RichParallelSourceFunction을 확장해 항상 자신만의 커스텀 소스를 작성할 수 있어요.
StreamExecutionEnvironment에서 접근할 수 있는 사전 정의된 스트림 소스가 몇 가지 있어요:
파일 기반:
-
fromSource(FileSource.forRecordStreamFormat(format, paths).build())- 파일에서 레코드 단위로 읽기. -
readFile(fileInputFormat, path)- 지정된 파일 입력 포맷에 따라 파일을 (한 번) 읽기. -
readFile(fileInputFormat, path, watchType, interval, pathFilter, typeInfo)- 앞의 두 메서드가 내부적으로 호출하는 메서드예요. 주어진fileInputFormat에 따라path의 파일을 읽어요. 제공된watchType에 따라 이 소스는 주기적으로(intervalms마다) 새 데이터를 위해path를 모니터링하거나(FileProcessingMode.PROCESS_CONTINUOUSLY), 현재path의 데이터를 한 번 처리하고 종료해요(FileProcessingMode.PROCESS_ONCE).pathFilter를 사용해 파일이 처리되지 않도록 추가로 제외할 수 있어요.구현(IMPLEMENTATION):
내부적으로 Flink는 파일 읽기 과정을 디렉터리 모니터링과 데이터 읽기라는 두 개의 하위 태스크로 분할해요. 각 하위 태스크는 별도의 엔티티로 구현돼요. 모니터링은 단일 비병렬(병렬도 = 1) 태스크로 구현되고, 읽기는 병렬로 실행되는 여러 태스크로 수행돼요. 후자의 병렬도는 작업 병렬도와 같아요. 단일 모니터링 태스크의 역할은 디렉터리를 스캔하고(
watchType에 따라 주기적으로 또는 한 번만), 처리할 파일을 찾고, 이를 splits로 나누고, 이 splits를 다운스트림 리더에 할당하는 것이에요. 리더는 실제 데이터를 읽는 사람들이에요. 각 split은 하나의 리더에 의해서만 읽히지만, 리더는 여러 split을 하나씩 읽을 수 있어요.중요 참고 사항(IMPORTANT NOTES):
watchType가FileProcessingMode.PROCESS_CONTINUOUSLY로 설정되면, 파일이 수정될 때 그 내용이 완전히 재처리돼요. 파일 끝에 데이터를 추가하면 모든 내용이 재처리되므로 이것은 "exactly-once" 의미론을 깨뜨릴 수 있어요.watchType가FileProcessingMode.PROCESS_ONCE로 설정되면, 소스는path를 한 번 스캔하고, 리더가 파일 내용을 읽는 것을 기다리지 않고 종료해요. 물론 리더는 모든 파일 내용이 읽힐 때까지 계속 읽을 거예요. 소스를 닫으면 그 시점 이후로 체크포인트가 없어져요. 이것은 작업이 마지막 체크포인트에서 읽기를 재개하므로 노드 장애 후 복구가 느려질 수 있어요.
소켓 기반:
socketTextStream- 소켓에서 읽기. 요소는 구분자로 분리될 수 있어요.
컬렉션 기반:
fromData(Collection)- Javajava.util.Collection에서 데이터 스트림을 만들어요. 컬렉션의 모든 요소는 같은 타입이어야 해요.fromData(T ...)- 주어진 객체 시퀀스에서 데이터 스트림을 만들어요. 모든 객체는 같은 타입이어야 해요.fromParallelCollection(SplittableIterator, Class)- 반복자에서 병렬로 데이터 스트림을 만들어요. 클래스는 반복자가 반환하는 요소의 데이터 타입을 지정해요.fromSequence(from, to)- 주어진 구간의 숫자 시퀀스를 병렬로 생성해요.
커스텀:
addSource- 새 소스 함수를 연결해요. 예를 들어 Apache Kafka에서 읽으려면addSource(new FlinkKafkaConsumer<>(...))을 사용할 수 있어요. 자세한 내용은 커넥터를 참조하세요.
DataStream 변환 (DataStream Transformations)
사용 가능한 스트림 변환에 대한 개요는 operators를 참조하세요.
데이터 싱크 (Data Sinks)
데이터 싱크는 DataStream을 소비하고 파일, 소켓, 외부 시스템으로 전달하거나 출력(print)해요. Flink는 DataStream의 연산 뒤에 캡슐화된 다양한 내장 출력 포맷을 제공해요:
sinkTo(FileSink.forRowFormat(new Path("outputPath"), new SimpleStringEncoder<>()).build())- 요소를 라인별로 String으로 작성해요. String은 각 요소의 toString() 메서드를 호출해 얻어져요.print()/printToErr()- 각 요소의 toString() 값을 표준 출력/표준 오류 스트림에 출력해요. 선택적으로 출력 앞에 붙는 접두사(msg)를 제공할 수 있어요. 이것은 print 호출을 구분하는 데 도움이 돼요. 병렬도가 1보다 크면 출력 앞에 출력을 생성한 태스크의 식별자도 붙어요.writeUsingOutputFormat()/FileOutputFormat- 커스텀 파일 출력을 위한 메서드와 기본 클래스. 커스텀 객체-바이트 변환을 지원해요.writeToSocket-SerializationSchema에 따라 요소를 소켓에 작성해요.addSink- 커스텀 싱크 함수를 호출해요. Flink에는 싱크 함수로 구현된 외부 시스템에 대한 커넥터(Apache Kafka 같은)가 함께 제공돼요.
DataStream의 write*() 메서드는 주로 디버깅 목적으로 의도됐다는 점에 유의하세요. 이들은 Flink의 체크포인팅에 참여하지 않으므로, 이런 함수들은 보통 at-least-once 의미론을 가져요. 대상 시스템으로의 데이터 플러시는 OutputFormat의 구현에 따라 달라져요. 즉 OutputFormat에 전송된 모든 요소가 즉시 대상 시스템에 나타나는 것은 아니에요. 또한 장애 상황에서 이런 레코드들은 손실될 수 있어요.
파일시스템으로의 신뢰할 수 있는 exactly-once 전달을 위해서는 FileSink를 사용하세요. 또한 .addSink(...) 메서드를 통한 커스텀 구현은 exactly-once 의미론을 위해 Flink의 체크포인팅에 참여할 수 있어요.
실행 파라미터 (Execution Parameters)
StreamExecutionEnvironment는 런타임에 작업별 구성 값을 설정할 수 있는 ExecutionConfig를 포함해요.
대부분의 파라미터에 대한 설명은 execution configuration을 참조하세요. 다음 파라미터들은 특히 DataStream API에 해당해요:
setAutoWatermarkInterval(long milliseconds): 자동 워터마크 방출 간격을 설정해요.long getAutoWatermarkInterval()로 현재 값을 얻을 수 있어요.
결함 허용 (Fault Tolerance)
State & Checkpointing은 Flink의 체크포인팅 메커니즘을 활성화하고 구성하는 방법을 설명해요.
지연 시간 제어 (Controlling Latency)
기본적으로 요소는 하나씩 네트워크로 전송되지 않고(불필요한 네트워크 트래픽을 유발하므로) 버퍼링돼요. 버퍼의 크기(실제로 머신 간에 전송되는)는 Flink 구성 파일에서 설정할 수 있어요. 이 방법은 처리량을 최적화하는 데 좋지만, 들어오는 스트림이 충분히 빠르지 않으면 지연 문제를 일으킬 수 있어요. 처리량과 지연 시간을 제어하려면 실행 환경(또는 개별 연산자)에서 env.setBufferTimeout(timeoutMillis)을 사용해 버퍼가 채워지는 최대 대기 시간을 설정할 수 있어요. 이 시간이 지나면 버퍼가 가득 차지 않아도 자동으로 전송돼요. 이 타임아웃의 기본값은 100ms예요.
사용법:
LocalStreamEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
env.setBufferTimeout(timeoutMillis);
env.generateSequence(1,10).map(new MyMapper()).setBufferTimeout(timeoutMillis);
최대 처리량을 위해 setBufferTimeout(-1)을 설정해 타임아웃을 제거하고 버퍼가 가득 찼을 때만 플러시되게 해요. 최소 지연을 위해 타임아웃을 0에 가까운 값(예: 5 또는 10ms)으로 설정해요. 버퍼 타임아웃 0은 심각한 성능 저하를 일으킬 수 있으므로 피해야 해요.
디버깅 (Debugging)
분산 클러스터에서 스트리밍 프로그램을 실행하기 전에, 구현한 알고리즘이 원하는 대로 동작하는지 확인하는 것이 좋아요. 따라서 데이터 분석 프로그램을 구현하는 것은 보통 결과 확인, 디버깅, 개선의 점진적 과정이에요.
Flink는 IDE 내 로컬 디버깅, 테스트 데이터 주입, 결과 데이터 수집을 지원해 데이터 분석 프로그램의 개발 과정을 크게 수월하게 해주는 기능을 제공해요. 이 섹션은 Flink 프로그램의 개발을 수월하게 하는 몇 가지 힌트를 제공해요.
로컬 실행 환경 (Local Execution Environment)
LocalStreamEnvironment는 생성된 동일한 JVM 프로세스 내에서 Flink 시스템을 시작해요. IDE에서 LocalEnvironment를 시작하면 코드에 중단점(breakpoint)을 설정하고 프로그램을 쉽게 디버깅할 수 있어요.
LocalEnvironment는 다음과 같이 생성되고 사용돼요:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
DataStream<String> lines = env.addSource(/* some source */);
// build your program
env.execute();
컬렉션 데이터 소스 (Collection Data Sources)
Flink는 테스트를 쉽게 하기 위해 Java 컬렉션에 의해 뒷받침되는 특수 데이터 소스를 제공해요. 프로그램이 테스트된 후, 소스와 싱크는 외부 시스템에서 읽고/쓰는 소스와 싱크로 쉽게 대체될 수 있어요.
컬렉션 데이터 소스는 다음과 같이 사용될 수 있어요:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
// Create a DataStream from a list of elements
DataStream<Integer> myInts = env.fromElements(1, 2, 3, 4, 5);
// Create a DataStream from any Java collection
List<Tuple2<String, Integer>> data = ...
DataStream<Tuple2<String, Integer>> myTuples = env.fromCollection(data);
// Create a DataStream from an Iterator
Iterator<Long> longIt = ...;
DataStream<Long> myLongs = env.fromCollection(longIt, Long.class);
참고: 현재 컬렉션 데이터 소스는 데이터 타입과 반복자가 Serializable을 구현해야 해요. 또한 컬렉션 데이터 소스는 병렬로 실행될 수 없어요(병렬도 = 1).
반복자 데이터 싱크 (Iterator Data Sink)
Flink는 테스트와 디버깅 목적으로 DataStream 결과를 수집하는 싱크도 제공해요. 다음과 같이 사용할 수 있어요:
DataStream<Tuple2<String, Integer>> myResult = ...;
Iterator<Tuple2<String, Integer>> myOutput = myResult.collectAsync();
다음 단계는 어디로? (Where to go next?)
- Operators: 사용 가능한 스트리밍 연산자의 사양.
- Event Time: Flink의 시간 개념에 대한 소개.
- State & Fault Tolerance: 상태를 가진 애플리케이션을 개발하는 방법 설명.
- Connectors: 사용 가능한 입력 및 출력 커넥터 설명.