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 커넥터 문서를 이어서 읽어 보세요.