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는 아래와 같아요. 포함 사항:

  • setInput
  • setOutputSchema
  • getTableSchema

HCatInputFormat으로 데이터를 읽으려면 먼저 읽는 테이블의 필요한 정보로 InputJobInfo를 인스턴스화한 다음 InputJobInfosetInput을 호출해요.

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는 아래와 같아요. 포함 사항:

  • setOutput
  • setSchema
  • getTableSchema

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 TypesType 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//hive.log"에서 실패로 인해 "2010-11-03 16:17:28,225 WARN hive.metastore ... - Unable to connect metastore with URI thrift://..." 같은 메시지가 발생하면, Kerberos 티켓을 얻고 HCatalog 서버에 인증할 수 있도록 "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)