Hive 커넥터

Hive 커넥터 (Hive connector)

Hive 커넥터는 Apache Hive 데이터 웨어하우스에 저장된 데이터를 질의할 수 있게 해 줘요. Hive는 세 가지 컴포넌트의 조합이에요:

출처: 문서

본문

Hive의 세 컴포넌트는 다음과 같아요:

  • 다양한 포맷의 데이터 파일 — 보통 Hadoop Distributed File System(HDFS)이나 Amazon S3 같은 오브젝트 스토리지 시스템에 저장돼요.
  • 데이터 파일이 스키마와 테이블에 어떻게 매핑되는지에 대한 메타데이터 — MySQL 같은 데이터베이스에 저장되고 Hive 메타스토어 서비스를 통해 접근돼요.
  • HiveQL이라는 쿼리 언어 — MapReduce나 Tez 같은 분산 컴퓨팅 프레임워크에서 실행돼요.

Trino는 처음 두 컴포넌트인 데이터와 메타데이터만 사용해요. HiveQL이나 Hive 실행 환경의 어떤 부분도 사용하지 않아요.

요구사항 (Requirements)

Hive 커넥터는 Hive 메타스토어 서비스(HMS) 또는 AWS Glue 같은 Hive 메타스토어의 호환 구현이 필요해요.

카탈로그 구성 파일에서 지원되는 파일 시스템을 선택해 구성해야 해요.

코디네이터와 모든 워커는 Hive 메타스토어와 스토리지 시스템에 네트워크로 접근할 수 있어야 해요. Thrift 프로토콜의 Hive 메타스토어 접근은 기본적으로 포트 9083을 사용해요.

데이터 파일은 지원되는 파일 포맷이어야 해요. 파일 포맷은 format 테이블 속성과 다른 특정 속성으로 구성할 수 있어요:

직렬화 가능한 포맷의 경우 특정 SerDe만 허용돼요:

  • RCText — ColumnarSerDe를 사용하는 RCFile
  • RCBinary — LazyBinaryColumnarSerDe를 사용하는 RCFile
  • org.apache.hadoop.io.Text를 사용하는 SequenceFile
  • com.twitter.elephantbird.hive.serde.ProtobufDeserializer를 사용해 프로토콜 버퍼 레코드를 담는 org.apache.hadoop.io.BytesWritable을 사용하는 SequenceFile
  • CSV — org.apache.hadoop.hive.serde2.OpenCSVSerde 사용
  • JSON — org.apache.hive.hcatalog.data.JsonSerDe 사용
  • OPENX_JSON — org.openx.data.jsonserde.JsonSerDe의 OpenX JSON SerDe. 소스 저장소에서 Trino 구현에 대한 자세한 내용을 확인하세요.
  • TextFile
  • ESRI — com.esri.hadoop.hive.serde.EsriJsonSerDe 사용
  • ESRI_GEO_JSON — com.esri.hadoop.hive.serde.GeoJsonSerDe 사용

일반 설정 (General configuration)

Hive 커넥터를 구성하려면 hive 커넥터를 참조하는 etc/catalog/example.properties 카탈로그 속성 파일을 만드세요.

메타데이터용 메타스토어를 구성해야 해요.

지원되는 파일 시스템 중 하나를 선택해 구성해야 해요.

connector.name=hive
hive.metastore.uri=thrift://example.net:9083
fs.x.enabled=true

Glue 메타스토어를 쓰면:

connector.name=hive
hive.metastore=glue

테이블 생성 (Table creation)

hive 테이블 생성 예시. format, partitioned_by, bucketed_by, bucket_count 지정:

CREATE TABLE example.web.page_views (
  view_time TIMESTAMP,
  user_id BIGINT,
  page_url VARCHAR,
  ds DATE,
  country VARCHAR
)
WITH (
  format = 'ORC',
  partitioned_by = ARRAY['ds', 'country'],
  bucketed_by = ARRAY['user_id'],
  bucket_count = 50
)

스키마 생성·삭제:

CREATE SCHEMA example.web
WITH (location = 's3://my-bucket/')
DROP SCHEMA example.web

삭제 연산:

DELETE FROM example.web.page_views
WHERE ds = DATE '2016-08-09'
  AND country = 'US'

쿼리:

SELECT * FROM example.web.page_views

파티션 나열:

SELECT * FROM example.web."page_views$partitions"

외부 테이블 생성(external_location으로 데이터 위치 지정):

