Delta Lake 커넥터

Delta Lake 커넥터 (Delta Lake connector)

Delta Lake 커넥터는 Delta Lake 포맷으로 저장된 데이터를 질의할 수 있게 해 줘요. Databricks Delta Lake도 포함해요. 이 커넥터는 Delta Lake 트랜잭션 로그를 네이티브로 읽어 외부 시스템이 데이터를 바꾸는 것을 감지할 수 있어요.

출처: 문서

본문

요구사항 (Requirements)

Databricks Delta Lake에 연결하려면 다음이 필요해요:

  • Databricks Runtime 7.3 LTS, 9.1 LTS, 10.4 LTS, 11.3 LTS, 12.2 LTS, 13.3 LTS, 14.3 LTS, 15.4 LTS, 16.4 LTS, 17.3 LTS로 작성된 테이블을 지원해요.
  • AWS, HDFS, Azure Storage, Google Cloud Storage(GCS) 배포를 완전히 지원해요.
  • 코디네이터와 워커에서 Delta Lake 스토리지로 네트워크 접근.
  • Delta Lake의 Hive 메타스토어 서비스(HMS) 또는 별도 HMS, 또는 Glue 메타스토어 접근.
  • 코디네이터와 워커에서 HMS로 네트워크 접근. HMS가 쓰는 Thrift 프로토콜의 기본 포트는 9083이에요.
  • 지원되는 파일 시스템Parquet 파일 포맷으로 저장된 데이터 파일.

일반 설정 (General configuration)

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

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

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

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

Glue 메타스토어를 쓰면:

connector.name=delta_lake
hive.metastore=glue

시간 여행 (Time travel)

Delta Lake 테이블의 이전 버전을 조회할 수 있어요. 버전 기준:

SELECT *
FROM example.testdb.customer_orders FOR VERSION AS OF 3

타임스탬프 기준:

SELECT *
FROM example.testdb.customer_orders FOR TIMESTAMP AS OF TIMESTAMP '2022-03-23 09:59:29.803 America/Los_Angeles';

과거 스냅샷으로 테이블 생성:

CREATE OR REPLACE TABLE example.testdb.customer_orders AS
SELECT *
FROM example.testdb.customer_orders FOR TIMESTAMP AS OF TIMESTAMP '2022-03-23 09:59:29.803 America/Los_Angeles';

날짜 시간대로도 지정할 수 있어요:

SELECT *
FROM example.testdb.customer_orders FOR TIMESTAMP AS OF DATE '2022-03-23';
SELECT *
FROM example.testdb.customer_orders FOR TIMESTAMP AS OF TIMESTAMP '2022-03-23 00:00:00';
SELECT *
FROM example.testdb.customer_orders FOR TIMESTAMP AS OF TIMESTAMP '2022-03-23 00:00:00.000 America/Los_Angeles';

테이블 히스토리에서 버전과 연산을 확인할 수 있어요:

SELECT version, operation
FROM example.testdb."customer_orders$history"
ORDER BY version DESC

시스템 프로시저 (System procedures)

메타스토어에 없는 기존 Delta Lake 테이블을 Trino에 등록:

CALL example.system.register_table(schema_name => 'testdb', table_name => 'customer_orders', table_location => 's3://my-bucket/a/path')

등록 해제:

CALL example.system.unregister_table(schema_name => 'testdb', table_name => 'customer_orders')

오래된 로그·데이터 파일 정리(vacuum):

CALL example.system.vacuum('exampleschemaname', 'exampletablename', '7d');

스키마·테이블 생성

위치를 지정해 스키마 생성:

CREATE SCHEMA example.example_schema
WITH (location = 's3://my-bucket/a/path');
CREATE SCHEMA example.example_schema;

기존 테이블 등록 후 새 테이블을 만들 수 있어요:

CALL example.system.register_table(schema_name => 'testdb', table_name => 'example_table', table_location => 's3://my-bucket/a/path')
CREATE TABLE example.default.new_table (id BIGINT, address VARCHAR);

CREATE OR REPLACE TABLE로 파티셔닝 지정:

