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

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/에서 확인하실 수 있어요.