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

객체 및 데이터 연산 다루기

원문 보기 위키 갱신

이 가이드는 lakeFS의 객체 연산을 다뤄요. 업로드, 다운로드, 배치(batch) 연산, 메타데이터 관리까지 포함합니다.

출처: 객체 및 데이터 연산 다루기

본문

기본 객체 연산

객체 업로드

lakeFS에 데이터를 업로드해요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# Upload text data
branch.object("data/simple.txt").upload(
    data=b"Hello, lakeFS!"
)

# Upload with content type
branch.object("data/data.json").upload(
    data=b'{"key": "value"}',
    content_type="application/json"
)

# Upload larger data
csv_data = b"id,name,value\n1,Alice,100\n2,Bob,200\n3,Carol,300"
branch.object("data/records.csv").upload(data=csv_data)

print("Objects uploaded successfully")

객체 다운로드

lakeFS에서 객체 데이터를 읽어와요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# Read as text
with branch.object("data/simple.txt").reader(mode='r') as f:
    content = f.read()
    print(f"Content: {content}")

# Read as binary
with branch.object("data/data.json").reader(mode='rb') as f:
    binary_content = f.read()
    print(f"Binary size: {len(binary_content)} bytes")

# Read CSV and process
import csv
import io

with branch.object("data/records.csv").reader(mode='r') as f:
    reader = csv.DictReader(f)
    for row in reader:
        print(f"  {row['name']}: {row['value']}")

객체 정보 & 메타데이터

객체 세부 정보와 메타데이터를 가져와요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

obj = branch.object("data/records.csv")

# Check if object exists
try:
    if obj.exists():
        print("Object exists")
except:
    print("Object not found")

# Get object statistics
stat = obj.stat()
print(f"Size: {stat.size_bytes} bytes")
print(f"Modified: {stat.mtime}")
print(f"Checksum: {stat.checksum}")
print(f"Content Type: {stat.content_type}")
print(f"Path: {stat.path}")

객체 삭제

lakeFS에서 객체를 제거해요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# Delete a single object
obj = branch.object("data/temp_file.txt")
obj.delete()
print("Object deleted")

# Handle non-existent objects gracefully
try:
    obj.delete()
except Exception as e:
    print(f"Delete failed: {e}")

배치 연산

여러 객체 배치 삭제

많은 객체를 효율적으로 삭제해요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# Delete multiple objects by path
paths_to_delete = [
    "data/file1.csv",
    "data/file2.csv",
    "data/file3.csv",
    "logs/temp.log"
]

try:
    branch.delete_objects(paths_to_delete)
    print(f"Deleted {len(paths_to_delete)} objects")
except Exception as e:
    print(f"Batch delete failed: {e}")

객체 목록 조회와 필터링

프리픽스로 객체 목록 조회

경로 아래의 모든 객체를 나열해요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# List all objects in data/ folder
print("Objects in data/:")
for obj in branch.objects(prefix="data/"):
    print(f"  {obj.path} ({obj.size_bytes} bytes)")

# Count total objects
total_objects = 0
for _ in branch.objects(prefix="data/"):
    total_objects += 1
print(f"Total objects: {total_objects}")

구분자(Delimiter)로 목록 조회 (폴더 뷰)

delimiter를 사용하면 폴더 구조를 볼 수 있어요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# List with folder delimiter
print("Folder structure (with /):")
for item in branch.objects(prefix="", delimiter="/"):
    if hasattr(item, 'path'):
        # It's a file
        print(f"  FILE: {item.path}")
    else:
        # It's a folder
        print(f"  FOLDER: {item.name}")

객체 메타데이터 다루기

커스텀 객체 메타데이터 설정

객체에 커스텀 메타데이터를 붙여요:

import lakefs
import json

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# Create object with metadata
obj = branch.object("data/important.csv")
obj.upload(
    data=b"id,value\n1,100",
    metadata={
        "owner": "data-team",
        "sensitivity": "public",
        "version": "1.0"
    }
)

print("Object uploaded with metadata")

객체 메타데이터 읽기

객체 메타데이터를 조회해요:

import lakefs

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

obj = branch.object("data/important.csv")
stat = obj.stat()

print(f"Object: {stat.path}")
print(f"Size: {stat.size_bytes}")
print(f"Metadata: {stat.metadata}")

실전 워크플로

데이터 정리

오래된 파일과 임시 파일을 제거해요:

import lakefs
from datetime import datetime, timedelta

