DataStream API 배우기
DataStream API 배우기 (Learn the DataStream API)
이 훈련의 초점은 스트리밍 애플리케이션 작성을 시작할 수 있을 만큼 DataStream API를 폭넓게 다루는 것이에요.
출처: 문서
본문
무엇을 스트리밍할 수 있나요? (What can be Streamed?)
Flink의 DataStream API는 직렬화 가능한 모든 것을 스트리밍하게 해줘요. Flink의 자체 직렬화기는 다음에 사용돼요:
- 기본 타입, 즉 String, Long, Integer, Boolean, Array
- 복합 타입: Tuples, POJOs
그리고 다른 타입에는 Flink가 Kryo로 대체(fall back)해요. Flink와 함께 다른 직렬화기를 사용하는 것도 가능해요. 특히 Avro는 잘 지원돼요.
Java tuples와 POJOs
Flink의 네이티브 직렬화기는 튜플과 POJO에 효율적으로 동작할 수 있어요.
Tuples
Java의 경우 Flink는 자체 Tuple0부터 Tuple25 타입을 정의해요.
Tuple2<String, Integer> person = Tuple2.of("Fred", 35);
// zero based index!
String name = person.f0;
Integer age = person.f1;
POJOs
Flink는 다음 조건이 충족되면 데이터 타입을 POJO 타입으로 인식하고(그리고 "이름 기준" 필드 참조를 허용하고):
- 클래스가 public이고 standalone(비정적 내부 클래스 아님)
- 클래스에 public 무인자 생성자
- 클래스(및 모든 상위 클래스)의 모든 비정적, 비일시적 필드가 public(그리고 비-final)이거나, getter와 setter에 Java beans 명명 규칙을 따르는 public getter- 및 setter- 메서드가 있음
예:
public class Person {
public String name;
public Integer age;
public Person() {}
public Person(String name, Integer age) {
. . .
}
}
Person person = new Person("Fred Flintstone", 35);
Flink의 직렬화기는 POJO 타입에 대한 스키마 진화를 지원해요.
완전한 예제 (A Complete Example)
이 예제는 사람에 대한 레코드 스트림을 입력으로 받아 성인만 포함하도록 필터링해요.
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.api.common.functions.FilterFunction;
public class Example {
public static void main(String[] args) throws Exception {
final StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<Person> flintstones = env.fromData(
new Person("Fred", 35),
new Person("Wilma", 35),
new Person("Pebbles", 2));
DataStream<Person> adults = flintstones.filter(new FilterFunction<Person>() {
@Override
public boolean filter(Person person) throws Exception {
return person.age >= 18;
}
});
adults.print();
env.execute();
}
public static class Person {
public String name;
public Integer age;
public Person() {}
public Person(String name, Integer age) {
this.name = name;
this.age = age;
}
public String toString() {
return this.name.toString() + ": age " + this.age.toString();
}
}
}
스트림 실행 환경 (Stream execution environment)
모든 Flink 애플리케이션에는 실행 환경(이 예제의 env)이 필요해요. 스트리밍 애플리케이션은 StreamExecutionEnvironment를 사용해야 해요. 애플리케이션에서 만든 DataStream API 호출은 StreamExecutionEnvironment에 연결되는 작업 그래프를 만듭니다. env.execute()가 호출되면 이 그래프는 패키징되어 JobManager로 전송되고, JobManager는 작업을 병렬화해 그 조각을 TaskManager에 분산해 실행해요. 작업의 각 병렬 조각은 task slot에서 실행돼요.
execute()를 호출하지 않으면 애플리케이션이 실행되지 않는다는 점에 주의해요. 이 분산 런타임은 애플리케이션이 직렬화 가능함과 모든 의존성이 클러스터의 각 노드에 있어야함에 의존해요.
기본 스트림 소스 (Basic stream sources)
위 예제는 env.fromData(...)로 DataStream<Person>을 만듭니다. 이는 프로토타입이나 테스트에 쓸 간단한 스트림을 구성하는 편리한 방법이에요. 프로토타이핑 중에 스트림에 데이터를 넣는 또 다른 편리한 방법은 소켓을 사용하는 것이에요:
DataStream<String> lines = env.socketTextStream("localhost", 9999);
또는 파일:
FileSource<String> fileSource = FileSource.forRecordStreamFormat(
new TextLineInputFormat(), new Path("file:///path")
).build();
DataStream<String> lines = env.fromSource(
fileSource,
WatermarkStrategy.noWatermarks(),
"file-input"
);
실제 애플리케이션에서 가장 흔히 쓰는 데이터 소스는 저지연, 고처리량 병렬 읽기와 되감기·재생(rewind and replay)을 지원하는 것들로, 높은 성능과 장애 허용의 전제 조건이며 Apache Kafka, Kinesis, 다양한 파일 시스템이 그 예예요. REST API와 데이터베이스도 스트림 보강(stream enrichment)에 자주 사용돼요.
기본 스트림 싱크 (Basic stream sinks)
위 예제는 adults.print()로 결과를 task manager 로그에 출력해요(IDE에서 실행하면 IDE의 콘솔에 나타남). 이는 스트림의 각 요소에 toString()을 호출합니다. 출력은 다음과 같아요:
1> Fred: age 35
2> Wilma: age 35
여기서 1>과 2>는 어느 하위 작업(즉, 스레드)이 출력을 만들었는지 나타내요. 프로덕션에서 흔히 쓰이는 싱크에는 FileSink, 다양한 데이터베이스, 여러 pub-sub 시스템이 있어요.
디버깅 (Debugging)
프로덕션에서는 애플리케이션이 원격 클러스터나 컨테이너 집합에서 실행돼요. 실패하면 원격에서 실패합니다. JobManager와 TaskManager 로그는 그러한 실패를 디버깅하는 데 매우 유용하지만, IDE 안에서 로컬 디버깅을 하는 것이 훨씬 쉬우며 Flink가 이를 지원해요. 중단점을 설정하고, 로컬 변수를 검사하고, 코드를 단계별로 실행할 수 있어요. Flink의 코드로도 들어갈 수 있는데, Flink가 어떻게 동작하는지 궁금하다면 내부를 배우는 좋은 방법이 될 수 있어요.
실습 (Hands-on)
이제 간단한 DataStream 애플리케이션을 코딩하고 실행하기 시작할 만큼 충분히 알게 됐어요. flink-training-repo를 클론하고, README의 지침을 따른 후 첫 번째 연습인 Filtering a Stream (Ride Cleansing)을 해보세요.
추가 자료 (Further Reading)
- Flink Serialization Tuning Vol. 1: Choosing your Serializer — if you can
- Anatomy of a Flink Program
- Data Sources
- Data Sinks