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

lakefs-spec으로 파일 시스템 연산 다루기

원문 보기 위키 갱신

lakefs-spec 프로젝트는 fsspec 기반으로 lakeFS에 파일시스템 같은 API를 제공해요. 데이터 사이언스 워크플로, pandas 통합, S3 비슷한 연산을 원하는 시나리오에 딱 맞는 통합이에요.

Note

lakefs-spec은 lakeFS 커뮤니티가 유지보수하는 서드파티 패키지예요. 이슈와 질문은 lakefs-spec repository를 참고하세요.

출처: lakefs-spec으로 파일 시스템 연산 다루기

본문

언제 사용하나요?

다음 경우에 lakefs-spec을 사용해요:

  • 파일시스템 같은 연산(open, read, write, delete)이 필요할 때

  • 데이터 사이언스 도구(pandas, dask, polars)와 함께 작업할 때

  • 브랜치를 명시적으로 관리하지 않는 S3 호환 인터페이스가 필요할 때

  • fsspec 호환 라이브러리와 통합해야 할 때

  • 버저닝 추상화보다 익숙한 파일 연산을 선호할 때

버저닝 중심 워크플로(브랜치, 태그, 커밋)에는 대신 High-Level SDK를 사용하세요.

설치

lakefs-spec은 pip으로 설치해요:

pip install lakefs-spec

또는 최신 버전으로 업그레이드:

pip install --upgrade lakefs-spec

기본 설정

파일시스템 초기화

from lakefs_spec import LakeFSFileSystem

# Auto-discover credentials from ~/.lakectl.yaml
fs = LakeFSFileSystem()

# Or provide explicit credentials
fs = LakeFSFileSystem(
    host="http://localhost:8000",
    username="your-access-key",
    password="your-secret-key"
)

파일 연산

파일 쓰기

from pathlib import Path
from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# Write text file
fs.pipe("my-repo/main/data/text.txt", b"Hello, lakeFS!")

# Write from local file
local_file = Path("local_data.csv")
local_file.write_text("id,name\n1,Alice\n2,Bob")
fs.put(str(local_file), "my-repo/main/data/imported.csv")

# Write using context manager
with fs.open("my-repo/main/data/output.txt", "w") as f:
    f.write("Data written to lakeFS")

파일 읽기

from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# Read entire file
data = fs.cat("my-repo/main/data/text.txt")
print(data.decode())

# Read using context manager
with fs.open("my-repo/main/data/input.txt", "r") as f:
    content = f.read()
    print(content)

# Read in chunks (for large files)
with fs.open("my-repo/main/data/large_file.csv", "rb") as f:
    chunk = f.read(1024)
    while chunk:
        process(chunk)
        chunk = f.read(1024)

파일 목록 조회

from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# List files at path
files = fs.ls("my-repo/main/data/")
for file in files:
    print(file)

# Find files with glob pattern
csv_files = fs.glob("my-repo/main/**/*.csv")
for csv in csv_files:
    print(csv)

파일 삭제

from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# Delete single file
fs.rm("my-repo/main/data/temp.txt")

# Delete directory recursively
fs.rm("my-repo/main/data/temp_dir", recursive=True)

pandas 통합

pandas로 데이터 읽기

import pandas as pd
from lakefs_spec import LakeFSFileSystem

# Read CSV directly from lakeFS
df = pd.read_csv("lakefs://my-repo/main/data/dataset.csv")
print(df.head())

# Read Parquet
df = pd.read_parquet("lakefs://my-repo/main/data/data.parquet")

# Read JSON
df = pd.read_json("lakefs://my-repo/main/data/data.json")

pandas에서 데이터 쓰기

import pandas as pd

# Create sample data
df = pd.DataFrame({
    "id": [1, 2, 3, 4, 5],
    "name": ["Alice", "Bob", "Carol", "David", "Eve"],
    "value": [100, 200, 300, 400, 500]
})

# Write to lakeFS as CSV
df.to_csv("lakefs://my-repo/main/output/data.csv", index=False)

# Write as Parquet
df.to_parquet("lakefs://my-repo/main/output/data.parquet")

# Write as JSON
df.to_json("lakefs://my-repo/main/output/data.json")

데이터 사이언스 워크플로 예제

import pandas as pd
import numpy as np

# Read training data
train_df = pd.read_csv("lakefs://ml-repo/main/datasets/train.csv")
test_df = pd.read_csv("lakefs://ml-repo/main/datasets/test.csv")