CREATE OR REPLACE TABLE example_table
WITH (partitioned_by = ARRAY['a'])
AS SELECT * FROM another_table;

파일 최적화 (Optimize)

테이블 병합 최적화:

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'

파일 수정 시간 기준:

ALTER TABLE test_table EXECUTE optimize
WHERE "$file_modified_time" > date_trunc('day', CURRENT_TIMESTAMP);

경로 제외:

ALTER TABLE test_table EXECUTE optimize
WHERE "$path" <> 'skipping-file-path'

크기 기준:

-- optimze files smaller than 1MB
ALTER TABLE test_table EXECUTE optimize
WHERE "$file_size" <= 1024 * 1024

메타데이터 테이블

파티셔닝된 테이블 생성:

CREATE TABLE example.default.example_partitioned_table
WITH (
  location = 's3://my-bucket/a/path',
  partitioned_by = ARRAY['regionkey'],
  checkpoint_interval = 5,
  change_data_feed_enabled = false,
  column_mapping_mode = 'name',
  deletion_vectors_enabled = false
)
AS SELECT name, comment, regionkey FROM tpch.tiny.nation;

$history 테이블로 변경 이력을 확인:

SELECT * FROM "test_table$history"
 version |               timestamp               | user_id | user_name |  operation   |         operation_parameters          |                 cluster_id      | read_version |  isolation_level  | is_blind_append | operation_metrics
---------+---------------------------------------+---------+-----------+--------------+---------------------------------------+---------------------------------+--------------+-------------------+-----------------+-------------------
       2 | 2023-01-19 07:40:54.684 Europe/Vienna | trino   | trino     | WRITE        | {queryId=20230119_064054_00008_4vq5t} | trino-406-trino-coordinator     |            2 | WriteSerializable | true            | {}
       1 | 2023-01-19 07:40:41.373 Europe/Vienna | trino   | trino     | ADD COLUMNS  | {queryId=20230119_064041_00007_4vq5t} | trino-406-trino-coordinator     |            0 | WriteSerializable | true            | {}
       0 | 2023-01-19 07:40:10.497 Europe/Vienna | trino   | trino     | CREATE TABLE | {queryId=20230119_064010_00005_4vq5t} | trino-406-trino-coordinator     |            0 | WriteSerializable | true            | {}

$partitions 테이블로 파티션 정보를 확인:

SELECT * FROM "test_table$partitions"
           partition           | file_count | total_size |                     data                     |
-------------------------------+------------+------------+----------------------------------------------+
{_bigint=1, _date=2021-01-12}  |          2 |        884 | {_decimal={min=1.0, max=2.0, null_count=0}}  |
{_bigint=1, _date=2021-01-13}  |          1 |        442 | {_decimal={min=1.0, max=1.0, null_count=0}}  |

$properties 테이블로 테이블 속성을 확인:

SELECT * FROM "test_table$properties"
 key                        | value           |
----------------------------+-----------------+
delta.minReaderVersion      | 1               |
delta.minWriterVersion      | 4               |
delta.columnMapping.mode    | name            |
delta.feature.columnMapping | supported       |

체인지 데이터 피드 (Change Data Feed)

table_changes 테이블 함수로 변경 피드를 조회할 수 있어요:

SELECT
  *
FROM
  TABLE(
    system.table_changes(
      schema_name => 'test_schema',
      table_name => 'tableName',
      since_version => 0
    )
  );

change_data_feed_enabled = true로 테이블을 생성하고:

CREATE TABLE test_schema.pages (page_url VARCHAR, domain VARCHAR, views INTEGER)
    WITH (change_data_feed_enabled = true);

데이터를 삽입하고:

INSERT INTO test_schema.pages
    VALUES
        ('url1', 'domain1', 1),
        ('url2', 'domain2', 2),
        ('url3', 'domain1', 3);
INSERT INTO test_schema.pages
    VALUES
        ('url4', 'domain1', 400),
        ('url5', 'domain2', 500),
        ('url6', 'domain3', 2);

업데이트하면:

