Hive 카탈로그
Hive 카탈로그 (Hive Catalog)
Hive Metastore는 수년간 Hadoop 생태계에서 사실상의 메타데이터 허브로 진화했어요. 많은 회사가 프로덕션에서 단일 Hive Metastore 서비스 인스턴스를 운영하며 Hive 메타데이터든 비-Hive 메타데이터든 모든 메타데이터를 진실의 원천(source of truth)으로 관리해요. Hive와 Flink 배포를 모두 가진 사용자에게 HiveCatalog는 Hive Metastore로 Flink 메타데이터를 관리할 수 있게 해줘요.
출처: 문서
본문
Flink 배포만 있는 사용자에게 HiveCatalog는 Flink가 기본 제공하는 유일한 영속 카탈로그예요. 영속 카탈로그가 없으면 Flink SQL CREATE DDL을 쓰는 사용자가 Kafka 테이블 같은 메타 객체를 세션마다 반복해서 만들어야 해서 시간이 많이 낭비돼요. HiveCatalog는 테이블과 다른 메타 객체를 한 번만 만들고 이후 세션에 걸쳐 편리하게 참조·관리할 수 있게 해 이 공백을 메워요.
HiveCatalog 설정 (Set up HiveCatalog)
의존성 (Dependencies)
Flink에서 HiveCatalog를 설정하려면 전체 Flink-Hive 통합과 같은 의존성이 필요해요.
구성 (Configuration)
Flink에서 HiveCatalog를 설정하려면 전체 Flink-Hive 통합과 같은 구성이 필요해요.
HiveCatalog 사용 방법 (How to use HiveCatalog)
올바르게 구성되면 HiveCatalog는 기본적으로 동작해요. 사용자는 DDL로 Flink 메타 객체를 만들 수 있고 즉시 볼 수 있어요. HiveCatalog는 두 종류의 테이블을 다룰 수 있어요: Hive 호환 테이블과 일반(generic) 테이블. Hive 호환 테이블은 메타데이터·데이터 모두 저장 계층에서 Hive 호환 방식으로 저장된 것이에요. 따라서 Flink로 만든 Hive 호환 테이블은 Hive 쪽에서 조회할 수 있어요.
반면 일반 테이블은 Flink 고유의 것이에요. HiveCatalog로 일반 테이블을 만들 때는 그저 HMS로 메타데이터를 영속화할 뿐이에요. 이 테이블은 Hive에 보이지만 Hive가 메타데이터를 이해할 가능성은 낮아요. 따라서 Hive에서 이런 테이블을 사용하면 정의되지 않은 동작이 발생해요.
Hive 호환 테이블을 만들려면 Hive dialect로 전환하는 것이 좋아요. 기본 dialect로 Hive 호환 테이블을 만들려면 테이블 프로퍼티에 'connector'='hive'를 설정해야 해요. 그렇지 않으면 HiveCatalog에서 기본적으로 일반 테이블로 간주돼요. Hive dialect를 쓰면 connector 프로퍼티는 필요하지 않다는 점을 주의해요.
예제 (Example)
간단한 예제를 살펴볼게요.
step 1: Hive Metastore 설정
Hive Metastore를 실행해요. 여기서는 로컬 Hive Metastore와 hive-site.xml 파일을 로컬 경로 /opt/hive-conf/hive-site.xml에 설정해요. 다음과 같은 구성이 있어요:
<configuration>
<property>
<name>javax.jdo.option.ConnectionURL</name>
<value>jdbc:mysql://localhost/metastore?createDatabaseIfNotExist=true</value>
<description>metadata is stored in a MySQL server</description>
</property>
<property>
<name>javax.jdo.option.ConnectionDriverName</name>
<value>com.mysql.jdbc.Driver</value>
<description>MySQL JDBC driver class</description>
</property>
<property>
<name>javax.jdo.option.ConnectionUserName</name>
<value>...</value>
<description>user name for connecting to mysql server</description>
</property>
<property>
<name>javax.jdo.option.ConnectionPassword</name>
<value>...</value>
<description>password for connecting to mysql server</description>
</property>
<property>
<name>hive.metastore.uris</name>
<value>thrift://localhost:9083</value>
<description>IP address (or fully-qualified domain name) and port of the metastore host</description>
</property>
<property>
<name>hive.metastore.schema.verification</name>
<value>true</value>
</property>
</configuration>
Hive Cli로 HMS 연결을 테스트해요. 일부 명령을 실행하면 default라는 데이터베이스가 있고 테이블이 없는 것을 볼 수 있어요.
hive> show databases;
OK
default
Time taken: 0.032 seconds, Fetched: 1 row(s)
hive> show tables;
OK
Time taken: 0.028 seconds, Fetched: 0 row(s)
step 2: SQL Client 시작 및 Flink SQL DDL로 Hive 카탈로그 생성
모든 Hive 의존성을 Flink 배포의 /lib 디렉터리에 추가하고, Flink SQL CLI에서 Hive 카탈로그를 다음과 같이 만들어요:
Flink SQL> CREATE CATALOG myhive WITH (
'type' = 'hive',
'hive-conf-dir' = '/opt/hive-conf'
);
step 3: Kafka 클러스터 설정
이름과 나이의 튜플로 "test"라는 토픽을 가진 로컬 Kafka 클러스터를 부트스트랩하고 간단한 데이터를 토픽에 만듭니다:
localhost$ bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test
>tom,15
>john,21
이 메시지는 Kafka 콘솔 consumer를 시작하면 볼 수 있어요.
localhost$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning
tom,15
john,21
step 4: Flink SQL DDL로 Kafka 테이블 생성
Flink SQL DDL로 간단한 Kafka 테이블을 만들고 스키마를 확인해요.
Flink SQL> USE CATALOG myhive;
Flink SQL> CREATE TABLE mykafka (name String, age Int) WITH (
'connector' = 'kafka',
'topic' = 'test',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
);
[INFO] Table has been created.
Flink SQL> DESCRIBE mykafka;
root
|-- name: STRING
|-- age: INT
Hive Cli로 테이블이 Hive에도 보이는지 확인해요:
hive> show tables;
OK
mykafka
Time taken: 0.038 seconds, Fetched: 1 row(s)
step 5: Flink SQL로 Kafka 테이블 조회
Flink 클러스터(standalone 또는 yarn-session)에서 Flink SQL Client로 간단한 select 질의를 실행해요.
Flink SQL> select * from mykafka;
Kafka 토픽에 메시지를 더 만듭니다:
localhost$ bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning
tom,15
john,21
kitty,30
amy,24
kaiky,18
이제 SQL Client에서 Flink가 만든 결과를 다음과 같이 볼 수 있어요:
SQL Query Result (Table)
Refresh: 1 s Page: Last of 1
name age
tom 15
john 21
kitty 30
amy 24
kaiky 18
지원 타입 (Supported Types)
HiveCatalog는 일반 테이블에 모든 Flink 타입을 지원해요. Hive 호환 테이블의 경우 HiveCatalog는 다음 표에 설명된 대로 Flink 데이터 타입을 해당 Hive 타입으로 매핑해야 해요:
| Flink 데이터 타입 | Hive 데이터 타입 |
|---|---|
| CHAR(p) | CHAR(p) |
| VARCHAR(p) | VARCHAR(p) |
| STRING | STRING |
| BOOLEAN | BOOLEAN |
| TINYINT | TINYINT |
| SMALLINT | SMALLINT |
| INT | INT |
| BIGINT | LONG |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| DECIMAL(p, s) | DECIMAL(p, s) |
| DATE | DATE |
| TIMESTAMP(9) | TIMESTAMP |
| BYTES | BINARY |
| ARRAY<T> | LIST<T> |
| MAP | MAP |
| ROW | STRUCT |
타입 매핑에 대해 주의할 점:
- Hive의 CHAR(p)는 최대 길이가 255
- Hive의 VARCHAR(p)는 최대 길이가 65535
- Hive의 MAP은 기본 키 타입만 지원하는 반면 Flink의 MAP은 어떤 데이터 타입이든 될 수 있음
- Hive의 UNION 타입은 지원되지 않음
- Hive의 TIMESTAMP는 항상 정밀도 9이고 다른 정밀도를 지원하지 않음. 반면 Hive UDF는 정밀도가 9 이하인 TIMESTAMP 값을 처리할 수 있음
- Hive는 Flink의 TIMESTAMP_WITH_TIME_ZONE, TIMESTAMP_WITH_LOCAL_TIME_ZONE, MULTISET을 지원하지 않음
- Flink의 INTERVAL 타입은 아직 Hive INTERVAL 타입으로 매핑할 수 없음