# Process data
train_df["normalized_value"] = (train_df["value"] - train_df["value"].mean()) / train_df["value"].std()
test_df["normalized_value"] = (test_df["value"] - test_df["value"].mean()) / test_df["value"].std()

# Save processed data
train_df.to_parquet("lakefs://ml-repo/main/processed/train.parquet")
test_df.to_parquet("lakefs://ml-repo/main/processed/test.parquet")

print("Processing complete!")

트랜잭션

트랜잭션으로 원자적 연산 수행

from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# Perform atomic operations
with fs.transaction("my-repo", "main") as tx:
    # All operations happen on ephemeral branch
    fs.pipe(
        f"my-repo/{tx.branch.id}/data/file1.txt",
        b"Content 1"
    )
    fs.pipe(
        f"my-repo/{tx.branch.id}/data/file2.txt",
        b"Content 2"
    )

    # Commit when done
    tx.commit(message="Add files atomically")
    print(f"Committed: {tx.branch.id}")

에러 처리가 있는 트랜잭션

from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

try:
    with fs.transaction("my-repo", "main") as tx:
        # Perform operations
        fs.pipe(f"my-repo/{tx.branch.id}/file.txt", b"data")

        # Validate
        stat = fs.stat(f"my-repo/{tx.branch.id}/file.txt")
        if stat["size"] < 100:
            raise ValueError("File too small")

        tx.commit(message="Validated and committed")

except Exception as e:
    print(f"Transaction failed: {e}")
    print("Changes rolled back automatically")

트랜잭션 후 태깅

from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

with fs.transaction("ml-repo", "main") as tx:
    # Train and save model
    fs.pipe(f"ml-repo/{tx.branch.id}/models/model.pkl", model_data)

    # Save metrics
    fs.pipe(f"ml-repo/{tx.branch.id}/metrics.json", metrics_data)

    # Commit
    tx.commit(message="Model v1.0")

    # Tag as release
    tx.tag("v1.0.0")
    print("Model released as v1.0.0")

실전 예제

ETL 파이프라인

import pandas as pd
from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# Extract: Read from multiple sources
raw_files = fs.glob("my-repo/main/raw/*.csv")
dfs = [pd.read_csv(f"lakefs://{f}") for f in raw_files]
combined = pd.concat(dfs)

# Transform: Clean and process
combined = combined.dropna()
combined["timestamp"] = pd.to_datetime(combined["timestamp"])
combined["normalized"] = (combined["value"] - combined["value"].mean()) / combined["value"].std()

# Load: Write processed data
combined.to_parquet("lakefs://my-repo/main/processed/data.parquet")
print(f"ETL complete: {len(combined)} rows processed")

데이터 분석과 리포팅

import pandas as pd
import json
from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

# Read data for analysis
df = pd.read_parquet("lakefs://analytics-repo/main/data/raw_data.parquet")

# Perform analysis
summary = {
    "total_records": len(df),
    "mean_value": float(df["value"].mean()),
    "median_value": float(df["value"].median()),
    "std_value": float(df["value"].std())
}

# Save report
report = json.dumps(summary, indent=2)
fs.pipe("lakefs://analytics-repo/main/reports/summary.json", report.encode())

print("Analysis report saved")

모델 버저닝

import pickle
from datetime import datetime
from lakefs_spec import LakeFSFileSystem

fs = LakeFSFileSystem()

def save_model_version(repo, model, version, metrics):
    """Save model with version and metrics"""
    timestamp = datetime.now().isoformat()

    with fs.transaction(repo, "main") as tx:
        branch_id = tx.branch.id

        # Save model
        model_bytes = pickle.dumps(model)
        fs.pipe(
            f"{repo}/{branch_id}/models/{version}/model.pkl",
            model_bytes
        )

        # Save metrics
        metrics_json = json.dumps({
            "version": version,
            "timestamp": timestamp,
            **metrics
        })
        fs.pipe(
            f"{repo}/{branch_id}/models/{version}/metrics.json",
            metrics_json.encode()
        )

        # Commit
        tx.commit(message=f"Model {version}")

        # Tag for reference
        tx.tag(f"model-{version}")
        print(f"Model {version} saved and tagged")

# Usage:
model = train_model(training_data)
save_model_version(
    "ml-repo",
    model,
    "v2.1.0",
    {"accuracy": 0.95, "f1": 0.94}
)

더 볼 자료

  • lakefs-spec Project - 공식 프로젝트 문서

  • fsspec Documentation - 파일시스템 스펙 레퍼런스

더 알아보기 (Learn more)

공식 문서의 원문은 https://docs.lakefs.io/reference/python/lakefs-spec/ 에서 확인할 수 있어요.