lakeFS S3 Gateway에서 Boto3 사용하기
lakeFS는 S3 Gateway를 통해 S3 호환 API를 제공해요. 그래서 Boto3(AWS SDK for Python)를 lakeFS와 직접 사용할 수 있어요. 기존 S3 워크플로와 애플리케이션에 그대로 연결할 수 있다는 점이 이 통합의 큰 장점이에요.
출처: 문서
본문
lakeFS는 S3 Gateway를 통해 S3 호환 API를 노출해서, Boto3(AWS SDK for Python)를 lakeFS에서 직접 쓸 수 있게 해줘요. 기존 S3 워크플로와 애플리케이션에는 딱 맞는 통합이에요.
Info
lakeFS와 S3를 함께 Boto로 다루려면 Boto S3 Router를 확인해 보세요. 제공된 버킷 이름에 따라 요청을 S3나 lakeFS로 분기해 줘요.
언제 사용하나요
다음 경우에 Boto를 lakeFS와 함께 사용해요:
-
lakeFS와 함께 쓰고 싶은 기존 S3 워크플로가 있을 때
-
S3 호환 작업(put, get, list, delete)이 필요할 때
-
레거시 S3 애플리케이션을 다룰 때
-
코드 변경 없이 S3에서 마이그레이션하고 싶을 때
버전 관리 중심의 워크플로라면 High-Level SDK나 lakefs-spec을 사용하세요.
설치
pip로 Boto3를 설치해요:
pip install boto3
또는 최신 버전으로 업그레이드해요:
pip install --upgrade boto3
기본 설정
Boto3 클라이언트 초기화
import boto3
# Create S3 client pointing to lakeFS
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key',
region_name='us-east-1'
)
print("Client initialized")
체크섬 설정
최신 버전의 Boto3에서 HTTPS를 사용할 때 AccessDenied 오류가 나고 lakeFS 로그에 encoding/hex: invalid byte: U+0053 'S'가 표시되는 경우가 있어요. 체크섬 설정이 원인이에요.
체크섬 설정 구성하기
import boto3
from botocore.config import Config
# Configure checksum settings
config = Config(
request_checksum_calculation='when_required',
response_checksum_validation='when_required'
)
s3 = boto3.client(
's3',
endpoint_url='https://lakefs.example.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key',
config=config
)
print("Client with checksum configuration initialized")
기본 작업
오브젝트 업로드
import boto3
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# Upload from bytes
data = b"Hello, lakeFS!"
s3.put_object(
Bucket='my-repo',
Key='main/data/hello.txt',
Body=data
)
# Upload from file
with open('local_file.csv', 'rb') as f:
s3.put_object(
Bucket='my-repo',
Key='main/data/imported.csv',
Body=f
)
# Upload with metadata
s3.put_object(
Bucket='my-repo',
Key='main/data/data.csv',
Body=b'id,name\n1,Alice\n2,Bob',
Metadata={
'owner': 'data-team',
'version': '1.0'
}
)
print("Upload complete")
오브젝트 다운로드
import boto3
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# Download entire object
response = s3.get_object(
Bucket='my-repo',
Key='main/data/data.csv'
)
data = response['Body'].read()
print(f"Downloaded {len(data)} bytes")
# Download to file
s3.download_file(
Bucket='my-repo',
Key='main/data/large_file.parquet',
Filename='local_file.parquet'
)
# Stream download (for large files)
response = s3.get_object(Bucket='my-repo', Key='main/data/large.csv')
for chunk in iter(lambda: response['Body'].read(1024), b''):
process_chunk(chunk)
오브젝트 목록 조회
import boto3
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# List objects in branch
response = s3.list_objects_v2(
Bucket='my-repo',
Prefix='main/data/'
)
for obj in response.get('Contents', []):
print(f"{obj['Key']} ({obj['Size']} bytes)")
# List objects at commit
response = s3.list_objects_v2(
Bucket='my-repo',
Prefix='abc123def456/data/'
)
for obj in response.get('Contents', []):
print(f"{obj['Key']}")
오브젝트 메타데이터 조회
import boto3
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# Head object
response = s3.head_object(
Bucket='my-repo',
Key='main/data/file.csv'
)
print(f"Content Type: {response.get('ContentType')}")
print(f"Content Length: {response.get('ContentLength')}")
print(f"Last Modified: {response.get('LastModified')}")
print(f"Metadata: {response.get('Metadata')}")
오브젝트 삭제
import boto3
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# Delete single object
s3.delete_object(
Bucket='my-repo',
Key='main/data/temp.txt'
)
# Delete multiple objects
s3.delete_objects(
Bucket='my-repo',
Delete={
'Objects': [
{'Key': 'main/data/file1.txt'},
{'Key': 'main/data/file2.txt'},
{'Key': 'main/data/file3.txt'}
]
}
)
print("Delete complete")
실전 워크플로
S3 스타일 작업으로 ETL 하기
import boto3
import csv
import io
def etl_pipeline():
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# Extract: Read from source
response = s3.get_object(Bucket='my-repo', Key='main/raw/input.csv')
input_data = response['Body'].read().decode()
# Transform: Process data
reader = csv.DictReader(io.StringIO(input_data))
rows = list(reader)
# Clean: Remove duplicates
unique_rows = {row['id']: row for row in rows}.values()
# Load: Write processed data
output = io.StringIO()
writer = csv.DictWriter(output, fieldnames=['id', 'name', 'value'])
writer.writeheader()
writer.writerows(unique_rows)
s3.put_object(
Bucket='my-repo',
Key='main/processed/output.csv',
Body=output.getvalue()
)
print(f"ETL complete: {len(unique_rows)} unique records")
etl_pipeline()
백업과 동기화
import boto3
import os
from pathlib import Path
def backup_to_lakeFS(local_dir, repo, branch, prefix):
"""Backup local directory to lakeFS"""
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
count = 0
for local_file in Path(local_dir).rglob('*'):
if local_file.is_file():
# Calculate remote path
rel_path = local_file.relative_to(local_dir)
remote_path = f"{branch}/{prefix}/{rel_path}".replace("\\", "/")
# Upload
with open(local_file, 'rb') as f:
s3.put_object(
Bucket=repo,
Key=remote_path,
Body=f
)
count += 1
if count % 100 == 0:
print(f"Backed up {count} files...")
print(f"Backup complete: {count} files uploaded")
# Usage:
# backup_to_lakeFS("/path/to/local/data", "my-repo", "main", "backups/2024-01")
S3에서 lakeFS로 복사
import boto3
def migrate_s3_to_lakefs(s3_bucket, prefix, repo, branch):
"""Migrate data from S3 to lakeFS"""
# Connect to S3
s3_source = boto3.client(
's3',
region_name='us-east-1'
)
# Connect to lakeFS
s3_dest = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# List objects
paginator = s3_source.get_paginator('list_objects_v2')
pages = paginator.paginate(Bucket=s3_bucket, Prefix=prefix)
count = 0
for page in pages:
for obj in page.get('Contents', []):
# Download from S3
response = s3_source.get_object(
Bucket=s3_bucket,
Key=obj['Key']
)
data = response['Body'].read()
# Upload to lakeFS
s3_dest.put_object(
Bucket=repo,
Key=f"{branch}/{obj['Key']}",
Body=data
)
count += 1
if count % 100 == 0:
print(f"Migrated {count} objects...")
print(f"Migration complete: {count} objects")
# Usage:
# migrate_s3_to_lakefs("my-s3-bucket", "data/", "my-repo", "main")
버전별 접근
import boto3
def read_from_commit(repo, commit_id, key):
"""Read object from specific commit"""
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
# Use commit ID as prefix
response = s3.get_object(
Bucket=repo,
Key=f"{commit_id}/{key}"
)
return response['Body'].read()
# Usage:
# data = read_from_commit("my-repo", "abc123def456", "data/file.csv")
오류 처리
import boto3
from botocore.exceptions import ClientError
s3 = boto3.client(
's3',
endpoint_url='https://example.lakefs.io',
aws_access_key_id='your-access-key',
aws_secret_access_key='your-secret-key'
)
try:
s3.put_object(
Bucket='my-repo',
Key='main/data/file.txt',
Body=b'data'
)
except ClientError as e:
error_code = e.response['Error']['Code']
if error_code == 'AccessDenied':
print("Access denied - check credentials or permissions")
elif error_code == 'NoSuchBucket':
print("Bucket not found - check repository name")
else:
print(f"Error: {error_code}")
더 볼 자료
-
lakeFS S3 Gateway - S3 Gateway API 문서
-
Boto3 Documentation - Boto3 공식 레퍼런스
더 알아보기 (Learn more)
공식 문서의 자세한 내용은 https://docs.lakefs.io/reference/python/boto/에서 확인하실 수 있어요.