HCatalog 입출력 인터페이스
HCatalog 입출력 인터페이스 (HCatalog Input and Output Interfaces)
HCatInputFormat과 HCatOutputFormat은 MapReduce 잡에서 HCatalog 관리 테이블의 데이터를 읽고 쓰는 인터페이스예요. 별도의 HCatalog 전용 설정은 필요 없고, 테이블의 기본 OutputFormat을 사용하며 잡 완료 후 새 파티션을 테이블에 게시해줘요. HCatRecord 타입 매핑과 라이브러리 배송 방법까지 함께 다룬답니다.
출처: 문서
본문
설정 (Set Up)
HCatInputFormat과 HCatOutputFormat 인터페이스에는 HCatalog 전용 설정이 필요 없어요.
참고: HCatalog는 스레드 안전하지 않아요.
HCatInputFormat
HCatInputFormat은 HCatalog 관리 테이블의 데이터를 읽기 위해 MapReduce 잡과 함께 사용돼요.
HCatInputFormat은 데이터가 테이블에 게시된 것처럼 읽기 위한 Hadoop 0.20 MapReduce API를 노출해요.
API
HCatInputFormat이 노출하는 API는 아래와 같아요. 포함 사항:
setInputsetOutputSchemagetTableSchema
HCatInputFormat으로 데이터를 읽으려면 먼저 읽는 테이블의 필요한 정보로 InputJobInfo를 인스턴스화한 다음 InputJobInfo로 setInput을 호출해요.
setOutputSchema 메서드를 사용해 출력 필드를 지정하는 프로젝션 스키마를 포함할 수 있어요. 스키마를 지정하지 않으면 테이블의 모든 컬럼이 반환돼요.
getTableSchema 메서드를 사용해 지정된 입력 테이블의 테이블 스키마를 결정할 수 있어요.
/**
* 잡에 사용할 입력을 설정한다. 지정된 파티션 술어로 메타데이터 서버를
* 조회하고, 일치하는 파티션을 얻고, conf 객체에 정보를 넣는다.
* inputInfo 객체는 클라이언트 컨텍스트에 필요한 정보로 갱신된다.
* @param job 잡 객체
* @param inputJobInfo 읽을 테이블의 입력 정보
* @throws IOException 메타데이터 서버와의 통신 예외
*/
public static void setInput(Job job,
InputJobInfo inputJobInfo) throws IOException;
/**
* HCatInputFormat이 반환하는 HCatRecord 데이터의 스키마를 설정한다.
* @param job 잡 객체
* @param hcatSchema 통합 스키마로 사용할 스키마
*/
public static void setOutputSchema(Job job,HCatSchema hcatSchema)
throws IOException;
/**
* 지정된 잡 컨텍스트에서 HCatInputFormat.setInput 호출에 지정된 테이블의
* HCatTable 스키마를 얻는다. 이 정보는 JobContext에 대해 HCatInputFormat.setInput이
* 호출된 후에만 사용할 수 있다.
* @param context 컨텍스트
* @return 테이블 스키마
* @throws IOException 현재 컨텍스트에 대해 HCatInputFormat.setInput이 호출되지 않은 경우
*/
public static HCatSchema getTableSchema(JobContext context)
throws IOException;
HCatOutputFormat
HCatOutputFormat은 HCatalog 관리 테이블에 데이터를 쓰기 위해 MapReduce 잡과 함께 사용돼요.
HCatOutputFormat은 테이블에 데이터를 쓰기 위한 Hadoop 0.20 MapReduce API를 노출해요. MapReduce 잡이 HCatOutputFormat으로 출력을 쓰면 테이블에 대해 구성된 기본 OutputFormat이 사용되고, 잡 완료 후 새 파티션이 테이블에 게시돼요.
API
HCatOutputFormat이 노출하는 API는 아래와 같아요. 포함 사항:
setOutputsetSchemagetTableSchema
HCatOutputFormat에 대한 첫 번째 호출은 setOutput이어야 해요. 다른 호출은 출력 포맷이 초기화되지 않았다는 예외를 던져요. 쓰여지는 데이터의 스키마는 setSchema 메서드로 지정돼요. 쓰는 데이터의 스키마를 제공하는 이 메서드를 호출해야 해요. 데이터가 테이블 스키마와 같은 스키마를 가지면 HCatOutputFormat.getTableSchema()로 테이블 스키마를 얻은 다음 그것을 setSchema()에 전달할 수 있어요.
/**
* 잡에 대한 출력 정보를 설정한다. 사용할 테이블의 StorageHandler를 찾기 위해
* 메타데이터 서버를 조회한다. 파티션이 이미 게시되어 있으면 오류를 던진다.
* @param job 잡 객체
* @param outputJobInfo 잡의 테이블 출력 정보
* @throws IOException 메타데이터 서버와의 통신 예외
*/
@SuppressWarnings("unchecked")
public static void setOutput(Job job, OutputJobInfo outputJobInfo) throws IOException;
/**
* 파티션에 쓰여질 데이터의 스키마를 설정한다. 호출하지 않으면 파티션에 테이블 스키마가
* 기본적으로 사용된다.
* @param job 잡 객체
* @param schema 데이터의 스키마
* @throws IOException
*/
public static void setSchema(final Job job, final HCatSchema schema) throws IOException;
/**
* 지정된 잡 컨텍스트에서 HCatOutputFormat.setOutput 호출에 지정된 테이블의 테이블
* 스키마를 얻는다.
* @param context 컨텍스트
* @return 테이블 스키마
* @throws IOException 전달된 컨텍스트에 대해 HCatOutputFormat.setOutput이 호출되지 않은 경우
*/
public static HCatSchema getTableSchema(JobContext context) throws IOException;
HCatRecord
HCatRecord는 HCatalog 테이블에 값을 저장하는 데 지원되는 타입이에요.
HCatalog 테이블 스키마의 타입은 HCatRecord의 서로 다른 필드에 대해 반환되는 객체 타입을 결정해요. 이 테이블은 MapReduce 프로그램용 Java 클래스와 HCatalog 데이터 타입 사이의 매핑을 보여줘요:
| HCatalog 데이터 타입 | MapReduce의 Java 클래스 | 값 |
|---|---|---|
| TINYINT | java.lang.Byte | -128 to 127 |
| SMALLINT | java.lang.Short | -(2^15) to (2^15)-1, 즉 -32,768 to 32,767 |
| INT | java.lang.Integer | -(2^31) to (2^31)-1, 즉 -2,147,483,648 to 2,147,483,647 |
| BIGINT | java.lang.Long | -(2^63) to (2^63)-1, 즉 -9,223,372,036,854,775,808 to 9,223,372,036,854,775,807 |
| BOOLEAN | java.lang.Boolean | true 또는 false |
| FLOAT | java.lang.Float | 단정밀도 부동소수점 값 |
| DOUBLE | java.lang.Double | 배정밀도 부동소수점 값 |
| DECIMAL | java.math.BigDecimal | 38자리 정밀도의 정확한 부동소수점 값 |
| BINARY | byte[] | 이진 데이터 |
| STRING | java.lang.String | 문자열 |
| STRUCT | java.util.List | 구조화된 데이터 |
| ARRAY | java.util.List | 한 데이터 타입의 값들 |
| MAP | java.util.Map | 키-값 쌍 |
Hive 데이터 타입에 대한 일반 정보는 Hive Data Types 및 Type System을 참고해요.
HCatalog로 MapReduce 실행 (Running MapReduce with HCatalog)
MapReduce 프로그램에 Thrift 서버가 어디 있는지 알려줘야 해요. 가장 쉬운 방법은 위치를 Java 프로그램의 인자로 전달하는 것이에요. 또한 -libjars 인자로 Hive와 HCatalog jar를 MapReduce에 전달해야 해요.
export HADOOP_HOME=<path_to_hadoop_install>
export HCAT_HOME=<path_to_hcat_install>
export HIVE_HOME=<path_to_hive_install>
export LIB_JARS=$HCAT_HOME/share/hcatalog/hcatalog-core-0.5.0.jar,
$HIVE_HOME/lib/hive-metastore-0.10.0.jar,
$HIVE_HOME/lib/libthrift-0.7.0.jar,
$HIVE_HOME/lib/hive-exec-0.10.0.jar,
$HIVE_HOME/lib/libfb303-0.7.0.jar,
$HIVE_HOME/lib/jdo2-api-2.3-ec.jar,
$HIVE_HOME/lib/slf4j-api-1.6.1.jar
export HADOOP_CLASSPATH=$HCAT_HOME/share/hcatalog/hcatalog-core-0.5.0.jar:
$HIVE_HOME/lib/hive-metastore-0.10.0.jar:
$HIVE_HOME/lib/libthrift-0.7.0.jar:
$HIVE_HOME/lib/hive-exec-0.10.0.jar:
$HIVE_HOME/lib/libfb303-0.7.0.jar:
$HIVE_HOME/lib/jdo2-api-2.3-ec.jar:
$HIVE_HOME/conf:$HADOOP_HOME/conf:
$HIVE_HOME/lib/slf4j-api-1.6.1.jar
$HADOOP_HOME/bin/hadoop --config $HADOOP_HOME/conf jar <path_to_jar>
<main_class> -libjars $LIB_JARS <program_arguments>
이 방법은 동작하지만 Hadoop은 MapReduce 프로그램을 실행할 때마다 libjars를 배송하며, 파일을 다른 캐시 엔트리로 취급하므로 비효율적이고 Hadoop 분산 캐시를 고갈시킬 수 있어요.
대신 HDFS 위치를 사용해 libjars를 배송하도록 최적화할 수 있어요. 이렇게 하면 Hadoop은 분산 캐시의 엔트리를 재사용해요.
bin/hadoop fs -copyFromLocal $HCAT_HOME/share/hcatalog/hcatalog-core-0.5.0.jar /tmp
bin/hadoop fs -copyFromLocal $HIVE_HOME/lib/hive-metastore-0.10.0.jar /tmp
bin/hadoop fs -copyFromLocal $HIVE_HOME/lib/libthrift-0.7.0.jar /tmp
bin/hadoop fs -copyFromLocal $HIVE_HOME/lib/hive-exec-0.10.0.jar /tmp
bin/hadoop fs -copyFromLocal $HIVE_HOME/lib/libfb303-0.7.0.jar /tmp
bin/hadoop fs -copyFromLocal $HIVE_HOME/lib/jdo2-api-2.3-ec.jar /tmp
bin/hadoop fs -copyFromLocal $HIVE_HOME/lib/slf4j-api-1.6.1.jar /tmp
export LIB_JARS=hdfs:///tmp/hcatalog-core-0.5.0.jar,
hdfs:///tmp/hive-metastore-0.10.0.jar,
hdfs:///tmp/libthrift-0.7.0.jar,
hdfs:///tmp/hive-exec-0.10.0.jar,
hdfs:///tmp/libfb303-0.7.0.jar,
hdfs:///tmp/jdo2-api-2.3-ec.jar,
hdfs:///tmp/slf4j-api-1.6.1.jar
# (나머지 문은 같음)
인증 (Authentication)
"/tmp/kinit <username>@FOO.COM"을 실행했는지 확인해요.
읽기 예제 (Read Example)
다음의 매우 간단한 MapReduce 프로그램은 두 번째 컬럼("column 1")에 정수가 있다고 가정하는 한 테이블에서 데이터를 읽고, 발견한 각 고유 값의 인스턴스 수를 셉니다. 즉 "select col1, count(*) from $table group by col1;"을 수행하는 것과 같아요.
예를 들어 두 번째 컬럼의 값이 {1,1,1,3,3,5}이면 프로그램은 값과 개수의 다음 출력을 생성해요:
1, 3
3, 2
5, 1
public class GroupByAge extends Configured implements Tool {
public static class Map extends
Mapper<WritableComparable, HCatRecord, IntWritable, IntWritable> {
int age;
@Override
protected void map(
WritableComparable key,
HCatRecord value,
org.apache.hadoop.mapreduce.Mapper<WritableComparable, HCatRecord,
IntWritable, IntWritable>.Context context)
throws IOException, InterruptedException {
age = (Integer) value.get(1);
context.write(new IntWritable(age), new IntWritable(1));
}
}
public static class Reduce extends Reducer<IntWritable, IntWritable,
WritableComparable, HCatRecord> {
@Override
protected void reduce(
IntWritable key,
java.lang.Iterable<IntWritable> values,
...
더 알아보기 (Learn more)
- HCatalog 문서에서 HCatalog 전체 기능을 볼 수 있어요.
- HCatalog InputOutput 위키에서 API 예제를 확인할 수 있어요.