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 에서 파생된 입력 포맷에 사용되고, 후자는 범용 입력 포맷에 사용해야 합니다. 결과로 만들어진 InputFormatExecutionEnvironment#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.
[...]

더 알아보기 (Learn more)