본문 바로가기
WIKI 기술 지식 베이스

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 id
  • secret: lakeFS secret access key
  • region: 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): 저장소 ID
  • ref(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 가이드