객체 및 데이터 연산 다루기
이 가이드는 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/ 에서 확인할 수 있어요.