Hadoop 포맷
Hadoop 포맷
Hadoop에 대한 지원은 flink-hadoop-compatibility Maven 모듈에 포함되어 있습니다. 이 모듈을 사용하면 Flink 애플리케이션에서 Hadoop InputFormat 을 데이터 소스로 사용할 수 있습니다.
출처: 문서
본문
프로젝트 구성
Hadoop에 대한 지원은 flink-hadoop-compatibility Maven 모듈에 포함되어 있습니다.
hadoop을 사용하려면 pom.xml 에 다음 의존성을 추가하세요:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-hadoop-compatibility</artifactId>
<version>2.3.0</version>
</dependency>
Flink 애플리케이션을 로컬에서(예: IDE에서) 실행하려면 다음과 같은 hadoop-client 의존성도 추가해야 합니다:
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>2.10.2</version>
<scope>provided</scope>
</dependency>
Hadoop InputFormat 사용하기
Flink에서 Hadoop InputFormat 을 사용하려면 먼저 HadoopInputs 유틸리티 클래스의 readHadoopFile 또는 createHadoopInput 을 사용해 포맷을 감싸야 합니다. 전자는 FileInputFormat 에서 파생된 입력 포맷에 사용되고, 후자는 범용 입력 포맷에 사용해야 합니다. 결과로 만들어진 InputFormat 은 ExecutionEnvironment#createInput 을 사용해 데이터 소스를 만들 수 있습니다.
결과 DataStream 은 2-튜플을 포함하는데, 첫 번째 필드는 키(key)이고 두 번째 필드는 Hadoop InputFormat 에서 검색된 값(value)입니다.
다음 예제는 Hadoop의 TextInputFormat 을 사용하는 방법을 보여줍니다:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
KeyValueTextInputFormat textInputFormat = new KeyValueTextInputFormat();
DataStream<Tuple2<Text, Text>> input = env.createInput(HadoopInputs.readHadoopFile(
textInputFormat, Text.class, Text.class, textPath));
// Do something with the data.
[...]