CREATE TABLE example.web.request_logs (
  request_time TIMESTAMP,
  url VARCHAR,
  ip VARCHAR,
  user_agent VARCHAR
)
WITH (
  format = 'TEXTFILE',
  external_location = 's3://my-bucket/data/logs/'
)

통계 수집과 테이블 삭제:

ANALYZE example.web.request_logs;
DROP TABLE example.web.request_logs

트랜잭션 테이블 생성:

CREATE TABLE
WITH (
    format='ORC',
    transactional=true
)
AS

시스템 프로시저 (System procedures)

빈 파티션 생성:

CALL system.create_empty_partition(
    schema_name => 'web',
    table_name => 'page_views',
    partition_columns => ARRAY['ds', 'country'],
    partition_values => ARRAY['2016-08-09', 'US']);

통계 삭제:

CALL system.drop_stats(
    schema_name => 'web',
    table_name => 'page_views',
    partition_values => ARRAY[ARRAY['2016-08-09', 'US']]);

Avro 테이블엔 avro_schema_url로 스키마를 지정해요:

CREATE TABLE example.avro.avro_data (
   id BIGINT
 )
WITH (
   format = 'AVRO',
   avro_schema_url = '/usr/local/avro_data.avsc'
)

파일 최적화 (Optimize)

테이블 병합(compact) 최적화:

ALTER TABLE test_table EXECUTE optimize

파일 크기 임계값 지정:

ALTER TABLE test_table EXECUTE optimize(file_size_threshold => '128MB')

파티션별 최적화:

ALTER TABLE test_partitioned_table EXECUTE optimize
WHERE partition_key = 1

타임스탬프 기반 조건:

ALTER TABLE test_table EXECUTE optimize
WHERE CAST(timestamp_tz AS DATE) > DATE '2021-12-31'

비트랜잭션 테이블 최적화 활성화:

SET SESSION .non_transactional_optimize_enabled=true

CSV 이스케이프 지정:

CREATE TABLE tablename
WITH (format='CSV',
      csv_escape = '"')

메타데이터·하이든 컬럼

$properties 테이블로 테이블 속성을 확인할 수 있어요:

SELECT * FROM example.web."page_views$properties";
       stats_generated_via_stats_task        | auto.purge |       trino_query_id       | trino_version | transactional
---------------------------------------------+------------+-----------------------------+---------------+---------------
 workaround for potential lack of HIVE-12730 | false      | 20230705_152456_00001_nfugi | 434           | false

$partitions 테이블로 파티션을 확인할 수 있어요:

SELECT * FROM example.web."page_views$partitions";
     day    | country
------------+---------
 2023-07-01 | POL
 2023-07-02 | POL
 2023-07-03 | POL
 2023-03-01 | USA
 2023-03-02 | USA

하이든(hidden) 컬럼으로 파일 경로·크기·파티션을 확인할 수 있어요:

SELECT *, "$path", "$file_size", "$partition"
FROM example.web.page_views;
SELECT *, "$path", "$file_size"
FROM example.web.page_views
WHERE "$partition" = 'ds=2016-08-09/country=US'

특정 파티션만 분석:

ANALYZE table_name WITH (
    partitions = ARRAY[
        ARRAY['p1_value1', 'p1_value2'],
        ARRAY['p2_value1', 'p2_value2']])

파티션과 컬럼 지정:

ANALYZE table_name WITH (
    partitions = ARRAY[ARRAY['p2_value1', 'p2_value2']],
    columns = ARRAY['col_1', 'col_2'])

통계 삭제:

CALL system.drop_stats('schema_name', 'table_name')
CALL system.drop_stats(
    schema_name => 'schema',
    table_name => 'table',
    partition_values => ARRAY[ARRAY['p2_value1', 'p2_value2']])

멀티 카탈로그 액세스

기본 카탈로그와 다른 카탈로그의 테이블을 참조하려면 정규화된 catalog.schema.table 이름을 쓰세요:

USE example.example_schema;

EXPLAIN SELECT * FROM example_table;
                               Query Plan
-------------------------------------------------------------------------
Fragment 0 [SOURCE]
     ...
     Output[columnNames = [...]]
     │   ...
     └─ TableScan[table = another_catalog:example_schema:example_table]
            ...
EXPLAIN SELECT * FROM example.example_schema.example_table;

더 알아보기 (Learn more)

데이터 레이크 포맷 관련 커넥터가 궁금하다면 Iceberg 커넥터 문서를 이어서 읽어 보세요.