Lua 훅
Lua 훅 (Lua Hooks)
lakeFS는 내장 Lua VM을 사용해 외부 구성 요소 없이 훅을 실행할 수 있어요.
Lua 훅을 쓰면 액션이 발생했을 때 lakeFS 서버가 직접 실행할 Lua 스크립트를 넘길 수 있어요.
lakeFS에 내장된 Lua 런타임은 보안상 제한돼 있어요. 기본으로 다음을 허용하지 않는 좁은 API와 함수 집합만 제공해요:
-
실행 중인 lakeFS 서버 환경에 접근하기
-
lakeFS 프로세스가 볼 수 있는 로컬 파일시스템에 접근하기
출처: 문서
본문
액션 파일의 Lua 훅 속성
정보
전체 설정 스키마와 세부 사항은 Action configuration을 참고하세요.
| 속성 | 설명 | 데이터 타입 | 필수 | 기본값 |
|---|---|---|---|---|
| args | 훅에 전달할 하나 이상의 인자 | Dictionary | false | |
| script | 인라인 Lua 스크립트 | String | 이 값 또는 script_path 중 하나는 반드시 지정 | |
| script_path | lakeFS 안의 Lua 스크립트 경로 | String | 이 값 또는 script 중 하나는 반드시 지정 |
Lua 훅 예시
더 많은 예시와 설정 샘플은 lakeFS 저장소의 examples/hooks/ 디렉터리를 참고하세요. 훅이 실제로 돌아가는 단계별 예시는 lakeFS samples 저장소에도 있어요.
이벤트 정보 표시하기
이 예시는 발생한 이벤트를 JSON 표현으로 출력해요:
name: dump_all
on:
post-commit:
post-merge:
post-create-tag:
post-create-branch:
hooks:
- id: dump_event
type: lua
properties:
script: |
json = require("encoding/json")
print(json.marshal(action))
커밋에 필수 메타데이터 필드 강제하기
더 실용적인 예시: 모든 커밋이 필수 메타데이터 필드를 담고 있는지 확인해요:
name: pre commit metadata field check
on:
pre-commit:
branches:
- main
- dev
hooks:
- id: ensure_commit_metadata
type: lua
properties:
args:
notebook_url: {"pattern": "my-jupyter.example.com/.*"}
spark_version: {}
script_path: lua_hooks/ensure_metadata_field.lua
lakefs://repo/main/lua_hooks/ensure_metadata_field.lua의 Lua 코드:
regexp = require("regexp")
for k, props in pairs(args) do
current_value = action.commit.metadata[k]
if current_value == nil then
error("missing mandatory metadata field: " .. k)
end
if props.pattern and not regexp.match(props.pattern, current_value) then
error("current value for commit metadata field " .. k .. " does not match pattern: " .. props.pattern .. " - got: " .. current_value)
end
end
더 많은 예시와 설정 샘플은 lakeFS 저장소의 examples/hooks/ 디렉터리를 참고하세요.
Lua 라이브러리 레퍼런스
lakeFS에 내장된 Lua 런타임은 보안상 제한돼 있어요. 제공되는 API는 아래와 같아요. 네트워크에 도달하는 패키지는 lakeFS 서버에서 요청을 보내며, actions.network.blocked_addresses 설정에 구성된 주소 제한을 받아요.
array(table)
테이블 객체를 런타임이 배열로 인식하도록 _is_array: true 메타테이블 필드를 설정하는 헬퍼 함수예요.
aws
aws/s3_client
S3 클라이언트 라이브러리예요.
예시
local aws = require("aws")
-- pass valid AWS credentials
local client = aws.s3_client("ACCESS_KEY_ID", "SECRET_ACCESS_KEY", "REGION")
aws/s3_client.get_object(bucket, key)
요청한 객체의 본문(Lua 문자열)과, 객체가 존재하면 true인 불리언 값을 반환해요.
aws/s3_client.put_object(bucket, key, value)
지정한 버킷과 키의 객체를 전달한 value 문자열 값으로 설정해요.
aws/s3_client.delete_object(bucket [, key])
지정한 키의 객체를 삭제해요.
aws/s3_client.list_objects(bucket [, prefix, continuation_token, delimiter])
다음 구조를 담은 결과 테이블을 반환해요:
-
is_truncated: (boolean) continuation 토큰으로 더 페이징할 결과가 있는지 여부 -
next_continuation_token: (string) 다음 페이지 결과를 얻는 다음 요청에 넘길 토큰 -
results(테이블의 테이블): 객체(그리고 delimiter를 썼다면 접두사)에 대한 정보
하나의 결과는 다음 구조 중 하나예요
{
["key"] = "a/common/prefix/",
["type"] = "prefix"
}
또는:
{
["key"] = "path/to/object",
["type"] = "object",
["etag"] = "etagString",
["size"] = 1024,
["last_modified"] = "2023-12-31T23:10:00Z"
}
aws/s3_client.delete_recursive(bucket, prefix)
지정한 접두사 아래의 모든 객체를 삭제해요.
aws/glue
Glue 클라이언트 라이브러리예요.
예시
local aws = require("aws")
-- pass valid AWS credentials
local glue = aws.glue_client("ACCESS_KEY_ID", "SECRET_ACCESS_KEY", "REGION")
aws/glue.create_database(database, options)
Glue Catalog에 새 Database를 만들어요.
파라미터:
-
database(string): Glue Database 이름. -
options(table)(선택):-
error_on_already_exists(boolean): 같은 이름의 DB가 이미 있을 때 오류로 실패할지 여부 -
create_db_input(Table): AWS에 "그대로" 전달되는 테이블로, AWS SDK의 CreateDatabaseInput에 대응해요
-
예시
local opts = {
error_on_already_exists = false,
create_db_input = {DatabaseInput = {Description = "Created via LakeFS Action"}, Tags = {Owner = "Joe"}}
}
glue.create_database(db, opts)
aws/glue.delete_database(database, catalog_id)
Glue Catalog의 기존 Database를 삭제해요.
파라미터:
-
database(string): Glue Database 이름. -
catalog_id(string)(선택): Glue Catalog ID
예시
glue.delete_database(db, "461129977393")
aws/glue.get_table(database, table [, catalog_id)
Glue Catalog의 테이블을 조회(descrbe)해요.
예시
local table, exists = glue.get_table(db, table_name)
if exists then
print(json.marshal(table))
aws/glue.create_table(database, table_input, [, catalog_id])
Glue Catalog에 새 테이블을 만들어요. table_input 인자는 AWS에 "그대로" 전달되는 JSON으로, AWS SDK의 TableInput에 대응해요.
예시
local json = require("encoding/json")
local input = {
Name = "my-table",
PartitionKeys = array(partitions),
-- etc...
}
local json_input = json.marshal(input)
glue.create_table("my-db", table_input)
aws/glue.update_table(database, table_input, [, catalog_id, version_id, skip_archive])
Glue Catalog의 기존 테이블을 갱신해요. table_input은 glue.create_table 함수의 인자와 같아요.
aws/glue.delete_table(database, table_input, [, catalog_id])
Glue Catalog의 기존 테이블을 삭제해요.
azure
azure/blob_client
Azure blob 클라이언트 라이브러리예요.
예시
local azure = require("azure")
-- pass valid Azure credentials
local client = azure.blob_client("AZURE_STORAGE_ACCOUNT", "AZURE_ACCESS_KEY")
azure/blob_client.get_object(path_uri)
요청한 객체의 본문(Lua 문자열)과, 객체가 존재하면 true인 불리언 값을 반환해요.
path_uri - https://myaccount.blob.core.windows.net/mycontainer/myblob 형태의 유효한 Azure blob storage uri예요.
azure/blob_client.put_object(path_uri, value)
지정한 버킷과 키의 객체를 전달한 value 문자열 값으로 설정해요.
path_uri - https://myaccount.blob.core.windows.net/mycontainer/myblob 형태의 유효한 Azure blob storage uri예요.
azure/blob_client.delete_object(path_uri)
지정한 키의 객체를 삭제해요.
path_uri - https://myaccount.blob.core.windows.net/mycontainer/myblob 형태의 유효한 Azure blob storage uri예요.
azure/abfss_transform_path(path)
HTTPS Azure URL을 ABFSS 스킴으로 변환해요. delta_exporter 함수가 Azure Unity catalog 사용 사례를 지원할 때 써요. path - https://myaccount.blob.core.windows.net/mycontainer/myblob 형태의 유효한 Azure blob storage URL이에요.
crypto
crypto/aes/encryptCBC(key, plaintext)
AES로 암호화한 텍스트의 암호문을 반환해요.
crypto/aes/decryptCBC(key, ciphertext)
암호화된 암호문을 복호화한(평문) 문자열을 반환해요.
crypto/hmac/sign_sha256(message, key)
전달한 키와 메시지로 SHA256 hmac 시그니처를 반환해요(SHA256 해싱 알고리즘 사용).
crypto/hmac/sign_sha1(message, key)
전달한 키와 메시지로 SHA1 hmac 시그니처를 반환해요(SHA1 해싱 알고리즘 사용).
crypto/md5/digest(data)
주어진 데이터의 MD5 다이제스트(문자열)를 반환해요.
crypto/sha256/digest(data)
주어진 데이터의 SHA256 다이제스트(문자열)를 반환해요.
databricks/client(databricks_host, databricks_service_principal_token)
register_external_table과 create_or_get_schema 메서드를 가진 Databricks 클라이언트를 나타내는 테이블을 반환해요.
databricks/client.create_schema(schema_name, catalog_name, get_if_exists)
설정된 Databricks 호스트의 Unity catalog에서 스키마를 만들거나, 이미 있으면 가져와요. 스키마가 없다면 지정한 catalog_name 아래에 주어진 schema_name으로 새 스키마가 생겨요. 생성되거나 가져온 스키마 이름을 반환해요.
파라미터:
-
schema_name(string): 원하는 스키마 이름 -
catalog_name(string): 스키마가 생성될(또는 가져와질) 카탈로그 이름 -
get_if_exists(boolean): 지정한catalog_name에 같은schema_name의 스키마가 이미 있어 실패한 경우, 그 스키마를 반환해요.
예시
local databricks = require("databricks")
local client = databricks.client("https://my-host.cloud.databricks.com", "my-service-principal-token")
local schema_name = client.create_schema("main", "mycatalog", true)
databricks/client.execute_statement(statement, warehouse_id, catalog_name, schema_name)
파라미터:
-
statement(boolean): databricks 테이블에서 실행할 SQL 문 -
warehouse_id(string): Databricks에서CREATE TABLE쿼리를 실행하는 데 쓰이는 SQL warehouse ID(SQL warehouse에서 가져옴 -
catalog_name(string): 스키마가 생성될(또는 가져와질) 카탈로그 이름 -
schema_name(string): 원하는 스키마 이름status, SQL 상태(SUCCEEDED 또는 오류 코드/메시지)를 반환해요
예시
local databricks = require("databricks")
local client = databricks.client("https://my-host.cloud.databricks.com", "my-service-principal-token")
local statement = "ALTER TABLE " .. table_descriptor.name .. " ALTER COLUMN ID SET MASK mask_num"
databricks_client.execute_statement(statement, args.warehouse_id, table_descriptor.catalog, table_descriptor.schema)
databricks/client.register_external_table(table_name, physical_path, warehouse_id, catalog_name, schema_name, metadata)
지정한 warehouse ID, catalog 이름, schema 이름 아래에 외부 테이블을 등록해요. 이 메서드가 성공하려면 카탈로그에 외부 위치(external location)가 설정되어 있고, physical_path의 루트 스토리지 URI(예: s3://mybucket)를 가리켜야 해요. 테이블의 생성 상태를 반환해요.
파라미터:
-
table_name(string): 테이블 이름. -
physical_path(string): 외부 테이블이 참조할 위치, 예:s3://mybucket/the/path/to/mytable. -
warehouse_id(string): Databricks에서CREATE TABLE쿼리를 실행하는 데 쓰이는 SQL warehouse ID(SQL warehouse의Connection Details에서, 또는databricks warehouses get을 실행해 SQL warehouse를 고르고 ID를 가져와요). -
catalog_name(string): 스키마가 생성될(또는 가져와질) 카탈로그 이름. -
schema_name(string): 테이블이 생성될 스키마 이름. -
metadata(table): 테이블 등록에 추가할 메타데이터 테이블이에요. 형태는{key1 = "value1", key2 = "value2", ...}이어야 해요.
예시
local databricks = require("databricks")
local client = databricks.client("https://my-host.cloud.databricks.com", "my-service-principal-token")
local status = client.register_external_table("mytable", "s3://mybucket/the/path/to/mytable", "examwarehouseple", "my-catalog-name", "myschema")
- 이 메서드를 실행하는 데 필요한 Databricks 권한은 Unity Catalog Exporter 문서를 참고하세요.
encoding/base64/encode(data)
주어진 데이터를 base64 문자열로 인코딩해요.
encoding/base64/decode(data)
base64로 인코딩된 데이터를 디코딩해 문자열로 반환해요.
encoding/base64/url_encode(data)
주어진 데이터를 RFC 4648에 정의된 패딩 없는 대체 base64 인코딩으로 인코딩해요.
encoding/base64/url_decode(data)
RFC 4648에 정의된 패딩 없는 대체 base64 인코딩을 디코딩해 문자열로 반환해요.
encoding/hex/encode(value)
주어진 value 문자열을 16진수 값(문자열)으로 인코딩해요.
encoding/hex/decode(value)
16진수 문자열을 원래 표현하던 문자열(UTF-8)로 되돌려요.
encoding/json/marshal(table)
주어진 테이블을 JSON 문자열로 인코딩해요.
encoding/json/unmarshal(string)
주어진 문자열을 대응되는 Lua 구조로 디코딩해요.
encoding/yaml/marshal(table)
주어진 테이블을 YAML 문자열로 인코딩해요.
encoding/yaml/unmarshal(string)
YAML로 인코딩된 문자열을 대응되는 Lua 구조로 디코딩해요.
encoding/parquet/get_schema(payload)
payload(문자열)를 Parquet 파일 내용으로 읽어 다음 테이블 구조의 스키마를 반환해요:
{
{ ["name"] = "column_a", ["type"] = "INT32" },
{ ["name"] = "column_b", ["type"] = "BYTE_ARRAY" }
}
formats
formats/delta_client(key, secret, region)
lakeFS 서버와 상호작용하는 새 Delta Lake 클라이언트를 만들어요.
key: lakeFS access key idsecret: lakeFS secret access keyregion: lakeFS 서버가 설정된 리전이에요.
formats/delta_client.get_table(repository_id, reference_id, prefix)
지정한 저장소, 참조, 접두사 아래의 Delta Lake 테이블 표현을 반환해요. 응답은 두 개의 테이블이에요:
- 첫 번째는
{number, {string}}형태의 테이블이에요.number는 Delta Log의 버전이고, 매핑된{string}배열에는 해당 버전 항목에 나열된 여러 Delta Lake 로그 연산의 JSON 문자열이 들어가요. 예:
{
0 = {
"{\"commitInfo\":...}",
"{\"add\": ...}",
"{\"remove\": ...}"
},
1 = {
"{\"commitInfo\":...}",
"{\"add\": ...}",
"{\"remove\": ...}"
}
}
- 두 번째는 현재 테이블 스냅샷의 메타데이터 테이블이에요. 이 메타데이터 테이블은 외부 카탈로그에서 Delta Lake 테이블을 초기화할 때 쓸 수 있어요.
다음 필드로 이루어져요:
-
id: 테이블의 ID -
name: 테이블의 이름 -
description: 테이블의 설명 -
schema_string: 테이블의 스키마 문자열 -
partition_columns: 테이블의 파티션 컬럼 -
configuration: 테이블의 구성 -
created_time: 테이블의 생성 시각
gcloud
gcloud/gs_client(gcs_credentials_json_string)
유효한 credentials.json 파일 내용을 담은 문자열로 새 Google Cloud Storage 클라이언트를 만들어요.
gcloud/gs.write_fuse_symlink(source, destination, mount_info)
소스(보통 객체의 lakeFS 물리적 주소)에서 지정한 대상으로 gcsfuse symlink를 만들어요.
mount_info는 "from"과 "to" 키를 가진 Lua 테이블이에요 — symlink는 gs://... URI로는 동작하지 않아서 마운트된 위치를 가리켜야 하거든요. from이 source의 앞부분에서 제거되고, 그 자리에 destination이 들어가요.
예시
source = "gs://bucket/lakefs/data/abc/def"
destination = "gs://bucket/exported/path/to/object"
mount_info = {
["from"] = "gs://bucket",
["to"] = "/home/user/gcs-mount"
}
gs.write_fuse_symlink(source, destination, mount_info)
-- Symlink: "/home/user/gcs-mount/exported/path/to/object" -> "/home/user/gcs-mount/lakefs/data/abc/def"
hook
사용자 친화적인 훅을 작성하는 데 도움을 주는 유틸리티 집합이에요.
hook/fail(message)
전달한 메시지로 현재 훅의 실행을 중단시켜요. error()를 쓰는 것과 비슷하지만, 일반적인 런타임 오류(예상 밖의 응답을 반환한 API 호출)와 호출한 훅의 명시적 실패를 구분할 때 보통 써요.
호출되면 오류가 스택 트레이스 없이 표시되고, 오류 메시지가 message로 준 값 그대로 나와요.
예시
> hook = require("hook")
> hook.fail("this hook shall not pass because of: " .. reason)
lakefs
Lua Hook 라이브러리는 액션을 트리거한 사용자의 신원으로 lakeFS API를 다시 호출할 수 있게 해 줘요. 예컨대 사용자 A가 커밋을 시도해 pre-commit 훅을 트리거했다면, 그 훅 안에서 lakeFS API로 호출하는 모든 요청은 인증과 감사(auditing) 목적으로 자동으로 사용자 A의 신원을 써요.
lakefs/create_tag(repository_id, reference_id, tag_id)
지정한 참조에 새 태그를 만들어요.
lakefs/diff_refs(repository_id, left_reference_id, right_reference_id [, after, prefix, delimiter, amount])
left_reference_id와 right_reference_id 사이의 객체 단위 diff를 반환해요.
lakefs/list_objects(repository_id, reference_id [, after, prefix, delimiter, amount])
지정한 저장소와 참조(브랜치, 태그, 커밋 ID 등)에서 객체를 나열해요. delimiter가 비어 있으면 재귀 목록이 기본이고, 그렇지 않으면 delimiter까지의 공통 접두사가 하나의 항목으로 표시돼요.
lakefs/get_object(repository_id, reference_id, path)
2개의 값을 반환해요:
-
lakeFS API가 반환한 HTTP 상태 코드
-
지정한 객체의 내용(Lua 문자열)
lakefs/diff_branch(repository_id, branch_id [, after, amount, prefix, delimiter])
branch_id의 커밋되지 않은 변경에 대한 객체 단위 diff를 반환해요.
lakefs/stat_object(repository_id, ref_id, path[, user_metadata])
지정한 참조와 저장소에서 주어진 경로의 stat 객체를 반환해요. 2개의 값을 반환해요:
-
lakeFS API가 반환한 HTTP 상태 코드
-
JSON 문자열인 stat 응답 객체
파라미터:
-
repository_id: 저장소 ID -
ref_id: stat을 수행할 참조(브랜치, 태그, 커밋 ID) -
path: stat할 객체 경로 -
user_metadata: (선택) 응답에 사용자 메타데이터를 포함할지 나타내는 불리언 플래그
lakefs/update_object_user_metadata(repository_id, branch_id, path, metadata)
객체의 사용자 메타데이터를 갱신해요.
파라미터:
-
repository_id: 저장소 ID -
branch_id: 객체가 들어 있는 브랜치 -
path: 갱신할 객체 경로 -
metadata: 사용자 메타데이터로 설정할 키-값 쌍을 담은 테이블
2개의 값을 반환해요:
-
lakeFS API가 반환한 HTTP 상태 코드 (성공 시 204)
-
성공 시 비어 있고, 실패 시 오류
lakefs/catalogexport/glue_exporter.get_full_table_name(descriptor, action_info)
glue 테이블 이름을 생성해요.
파라미터:
-
descriptor(Table): (예: _lakefs_tables/my_table.yaml)의 객체예요. -
action_info(Table): 전역 action 객체예요.
lakefs/catalogexport/delta_exporter
참고
이 모듈과 그 별칭
lakefs/catalogexport/delta_exporter_v1은 폐기(deprecated)됐고 2026년 9월 1일에 동작을 멈춰요. 대신lakefs/catalogexport/delta_exporter_v2를 사용하세요.
lakeFS에서 외부 클라우드 스토리지로 Delta Lake 테이블을 내보내는 데 쓰는 패키지예요.
lakefs/catalogexport/delta_exporter.export_delta_log(action, table_def_names, write_object, delta_client, table_descriptors_path, path_transformer)
Delta Lake 테이블을 내보내는 함수예요. 반환값은 테이블 이름을 외부 테이블 위치(여기서 데이터를 쿼리할 수 있어요)와 최신 Delta 테이블 버전의 메타데이터에 매핑한 테이블이에요.
응답은 다음 형태예요:
{<table_name> = {path = "s3://mybucket/mypath/mytable", metadata = {id = "table_id", name = "table_name", ...}}}.
파라미터:
-
action: 전역 action 객체 -
table_def_names: Delta 테이블 이름 목록 (예:{"table1", "table2"}) -
write_object:function(bucket, key, data)시그니처를 가진 writer 함수로, 내보낸 Delta Log를 기록하는 데 쓰여요 (예:aws/s3_client.put_object또는azure/blob_client.put_object) -
delta_client:get_table: function(repo, ref, prefix)을 구현하는 Delta Lake 클라이언트 -
table_descriptors_path: 지정한table_def_names의 테이블 기술자가 있는 경로 -
path_transformer: (선택) 저장된 delta 로그의 path 필드와 저장된 테이블의 물리적 경로를 변환하는 function(path)예요 (Azure Unity catalog 사용 사례를 지원할 때 써요)
AWS S3용 Delta export 예시
---
name: delta_exporter
on:
post-commit: null
hooks:
- id: delta_export
type: lua
properties:
script: |
local aws = require("aws")
local formats = require("formats")
local delta_exporter = require("lakefs/catalogexport/delta_exporter")
local json = require("encoding/json")
local table_descriptors_path = "_lakefs_tables"
local sc = aws.s3_client(args.aws.access_key_id, args.aws.secret_access_key, args.aws.region)
local delta_client = formats.delta_client(args.lakefs.access_key_id, args.lakefs.secret_access_key, args.aws.region)
local delta_table_details = delta_exporter.export_delta_log(action, args.table_defs, sc.put_object, delta_client, table_descriptors_path)
for t, details in pairs(delta_table_details) do
print("Delta Lake exported table \"" .. t .. "\"'s location: " .. details["path"] .. "\n")
print("Delta Lake exported table \"" .. t .. "\"'s metadata:\n")
for k, v in pairs(details["metadata"]) do
if type(v) == "table" then
print("\t" .. k .. " = " .. json.marshal(v) .. "\n")
else
print("\t" .. k .. " = " .. v .. "\n")
end
end
end
args:
aws:
access_key_id: <AWS_ACCESS_KEY_ID>
secret_access_key: <AWS_SECRET_ACCESS_KEY>
region: us-east-1
lakefs:
access_key_id: <LAKEFS_ACCESS_KEY_ID>
secret_access_key: <LAKEFS_SECRET_ACCESS_KEY>
table_defs:
- mytable
_lakefs_tables/mytable.yaml의 테이블 기술자:
---
name: myTableActualName
type: delta
path: a/path/to/my/delta/table
Azure Blob Storage용 Delta export 예시:
name: Delta Exporter
on:
post-commit:
branches: ["{% raw %}{{ .Branch }}{% endraw %}*"]
hooks:
- id: delta_exporter
type: lua
properties:
script: |
local azure = require("azure")
local formats = require("formats")
local delta_exporter = require("lakefs/catalogexport/delta_exporter")
local table_descriptors_path = "_lakefs_tables"
local sc = azure.blob_client(args.azure.storage_account, args.azure.access_key)
local function write_object(_, key, buf)
return sc.put_object(key,buf)
end
local delta_client = formats.delta_client(args.lakefs.access_key_id, args.lakefs.secret_access_key)
local delta_table_details = delta_exporter.export_delta_log(action, args.table_defs, sc.put_object, delta_client, table_descriptors_path)
for t, details in pairs(delta_table_details) do
print("Delta Lake exported table \"" .. t .. "\"'s location: " .. details["path"] .. "\n")
print("Delta Lake exported table \"" .. t .. "\"'s metadata:\n")
for k, v in pairs(details["metadata"]) do
if type(v) == "table" then
print("\t" .. k .. " = " .. json.marshal(v) .. "\n")
else
print("\t" .. k .. " = " .. v .. "\n")
end
end
end
args:
azure:
storage_account: "{% raw %}{{ .AzureStorageAccount }}{% endraw %}"
access_key: "{% raw %}{{ .AzureAccessKey }}{% endraw %}"
lakefs: # provide credentials of a user that has access to the script and Delta Table
access_key_id: "{% raw %}{{ .LakeFSAccessKeyID }}{% endraw %}"
secret_access_key: "{% raw %}{{ .LakeFSSecretAccessKey }}{% endraw %}"
table_defs:
- mytable
lakefs/catalogexport/delta_exporter.changed_table_defs(table_def_names, table_descriptors_path, repository_id, ref, compare_ref)
변경된 테이블을 기준으로 테이블 defs 목록을 걸러 주는 유틸리티 함수예요. table_def_names 파라미터 중 변경된 테이블들의 부분 집합을 반환해요.
파라미터:
-
table_def_names(table of strings): diff를 기준으로 필터링할 테이블 이름 목록 -
table_descriptors_path(string): 지정한table_def_names의 테이블 기술자가 있는 경로 -
repository_id(string): 저장소 ID -
ref(string): 데이터의 특정 버전을 가리키는 기준 참조, 즉 브랜치, 커밋 ID, 태그 -
compare_ref(string): 어떤 테이블이 바뀌었는지 판단할 diff 비교 대상 참조
예시
local delta_export = require("lakefs/catalogexport/delta_exporter")
local ref = action.commit.parents[1]
local compare_ref = action.commit_id
local changed_table_defs = delta_export.changed_table_defs(args.table_defs, args.table_descriptors_path, action.repository_id, ref, compare_ref)
for i = 1, #changed_table_defs do
print(changed_table_defs[i])
end
lakefs/catalogexport/delta_exporter_v2
lakeFS Enterprise v1.92.0부터 사용 가능. 무료 평가판 시작.
delta_exporter 패키지의 후속 버전이에요. lakeFS에 내장된 Delta 클라이언트로 lakeFS에서 외부 클라우드 스토리지로 Delta Lake 테이블을 내보내요. 클라이언트가 내장됐기 때문에 훅이 더 이상 formats.delta_client로 클라이언트를 만들지 않고, args에 lakeFS 자격 증명도 필요하지 않아요. 이 익스포터는 Delta Lake 테이블 기능(table features)도 인식해요: 지원하지 않는 기능을 선언한 테이블은 건너뛰고, 인식하지 못하는 기능에는 경고를 줘요. 자세한 내용은 Unity Catalog 연동 문서에 있고, 이 모듈을 쓰는 완전한 훅 예시도 거기 있어요.
lakefs/catalogexport/delta_exporter_v2.export_delta_log(action, table_def_names, export_storage, delta_client, table_descriptors_path, path_transformer)
Delta Lake 테이블을 내보내는 함수예요. v1과 시그니처가 같아서 기존 훅은 require하는 모듈만 바꾸면 돼요. 반환값도 같은 매핑, 즉 테이블 기술자 이름을 내보낸 테이블 위치와 최신 Delta 테이블 버전의 메타데이터에 매핑한 {<table_def_name> = {path = "s3://mybucket/mypath/mytable", metadata = {id = "table_id", name = "table_name", ...}}} 형태예요.
파라미터:
-
action: 전역 action 객체 -
table_def_names:table_descriptors_path아래에 있는 테이블 기술자 파일 이름 목록..yaml확장자가 있든 없든 받아들여요 (예:{"table1", "table2"}) -
export_storage: 내보낸 Delta Log가 기록될 곳이에요. Choosing how the export is written에서 설명한 것처럼 스토리지 기술자 테이블이거나 writer 함수예요 -
delta_client: 무시되며nil이어도 돼요. 익스포터가 lakeFS에 내장된 Delta 클라이언트를 쓰기 때문이에요. 이 파라미터는 v1과의 시그니처 호환성을 위해 남아 있고, 클라이언트를 넘기면 경고를 출력해요. -
table_descriptors_path: 지정한table_def_names의 테이블 기술자가 있는 경로 -
path_transformer: (선택) 저장된 delta 로그의 path 필드와 저장된 테이블의 물리적 경로를 변환하는 function(path)예요 (Azure Unity catalog 사용 사례를 지원할 때 써요)
익스포트 기록 방식 고르기
선호되는 export_storage는 대상 스토리지와 그곳에 쓸 자격 증명을 지정하는 스토리지 기술자 테이블이에요. 이러면 훅에서 구현할 게 아무것도 남지 않아요:
{ type = "s3", access_key_id = "...", secret_access_key = "...", region = "...", endpoint = "..." }
{ type = "azure", storage_account = "...", access_key = "..." }
{ type = "gs", credentials_json = "..." }
하위 호환을 위해 export_storage는 function(bucket, key, data) 시그니처의 writer 함수도 받아들여요. aws/s3_client.put_object나 azure/blob_client.put_object처럼요. 익스포터는 기록하는 모든 파일에 대해 이 함수를 호출해요.
lakefs/catalogexport/delta_exporter_v2.changed_table_defs(table_def_names, table_descriptors_path, repository_id, ref, compare_ref)
테이블 defs 목록을 두 참조 사이에서 데이터가 바뀐 것들로 걸러 주는 유틸리티 함수예요. 파라미터와 동작은 v1 대응 함수와 같아요.
lakefs/catalogexport/table_extractor
_lakefs_tables/ 기술자를 파싱하는 유틸리티 패키지예요.
lakefs/catalogexport/table_extractor.list_table_descriptor_entries(client, repo_id, commit_id)
_lakefs_tables/* 아래의 모든 YAML 파일을 나열하고 [{physical_address, path}] 타입의 목록을 반환해요. 숨김 파일은 무시해요. client는 lakefs 클라이언트예요.
lakefs/catalogexport/table_extractor.get_table_descriptor(client, repo_id, ref, logical_path)
테이블 기술자를 읽고 YAML 객체로 파싱해요. 파티션이 정의되어 있지 않으면 partition_columns를 {}로 설정해요.
파라미터:
client:lakefs클라이언트repo_id(string): 저장소 IDref(string): 데이터의 특정 버전을 가리키는 참조, 즉 브랜치, 커밋 ID, 태그logical_path(string): 저장소 안에서 테이블 기술자 파일의 논리적 경로
lakefs/catalogexport/hive.extract_partition_pager(client, repo_id, commit_id, base_path, partition_cols, page_size)
Hive 포맷 파티션 이터레이터예요. 각 결과 집합은 lakeFS에서 같은 파티션에 속하는 파일들의 모음이에요.
예시
local lakefs = require("lakefs")
local pager = hive.extract_partition_pager(lakefs, repo_id, commit_id, prefix, partitions, 10)
for part_key, entries in pager do
print("partition: " .. part_key)
for _, entry in ipairs(entries) do
print("path: " .. entry.path .. " physical: " .. entry.physical_address)
end
end
lakefs/catalogexport/symlink_exporter
Hive의 SymlinkTextInputFormat을 사용해 테이블 메타데이터를 기록해요. 현재 S3만 지원돼요.
커밋별 기본 익스포트 경로:
${storageNamespace}
_lakefs/
exported/
${ref}/
${commitId}/
${tableName}/
p1=v1/symlink.txt
p1=v2/symlink.txt
p1=v3/symlink.txt
...
lakefs/catalogexport/symlink_exporter.export_s3(s3_client, table_src_path, action_info [, options])
테이블을 표현하는 Symlink 파일을 S3 위치로 내보내요.
파라미터:
-
s3_client: 설정된 클라이언트. -
table_src_path(string):_lakefs_tables안의 테이블 spec YAML 파일 경로 (예: _lakefs_tables/my_table.yaml). -
action_info(table): 전역 action 객체. -
options(table):-
debug(boolean): 추가 정보를 출력해요. -
export_base_uri(string): S3의 접두사를 덮어써요, 예:s3://other-bucket/path/. -
writer(function(bucket, key, data)): 이 값을 주면 s3 클라이언트를 쓰지 않아요. 디버깅에 유용해요.
-
예시
local exporter = require("lakefs/catalogexport/symlink_exporter")
local aws = require("aws")
-- args are user inputs from a lakeFS action.
local s3 = aws.s3_client(args.aws.aws_access_key_id, args.aws.aws_secret_access_key, args.aws.aws_region)
exporter.export_s3(s3, args.table_descriptor_path, action, {debug=true})
lakefs/catalogexport/glue_exporter
lakeFS에 저장된 테이블을 Glue catalog로 내보내는 과정을 자동화하는 패키지예요.
lakefs/catalogexport/glue_exporter.export_glue(glue, db, table_src_path, create_table_input, action_info, options)
lakeFS 테이블을 Glue Catalog에서 표현해요. 이 함수는 설정에 따라 Glue에 테이블을 만들어요. symlink 위치가 이미 만들어져 있다고 가정하고, 기본적으로 같은 커밋에 맞춰 설정만 합니다.
파라미터:
-
glue: AWS glue 클라이언트 -
db(string): glue 데이터베이스 이름 -
table_src_path(string): 테이블 spec 경로 (예: _lakefs_tables/my_table.yaml) -
create_table_input(table): AWS의 table_input에 대응하는 입력 맵핑이에요.glue.create_table에서 쓰는 것과 같아요. 데이터 포맷을 설명하는 입력(InputFormat, OutputFormat, SerdeInfo 등)이 들어 있어야 해요. 익스포터는 이 부분에 개입하지 않거든요. 기본적으로 이 함수는 테이블 위치와 스키마를 설정해요. -
action_info(table): 전역 action 객체. -
options(table):-
table_name(string): 기본 glue 테이블 이름을 덮어써요 -
debug(boolean -
export_base_uri(string): symlink 위치의 S3 기본 접두사를 덮어써요, 예: s3://other-bucket/path/ -
create_db_input(table): 이 값을 지정하면 테이블 익스포트를 위해 새 데이터베이스를 만들고 싶다는 뜻이에요. 이 파라미터는 JSON으로 변환되어 AWS에 "그대로" 전달되는 테이블을 기대하며, AWS SDK의 CreateDatabaseInput에 대응해요
-
glue 테이블을 만들 때 최종 테이블 입력은 create_table_input 입력 파라미터와, 그것을 덮어쓰는 lakeFS 계산 기본값으로 구성돼요:
-
Name—get_full_table_name(descriptor, action_info)로 생성한 glue 테이블 이름. -
PartitionKeys— 보통_lakefs_tables/${table_src_path}에서 유추한 파티션 컬럼. -
TableType= "EXTERNAL_TABLE" -
StorageDescriptor— 보통_lakefs_tables/${table_src_path}에서 유추한 컬럼. -
StorageDescriptor.Location= symlink_location
예시
local aws = require("aws")
local exporter = require("lakefs/catalogexport/glue_exporter")
local glue = aws.glue_client(args.aws_access_key_id, args.aws_secret_access_key, args.aws_region)
-- table_input can be passed as a simple Key-Value object in YAML as an argument from an action, this is inline example:
local table_input = {
StorageDescriptor:
InputFormat: "org.apache.hadoop.hive.ql.io.SymlinkTextInputFormat"
OutputFormat: "org.apache.hadoop.hive.ql.io.IgnoreKeyTextOutputFormat"
SerdeInfo:
SerializationLibrary: "org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe"
Parameters:
classification: "parquet"
EXTERNAL: "TRUE"
"parquet.compression": "SNAPPY"
}
exporter.export_glue(glue, "my-db", "_lakefs_tables/animals.yaml", table_input, action, {debug=true, create_db_input = {DatabaseInput = {Description="DB exported from LakeFS"}, Tags = {Owner = "Joe"}}})
lakefs/catalogexport/glue_exporter.get_full_table_name(descriptor, action_info)
glue 테이블 이름을 생성해요.
파라미터:
-
descriptor(Table): (예: _lakefs_tables/my_table.yaml)의 객체예요. -
action_info(Table): 전역 action 객체예요.
lakefs/catalogexport/unity_exporter
내보낸 Delta Lake 테이블을 Databricks의 Unity catalog에 등록하는 데 쓰는 패키지예요.
lakefs/catalogexport/unity_exporter.register_tables(action, table_descriptors_path, delta_table_details, databricks_client, warehouse_id)
내보낸 Delta Lake 테이블을 Databricks의 Unity Catalog에 등록하는 함수예요. 등록은 다음 경로를 써요: <catalog>.<branch name>.<table_name> — 브랜치 이름이 스키마 이름으로 쓰여요. 반환값은 테이블 이름을 등록 요청 상태에 매핑한 테이블이에요.
참고: (Azure 사용자) Databricks catalog external locations는 ADLS Gen2 스토리지 계정만 지원해요.
lakefs/catalogexport/delta_exporter.export_delta_log 함수로 Delta 테이블을 내보낼 때는 path_transformer를 사용해 경로 스킴을 abfss로 변환해야 해요. 내장 azure Lua 라이브러리가 transformPathToAbfss로 이 기능을 제공해요.
파라미터:
-
action(table): 전역 action 테이블 -
table_descriptors_path(string): 지정한table_paths의 테이블 기술자가 있는 경로. -
delta_table_details(table): 테이블 이름을 물리적 경로와 테이블 메타데이터에 매핑 (예:{table1 = {path = "s3://mybucket/mytable1", metadata = {id = "table_1_id", name = "table1", ...}}, table2 = {path = "s3://mybucket/mytable2", metadata = {id = "table_2_id", name = "table2", ...}}}.) -
databricks_client(table):create_or_get_schema: function(id, catalog_name)과register_external_table: function(table_name, physical_path, warehouse_id, catalog_name, schema_name)을 구현하는 Databricks 클라이언트 -
warehouse_id(string): Databricks warehouse ID.
예시
다음은 내보낸 Delta Lake 테이블을 Unity Catalog에 등록해요.
local databricks = require("databricks")
local unity_export = require("lakefs/catalogexport/unity_exporter")
local delta_table_locations = {
["table1"] = "s3://mybucket/mytable1",
}
-- Register the exported table in Unity Catalog:
local action_details = {
repository_id = "my-repo",
commit_id = "commit_id",
branch_id = "main",
}
local databricks_client = databricks.client("<DATABRICKS_HOST>", "<DATABRICKS_TOKEN>")
local registration_statuses = unity_export.register_tables(action_details, "_lakefs_tables", delta_table_locations, databricks_client, "<WAREHOUSE_ID>")
for t, status in pairs(registration_statuses) do
print("Unity catalog registration for table \"" .. t .. "\" completed with status: " .. status .. "\n")
end
_lakefs_tables/delta-table-descriptor.yaml의 테이블 기술자:
---
name: my_table_name
type: delta
path: path/to/delta/table/data
catalog: my-catalog
lakeFS 액션의 일부로 unity_exporter.register_tables를 쓰는 자세한 단계별 가이드는 Unity Catalog 문서를 참고하세요.
path/parse(path_string)
주어진 경로 문자열에 대해 다음 구조의 테이블을 반환해요.
예시
> require("path")
> path.parse("a/b/c.csv")
{
["parent"] = "a/b/"
["base_name"] = "c.csv"
}
path/join(*path_parts)
가변 개수의 문자열을 받아 경로를 표현하는 결합된 문자열을 반환해요:
예시
> path = require("path")
> path.join("/", "path/", "to", "a", "file.data")
path/o/a/file.data
path/is_hidden(path_string [, seperator, prefix])
불리언을 반환해요 — 주어진 경로 문자열이 숨김(prefix로 시작)이거나 그 부모 중 하나가 prefix로 시작하면 true예요.
예시
> require("path")
> path.is_hidden("a/b/c") -- false
> path.is_hidden("a/b/_c") -- true
> path.is_hidden("a/_b/c") -- true
> path.is_hidden("a/b/_c/") -- true
path/default_separator()
상수 문자열(/)을 반환해요.
regexp/match(pattern, s)
문자열 s가 pattern에 일치하면 true를 반환해요. Go의 regexp.MatchString을 얇게 감싼 래퍼예요.
regexp/quote_meta(s)
문자열 s의 메타 문자를 이스케이프하고 새 문자열을 반환해요.
regexp/compile(pattern)
주어진 패턴의 regexp 매치 객체를 반환해요.
regexp/compiled_pattern.find_all(s, n)
패턴에 매칭되는 모든 결과의 테이블 목록을 반환해요 (최대 n개, n == -1이면 가능한 모든 매치를 반환).
regexp/compiled_pattern.find_all_submatch(s, n)
패턴의 모든 서브매치의 테이블 목록을 반환해요 (최대 n개, n == -1이면 가능한 모든 매치를 반환). 서브매치는 정규식 안의 괄호로 묶인 하위 표현식(캡처 그룹)의 매치로, 여는 괄호 순서대로 왼쪽부터 번호가 매겨져요. 서브매치 0은 표현식 전체의 매치, 서브매치 1은 첫 괄호 하위 표현식의 매치예요.
regexp/compiled_pattern.find(s)
문자열 s에서 주어진 패턴의 가장 왼쪽 매치를 나타내는 문자열을 반환해요.
regexp/compiled_pattern.find_submatch(s)
find_submatch는 s에서 정규식의 가장 왼쪽 매치 텍스트와, 있다면 그 서브매치들의 매치를 담은 문자열 테이블을 반환해요.
strings/split(s, sep)
문자열 테이블을 반환해요. s를 sep으로 나눈 결과예요.
strings/trim(s)
유니코드 기준으로 앞뒤의 모든 공백을 제거한 문자열을 반환해요.
strings/replace(s, old, new, n)
문자열 s에서 겹치지 않는 old를 new로 바꾼 사본을 반환해요. 앞에서부터 n개만 바꿔요. old가 빈 문자열이면 문자열 시작과 각 UTF-8 시퀀스 뒤에 매치되어, k-rune 문자열에서 최대 k+1개의 치환이 일어나요.
n < 0이면 치환 횟수에 제한이 없어요.
strings/has_prefix(s, prefix)
s가 prefix로 시작하면 true를 반환해요.
strings/has_suffix(s, suffix)
s가 suffix로 끝나면 true를 반환해요.
strings/contains(s, substr)
substr이 s의 어디든 포함되어 있으면 true를 반환해요.
time/now()
유닉스 에포크(01/01/1970 00:00:00) 이후 경과 나노초를 나타내는 float64를 반환해요.
time/format(epoch_nano, layout, zone)
주어진 Timezone(예: "UTC", "America/Los_Angeles", ...)에서의 epoch_nano 타임스탬프 문자열 표현을 반환해요. layout 파라미터는 Go의 시간 레이아웃 형식을 따라야 해요.
time/format_iso(epoch_nano, zone)
주어진 Timezone(예: "UTC", "America/Los_Angeles", ...)에서의 epoch_nano 타임스탬프 문자열 표현을 반환해요. 반환 문자열은 ISO8601 형식이에요.
time/sleep(duration_ns)
duration_ns 나노초 동안 잠자요(sleep).
time/since(epoch_nano)
epoch_nano 이후 경과한 나노초 양을 반환해요.
time/add(epoch_time, duration_table)
주어진 duration에 대한 새 타임스탬프(01/01/1970 00:00:00 이후 경과 나노초)를 반환해요. duration은 다음 구조의 테이블이어야 해요.
예시
> require("time")
> time.add(time.now(), {
["hour"] = 1,
["minute"] = 20,
["second"] = 50
})
테이블에서 어떤 필드든 생략할 수 있고, 생략된 필드는 기본값 0이 돼요.
time/parse(layout, value)
유닉스 에포크(01/01/1970 00:00:00) 이후 경과 나노초를 나타내는 float64를 반환해요. 이 타임스탬프는 layout 형식으로 파싱한 날짜 value를 나타내요.
layout 파라미터는 Go의 시간 레이아웃 형식을 따라야 해요.
time/parse_iso(value)
유닉스 에포크(01/01/1970 00:00:00) 이후 경과 나노초를 나타내는 float64를 value에 대해 반환해요. value 문자열은 ISO8601 형식이어야 해요.
uuid/new()
문자열 표현의 새 128비트 RFC 4122 UUID를 반환해요.
net/url
URL 문자열을 부분으로 나누는 parse 함수를 제공해요. URL의 host, path, scheme, query, fragment를 담은 테이블을 반환해요.
예시
> local url = require("net/url")
> url.parse("https://example.com/path?p1=a#section")
{
["host"] = "example.com"
["path"] = "/path"
["scheme"] = "https"
["query"] = "p1=a"
["fragment"] = "section"
}
net/http (optional)
HTTP 요청을 수행하는 request 함수를 제공해요. 보안상의 이유로 이 패키지는 기본적으로 제공되지 않아요. lakeFS 인스턴스 네트워크에서 나가는 http 요청을 가능하게 하기 때문이에요. 이 기능은 actions.lua.net_http_enabled 설정 아래에서 활성화해야 해요. 요청은 30초 후 타임아웃되고, actions.network.blocked_addresses 아래에 설정된 주소 제한을 받아요.
예시
http.request(url [, body])
http.request{
url = string,
[method = string,]
[headers = header-table,]
[body = string,]
}
code(숫자), body(문자열), headers(테이블), status(문자열)를 반환해요.
-
code - 상태 코드 숫자
-
body - 응답 본문을 담은 문자열
-
headers - 응답 요청 헤더의 테이블 (키/값 또는 값의 테이블)
-
status - 상태 코드 텍스트
첫 번째 형태는 GET 요청을 수행하고, body 파라미터가 넘어가면 POST 요청을 해요.
두 번째 형태는 테이블을 받아 요청 메서드와 헤더를 커스터마이즈할 수 있게 해 줘요.
GET 요청 예시
예시
local http = require("net/http")
local code, body = http.request("<https://example.com>")
if code == 200 then
print(body)
else
print("Failed to get example.com - status code: " .. code)
end
POST 요청 예시
local http = require("net/http")
local code, body = http.request{
url="https://httpbin.org/post",
method="POST",
body="custname=tester",
headers={["Content-Type"]="application/x-www-form-urlencoded"},
}
if code == 200 then
print(body)
else
print("Failed to post data - status code: " .. code)
end
더 알아보기 (Learn more)
공식 문서: lakeFS Lua Hooks 가이드