UPDATE test_schema.pages
    SET domain = 'domain4'
    WHERE views = 2;

변경 피드를 조회하면 insert/update의 전후 이미지가 보여요:

SELECT
  *
FROM
  TABLE(
    system.table_changes(
      schema_name => 'test_schema',
      table_name => 'pages',
      since_version => 1
    )
  )
ORDER BY _commit_version ASC;
page_url    |     domain     |    views    |    _change_type     |    _commit_version    |    _commit_timestamp
url4        |     domain1    |    400      |    insert           |     2                 |    2023-03-10T21:22:23.000+0000
url5        |     domain2    |    500      |    insert           |     2                 |    2023-03-10T21:22:23.000+0000
url6        |     domain3    |    2        |    insert           |     2                 |    2023-03-10T21:22:23.000+0000
url2        |     domain2    |    2        |    update_preimage  |     3                 |    2023-03-10T22:23:24.000+0000
url2        |     domain4    |    2        |    update_postimage |     3                 |    2023-03-10T22:23:24.000+0000
url6        |     domain3    |    2        |    update_preimage  |     3                 |    2023-03-10T22:23:24.000+0000
url6        |     domain4    |    2        |    update_postimage |     3                 |    2023-03-10T22:23:24.000+0000

since_version을 생략하면 전체 변경 이력을 반환해요:

SELECT
  *
FROM
  TABLE(
    system.table_changes(
      schema_name => 'test_schema',
      table_name => 'pages'
    )
  )
ORDER BY _commit_version ASC;
page_url    |     domain     |    views    |    _change_type     |    _commit_version    |    _commit_timestamp
url1        |     domain1    |    1        |    insert           |     1                 |    2023-03-10T20:21:22.000+0000
url2        |     domain2    |    2        |    insert           |     1                 |    2023-03-10T20:21:22.000+0000
url3        |     domain1    |    3        |    insert           |     1                 |    2023-03-10T20:21:22.000+0000
url4        |     domain1    |    400      |    insert           |     2                 |    2023-03-10T21:22:23.000+0000
url5        |     domain2    |    500      |    insert           |     2                 |    2023-03-10T21:22:23.000+0000
url6        |     domain3    |    2        |    insert           |     2                 |    2023-03-10T21:22:23.000+0000
url2        |     domain2    |    2        |    update_preimage  |     3                 |    2023-03-10T22:23:24.000+0000
url2        |     domain4    |    2        |    update_postimage |     3                 |    2023-03-10T22:23:24.000+0000
url6        |     domain3    |    2        |    update_preimage  |     3                 |    2023-03-10T22:23:24.000+0000
url6        |     domain4    |    2        |    update_postimage |     3                 |    2023-03-10T22:23:24.000+0000

통계 (Statistics)

테이블 분석:

ANALYZE table_schema.table_name;

특정 시점 이후 수정된 파일만 분석:

ANALYZE example_table WITH(files_modified_after = TIMESTAMP '2021-08-23
16:43:01.321 Z')

특정 컬럼만 분석:

ANALYZE example_table WITH(columns = ARRAY['nationkey', 'regionkey'])

확장 통계 삭제:

CALL example.system.drop_extended_stats('example_schema', 'example_table')

트랜잭션 로그 모니터링

JMX를 통해 트랜잭션 로그 접근 캐시 상태를 모니터링할 수 있어요:

SELECT * FROM jmx.current."*.plugin.deltalake.transactionlog:name=,type=transactionlogaccess"
datafilemetadatacachestats.hitrate      | 0.97
datafilemetadatacachestats.missrate     | 0.03
datafilemetadatacachestats.requestcount | 3232
metadatacachestats.hitrate              | 0.98
metadatacachestats.missrate             | 0.02
metadatacachestats.requestcount         | 6783
node                                    | trino-master
object_name                             | io.trino.plugin.deltalake.transactionlog:type=TransactionLogAccess,name=delta

멀티 카탈로그 액세스

기본 카탈로그와 다른 카탈로그의 테이블을 참조하려면 정규화된 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 커넥터 문서를 이어서 읽어 보세요.