def cleanup_old_files(repo_name, branch_name, days_old=7):
    """Delete files older than specified days"""
    repo = lakefs.repository(repo_name)
    branch = repo.branch(branch_name)

    cutoff_time = datetime.now().timestamp() - (days_old * 24 * 60 * 60)

    old_files = []

    for obj in branch.objects():
        if hasattr(obj, 'mtime') and obj.mtime < cutoff_time:
            old_files.append(obj.path)

    if old_files:
        print(f"Found {len(old_files)} files older than {days_old} days")
        branch.delete_objects(old_files)
        print(f"Deleted {len(old_files)} old files")
        return len(old_files)
    else:
        print("No old files to delete")
        return 0

# Usage:
deleted_count = cleanup_old_files("archive-repo", "main", days_old=30)
print(f"Cleanup complete: {deleted_count} files removed")

대량 데이터 임포트

여러 파일을 효율적으로 임포트해요:

import lakefs
import os

def bulk_import_files(repo_name, branch_name, local_dir, lakeFS_prefix):
    """Import all files from local directory"""
    repo = lakefs.repository(repo_name)
    branch = repo.branch(branch_name)

    imported = 0
    errors = 0

    # Walk local directory
    for root, dirs, files in os.walk(local_dir):
        for filename in files:
            local_path = os.path.join(root, filename)

            # Calculate lakeFS path
            rel_path = os.path.relpath(local_path, local_dir)
            lakeFS_path = f"{lakeFS_prefix}/{rel_path}".replace("\\", "/")

            try:
                # Read and upload file
                with open(local_path, 'rb') as f:
                    data = f.read()

                branch.object(lakeFS_path).upload(data=data)
                print(f"  Imported: {lakeFS_path}")
                imported += 1

            except Exception as e:
                print(f"  Error importing {lakeFS_path}: {e}")
                errors += 1

    return imported, errors

# Usage (pseudo-code - adjust for your environment):
# imported, errors = bulk_import_files(
#     "my-repo",
#     "main",
#     "/local/data/directory",
#     "data/imports"
# )
# print(f"Imported: {imported}, Errors: {errors}")

스트림 처리

대용량 파일을 효율적으로 처리해요:

import lakefs
import io

def process_csv_stream(repo_name, branch_name, file_path, processor_func):
    """Process large CSV file line by line"""
    repo = lakefs.repository(repo_name)
    branch = repo.branch(branch_name)

    processed = 0

    with branch.object(file_path).reader(mode='r') as f:
        for line in f:
            processor_func(line.strip())
            processed += 1

    return processed

# Usage:
def count_records(line):
    pass  # Do something with each line

count = process_csv_stream(
    "data-repo",
    "main",
    "data/large_file.csv",
    count_records
)

완전한 데이터 파이프라인 만들기

데이터 연산, 트랜잭션, 머지를 결합한 엔드투엔드 파이프라인을 구현해요:

import lakefs

# Get repository and create experiment branch
repo = lakefs.repository("analytics-repo")
branch = repo.branch("processing-v2").create(source_reference="main")

try:
    # Upload raw data
    branch.object("raw/input.csv").upload(data=raw_data)

    # Perform transformations with transactions
    with branch.transact(commit_message="Process raw data") as tx:
        # Read and transform
        with tx.object("raw/input.csv").reader() as f:
            processed = transform(f.read())

        # Write processed data
        tx.object("processed/output.csv").upload(data=processed)

    # Review changes before merging
    changes = list(branch.uncommitted())
    print(f"Changes: {len(changes)} objects")

    # Merge to main if satisfied
    branch.merge_into(repo.branch("main"))

except Exception as e:
    print(f"Error in pipeline: {e}")
    branch.delete()  # Clean up on failure

이 패턴은 다음을 보장해요:

  • 원본 데이터가 격리된 상태로 보존돼요

  • 변환 작업이 원자적으로(all-or-nothing) 수행돼요

  • 통합하기 전에 변경 사항을 검토할 수 있어요

  • 실패한 파이프라인도 안전하게 정리할 수 있어요

에러 처리

객체 에러 다루기

import lakefs
from lakefs.exceptions import NotFoundException, ObjectNotFoundException

repo = lakefs.repository("my-data-repo")
branch = repo.branch("main")

# Object not found
try:
    obj = branch.object("non-existent.csv")
    obj.delete()
except (NotFoundException, ObjectNotFoundException):
    print("Object not found")

# Permission denied
try:
    obj = branch.object("data/file.csv")
    obj.upload(data=b"data")
except Exception as e:
    print(f"Upload failed: {e}")

더 알아보기 (Learn more)

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