Boto3: Python SDK для S3-операций - листинг, загрузка, Presigned URLs и управление метаданными

Глубокий разбор Boto3 как инструмента администрирования S3-совместимого хранилища из Python: архитектура клиента и ресурса, операции с объектами и бакетами, Presigned URLs для временного доступа, управление метаданными, жизненный цикл объектов и интеграция с Spark-пайплайнами.

storage

Boto3 в контексте Data Engineering

Apache Spark - это инструмент для обработки данных. Но в реальных пайплайнах вокруг Spark-задач существует целый слой административных операций, которые Spark не делает по умолчанию:

  • Проверить, существует ли файл перед запуском Spark Job
  • Удалить временные файлы после завершения задачи
  • Скопировать данные из «входящей» папки в «архивную» после обработки
  • Выдать временную ссылку на файл для партнёра (Presigned URL)
  • Настроить автоматическое удаление данных старше N дней (Lifecycle Policy)
  • Проверить метаданные файла: размер, дату изменения, тег, ETag
  • Восстановить объект из Glacier в Instant Retrieval для срочной аналитики

Всё это - задачи Boto3, официального Python SDK для AWS и S3-совместимых хранилищ.

Boto3 - это не просто «библиотека для AWS». Это стандартный инструмент взаимодействия с любым S3-совместимым хранилищем: MinIO, Ceph RGW, Selectel Object Storage, VK Cloud Object Storage, Yandex Object Storage, Cloud.ru Object Storage. Все они реализуют AWS S3 API, и Boto3 работает с ними через параметр endpoint_url.

Ключевая разница между Boto3 и Spark S3A Connector: Spark использует S3 как distributed filesystem для параллельного чтения/записи больших данных; Boto3 - для административных операций: управление объектами, метаданными, правами доступа, lifecycle.


Архитектура Boto3: Client vs Resource

Boto3 предоставляет два уровня API для работы с S3:

Client - низкоуровневый, соответствует AWS REST API 1:1. Каждый метод - один HTTP-запрос. Возвращает словари Python. Подходит для точного контроля над запросами и работы с paginator'ами.

Resource - высокоуровневый, объектно-ориентированный. Скрывает HTTP-детали за Pythonic интерфейсом (bucket.objects.all(), s3_object.download_file()). Удобнее для простых операций.

import boto3
from botocore.config import Config


def create_s3_client(
    endpoint_url: str,
    access_key: str,
    secret_key: str,
    region: str = "us-east-1",
    max_pool_connections: int = 50,
    connect_timeout: int = 10,
    read_timeout: int = 30
) -> boto3.client:
    """
    Создаёт production-ready S3 client с правильными таймаутами
    и connection pool для высоконагруженных пайплайнов.

    max_pool_connections: размер пула HTTP-соединений.
    По умолчанию 10 - мало для параллельных операций.
    50 = баланс между memory usage и throughput.

    connect_timeout: максимальное время ожидания установки соединения.
    read_timeout: максимальное время ожидания ответа от сервера.
    """
    config = Config(
        max_pool_connections=max_pool_connections,
        connect_timeout=connect_timeout,
        read_timeout=read_timeout,
        retries={
            "max_attempts": 3,         # автоматический retry при 5xx ошибках
            "mode": "adaptive"         # экспоненциальная задержка между попытками
        }
    )

    return boto3.client(
        "s3",
        endpoint_url=endpoint_url,
        aws_access_key_id=access_key,
        aws_secret_access_key=secret_key,
        region_name=region,
        config=config
    )


def create_s3_resource(
    endpoint_url: str,
    access_key: str,
    secret_key: str,
    region: str = "us-east-1"
) -> boto3.resource:
    """
    Создаёт S3 Resource для высокоуровневых операций.
    Resource удобнее для итерации по объектам, загрузки/скачивания файлов.
    """
    return boto3.resource(
        "s3",
        endpoint_url=endpoint_url,
        aws_access_key_id=access_key,
        aws_secret_access_key=secret_key,
        region_name=region
    )


# ==========================================
# Примеры для MinIO (self-hosted)
# ==========================================
MINIO_CLIENT = create_s3_client(
    endpoint_url="http://minio:9000",
    access_key="minioadmin",
    secret_key="minioadmin",
    region="us-east-1"   # MinIO принимает любой регион
)

MINIO_RESOURCE = create_s3_resource(
    endpoint_url="http://minio:9000",
    access_key="minioadmin",
    secret_key="minioadmin"
)

# ==========================================
# Примеры для AWS S3 (managed)
# ==========================================
# В AWS не нужен endpoint_url - boto3 использует стандартные endpoints
AWS_S3 = boto3.client(
    "s3",
    region_name="eu-central-1"
    # Credentials берутся из env AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY
    # или из ~/.aws/credentials, или из IAM Role на EC2/EKS
)

Операции с бакетами (Bucket Operations)

Бакет - контейнер верхнего уровня в S3. В production обычно бакеты создаются один раз через Terraform или Admin UI, но иногда нужно создавать и управлять ими из кода.

import boto3
from botocore.exceptions import ClientError


def ensure_bucket_exists(s3_client, bucket_name: str, region: str = "us-east-1") -> bool:
    """
    Создаёт бакет, если он не существует.
    Возвращает True если бакет создан, False если уже существовал.

    ClientError с кодом 'BucketAlreadyOwnedByYou': бакет уже ваш - OK.
    ClientError с кодом 'BucketAlreadyExists': бакет чужой - ошибка.
    """
    try:
        # head_bucket: проверяет существование без получения содержимого
        # Не требует List permissions, только s3:GetBucketLocation
        s3_client.head_bucket(Bucket=bucket_name)
        print(f"Бакет уже существует: {bucket_name}")
        return False
    except ClientError as e:
        error_code = e.response["Error"]["Code"]

        if error_code == "404":
            # Бакет не существует - создаём
            if region == "us-east-1":
                # us-east-1 - особый случай: нельзя передавать LocationConstraint
                s3_client.create_bucket(Bucket=bucket_name)
            else:
                s3_client.create_bucket(
                    Bucket=bucket_name,
                    CreateBucketConfiguration={"LocationConstraint": region}
                )
            print(f"Бакет создан: {bucket_name}")
            return True
        elif error_code == "403":
            raise PermissionError(f"Нет доступа к бакету: {bucket_name}")
        else:
            raise


def list_buckets(s3_client) -> list[dict]:
    """
    Возвращает список всех бакетов в аккаунте.
    Каждый элемент: {'Name': str, 'CreationDate': datetime}
    """
    response = s3_client.list_buckets()
    buckets = response.get("Buckets", [])

    print(f"Всего бакетов: {len(buckets)}")
    for bucket in buckets:
        print(f"  {bucket['Name']} (создан: {bucket['CreationDate'].strftime('%Y-%m-%d')})")

    return buckets


def get_bucket_size(s3_client, bucket_name: str) -> dict:
    """
    Подсчитывает суммарный размер и число объектов в бакете.
    Использует CloudWatch метрики (AWS) или обходит все объекты (S3-compatible).
    """
    total_size = 0
    total_count = 0

    paginator = s3_client.get_paginator("list_objects_v2")
    pages = paginator.paginate(Bucket=bucket_name)

    for page in pages:
        for obj in page.get("Contents", []):
            total_size += obj["Size"]
            total_count += 1

    return {
        "bucket": bucket_name,
        "object_count": total_count,
        "total_size_bytes": total_size,
        "total_size_gb": total_size / 1024**3
    }

Листинг объектов: ListObjectsV2 и Paginator

Листинг - одна из самых частых операций. S3 возвращает максимум 1000 объектов за один запрос. Для получения всех объектов нужен paginator.

import boto3
from datetime import datetime


def list_objects_paginated(
    s3_client,
    bucket: str,
    prefix: str = "",
    max_keys_per_page: int = 1000
) -> list[dict]:
    """
    Полный листинг объектов с автоматической пагинацией.

    prefix: фильтр по префиксу ключа (например "events/2024/01/")
    Каждый объект: {Key, Size, LastModified, ETag, StorageClass}
    """
    paginator = s3_client.get_paginator("list_objects_v2")

    all_objects = []
    page_count = 0

    for page in paginator.paginate(
        Bucket=bucket,
        Prefix=prefix,
        PaginationConfig={"MaxItems": 1_000_000, "PageSize": max_keys_per_page}
    ):
        page_count += 1
        objects = page.get("Contents", [])
        all_objects.extend(objects)

        if page_count % 10 == 0:
            print(f"Обработано страниц: {page_count}, объектов: {len(all_objects)}")

    print(f"Итого: {len(all_objects)} объектов в {page_count} страницах")
    return all_objects


def list_objects_by_date_range(
    s3_client,
    bucket: str,
    prefix: str,
    start_date: datetime,
    end_date: datetime
) -> list[dict]:
    """
    Фильтрует объекты по дате последнего изменения.
    S3 не поддерживает серверный фильтр по дате - фильтруем на клиенте.
    """
    all_objects = list_objects_paginated(s3_client, bucket, prefix)

    filtered = [
        obj for obj in all_objects
        if start_date <= obj["LastModified"].replace(tzinfo=None) <= end_date
    ]

    return filtered


def list_common_prefixes(s3_client, bucket: str, prefix: str, delimiter: str = "/") -> list[str]:
    """
    Листинг «виртуальных директорий» через delimiter.

    Пример: prefix="events/", delimiter="/"
    Вернёт: ["events/2024/", "events/2023/", "events/2022/"]
    Это аналог ls без рекурсии.

    Используется для обхода партиционированных таблиц без полного listing'а.
    """
    response = s3_client.list_objects_v2(
        Bucket=bucket,
        Prefix=prefix,
        Delimiter=delimiter,
        MaxKeys=1000
    )

    prefixes = [cp["Prefix"] for cp in response.get("CommonPrefixes", [])]
    print(f"Найдено «директорий»: {len(prefixes)}")
    for p in prefixes[:10]:
        print(f"  {p}")

    return prefixes


def find_small_files(
    s3_client,
    bucket: str,
    prefix: str,
    threshold_mb: float = 32.0
) -> list[dict]:
    """
    Находит файлы меньше threshold_mb мегабайт.
    Используется для аудита и принятия решения о компакции.
    """
    threshold_bytes = int(threshold_mb * 1024 * 1024)
    all_objects = list_objects_paginated(s3_client, bucket, prefix)

    small_files = [
        obj for obj in all_objects
        if obj["Size"] < threshold_bytes
        and not obj["Key"].endswith("/")   # исключаем «пустые директории»
    ]

    total_size = sum(obj["Size"] for obj in small_files)
    print(f"Мелких файлов (< {threshold_mb} МБ): {len(small_files)}")
    print(f"Суммарный объём: {total_size / 1024**2:.2f} МБ")

    return small_files

Загрузка и скачивание объектов

import os
from pathlib import Path


def upload_file(
    s3_client,
    local_path: str,
    bucket: str,
    s3_key: str,
    metadata: dict = None,
    content_type: str = None,
    storage_class: str = "STANDARD"
) -> None:
    """
    Загружает локальный файл в S3.

    storage_class:
      STANDARD          - горячие данные (частый доступ)
      STANDARD_IA       - холодные данные (редкий доступ, дешевле хранение)
      ONEZONE_IA        - один AZ, ещё дешевле, меньше надёжность
      GLACIER           - архив, доступ через 3-5 часов
      GLACIER_IR        - архив с мгновенным доступом (Instant Retrieval)
      DEEP_ARCHIVE      - долгосрочный архив, доступ через 12 часов

    metadata: произвольные key-value пары, хранятся с объектом
    Ограничение: до 2 КБ суммарно, только ASCII-ключи
    """
    extra_args = {"StorageClass": storage_class}

    if metadata:
        # Метаданные преобразуются в строки
        extra_args["Metadata"] = {k: str(v) for k, v in metadata.items()}

    if content_type:
        extra_args["ContentType"] = content_type

    file_size = os.path.getsize(local_path)
    print(f"Загружаем: {local_path} ({file_size / 1024**2:.2f} МБ) → s3://{bucket}/{s3_key}")

    s3_client.upload_file(
        Filename=local_path,
        Bucket=bucket,
        Key=s3_key,
        ExtraArgs=extra_args
    )
    print(f"Загружено: s3://{bucket}/{s3_key}")


def upload_fileobj(s3_client, file_obj, bucket: str, s3_key: str) -> None:
    """
    Загружает file-like объект (io.BytesIO, открытый файл) в S3.
    Удобно для потоковой записи без создания временных файлов.
    """
    s3_client.upload_fileobj(
        Fileobj=file_obj,
        Bucket=bucket,
        Key=s3_key
    )


def download_file(
    s3_client,
    bucket: str,
    s3_key: str,
    local_path: str
) -> dict:
    """
    Скачивает объект из S3 в локальный файл.
    Возвращает метаданные объекта.
    """
    Path(local_path).parent.mkdir(parents=True, exist_ok=True)

    # Сначала получаем размер для прогресс-бара
    head = s3_client.head_object(Bucket=bucket, Key=s3_key)
    size_bytes = head["ContentLength"]
    print(f"Скачиваем: s3://{bucket}/{s3_key} ({size_bytes / 1024**2:.2f} МБ)")

    s3_client.download_file(
        Bucket=bucket,
        Key=s3_key,
        Filename=local_path
    )
    print(f"Скачано: {local_path}")

    return head


def copy_object(
    s3_client,
    source_bucket: str,
    source_key: str,
    dest_bucket: str,
    dest_key: str,
    storage_class: str = "STANDARD"
) -> None:
    """
    Копирует объект внутри S3 (server-side copy).
    НЕ скачивает данные через клиент - S3 сам копирует на стороне сервера.
    Это единственный способ «переместить» объект (move = copy + delete source).

    Ограничение: server-side copy работает только для объектов до 5 ГБ.
    Для больших объектов нужен multipart copy.
    """
    copy_source = {"Bucket": source_bucket, "Key": source_key}

    s3_client.copy_object(
        CopySource=copy_source,
        Bucket=dest_bucket,
        Key=dest_key,
        StorageClass=storage_class
    )
    print(f"Скопировано: s3://{source_bucket}/{source_key} → s3://{dest_bucket}/{dest_key}")


def move_object(
    s3_client,
    source_bucket: str,
    source_key: str,
    dest_bucket: str,
    dest_key: str
) -> None:
    """
    Перемещает объект (copy + delete source).
    На S3 нет нативной операции move - это всегда два шага.
    Это та же проблема что и rename, только мы делаем её явно.
    """
    copy_object(s3_client, source_bucket, source_key, dest_bucket, dest_key)
    s3_client.delete_object(Bucket=source_bucket, Key=source_key)
    print(f"Перемещено: {source_key}{dest_key}")

Удаление объектов: delete_object и delete_objects

from typing import Union


def delete_object(s3_client, bucket: str, key: str) -> None:
    """
    Удаляет один объект.
    S3 не возвращает ошибку если объект не существует - это идемпотентная операция.
    """
    s3_client.delete_object(Bucket=bucket, Key=key)
    print(f"Удалено: s3://{bucket}/{key}")


def delete_objects_batch(
    s3_client,
    bucket: str,
    keys: list[str],
    batch_size: int = 1000
) -> dict:
    """
    Пакетное удаление объектов через delete_objects API.
    S3 принимает максимум 1000 объектов за один запрос.

    delete_objects значительно эффективнее, чем цикл из delete_object:
    1000 удалений = 1 HTTP-запрос вместо 1000.

    Возвращает: {'deleted': count, 'errors': list_of_errors}
    """
    total_deleted = 0
    all_errors = []

    # Разбиваем на батчи по 1000
    for i in range(0, len(keys), batch_size):
        batch = keys[i:i + batch_size]
        objects_to_delete = [{"Key": k} for k in batch]

        response = s3_client.delete_objects(
            Bucket=bucket,
            Delete={
                "Objects": objects_to_delete,
                "Quiet": True   # Quiet=True: в ответе только ошибки, не успехи
                                # Экономит трафик при удалении тысяч объектов
            }
        )

        deleted_count = len(batch) - len(response.get("Errors", []))
        total_deleted += deleted_count

        if response.get("Errors"):
            all_errors.extend(response["Errors"])
            for err in response["Errors"]:
                print(f"Ошибка удаления {err['Key']}: {err['Message']}")

        print(f"Удалено {total_deleted}/{len(keys)} объектов...")

    print(f"Итого удалено: {total_deleted}, ошибок: {len(all_errors)}")
    return {"deleted": total_deleted, "errors": all_errors}


def delete_prefix(s3_client, bucket: str, prefix: str, dry_run: bool = True) -> int:
    """
    Удаляет все объекты с указанным префиксом.
    По умолчанию dry_run=True: только показывает что будет удалено.

    ВНИМАНИЕ: операция необратима! Всегда проверяйте с dry_run=True первым.
    """
    all_objects = list_objects_paginated(s3_client, bucket, prefix)
    keys_to_delete = [obj["Key"] for obj in all_objects]

    print(f"Объектов с префиксом '{prefix}': {len(keys_to_delete)}")

    if dry_run:
        print("DRY RUN: объекты НЕ удалены. Запустите с dry_run=False для удаления.")
        for key in keys_to_delete[:10]:
            print(f"  {key}")
        if len(keys_to_delete) > 10:
            print(f"  ... и ещё {len(keys_to_delete) - 10}")
        return len(keys_to_delete)

    result = delete_objects_batch(s3_client, bucket, keys_to_delete)
    return result["deleted"]


# Пример: удаление временных файлов после Spark Job
delete_prefix(
    MINIO_CLIENT,
    bucket="analytics",
    prefix="events/_spark_tmp/",
    dry_run=False   # только после проверки с dry_run=True!
)

Работа с метаданными объектов

S3 позволяет хранить произвольные метаданные вместе с каждым объектом. Это полезно для tracking'а версий, пометки статуса обработки, хранения checksums.

import hashlib
import json


def get_object_metadata(s3_client, bucket: str, key: str) -> dict:
    """
    Получает все атрибуты объекта через head_object.
    head_object: только заголовки ответа, без тела объекта (без данных).
    Дешевле GET, так как S3 не передаёт содержимое файла.

    Возвращает:
    - ContentLength: размер в байтах
    - ContentType: MIME-тип
    - LastModified: дата изменения (datetime)
    - ETag: MD5-хеш (кавычки в значении - это нормально: '"abc123"')
    - Metadata: пользовательские метаданные (dict)
    - StorageClass: класс хранения
    - ServerSideEncryption: тип шифрования (если есть)
    """
    try:
        response = s3_client.head_object(Bucket=bucket, Key=key)
        return {
            "key": key,
            "size_bytes": response["ContentLength"],
            "size_mb": response["ContentLength"] / 1024**2,
            "content_type": response.get("ContentType", ""),
            "last_modified": response["LastModified"],
            "etag": response["ETag"].strip('"'),   # убираем кавычки из ETag
            "storage_class": response.get("StorageClass", "STANDARD"),
            "metadata": response.get("Metadata", {}),
            "encryption": response.get("ServerSideEncryption", "None")
        }
    except ClientError as e:
        if e.response["Error"]["Code"] == "404":
            raise FileNotFoundError(f"Объект не найден: s3://{bucket}/{key}")
        raise


def set_object_metadata(
    s3_client,
    bucket: str,
    key: str,
    new_metadata: dict
) -> None:
    """
    Обновляет пользовательские метаданные объекта.
    S3 не поддерживает partial metadata update - нужно скопировать объект
    с новыми метаданными (server-side copy на тот же ключ).

    Это COPY-на-месте: данные не перемещаются физически.
    Стоимость: один CopyObject запрос.
    """
    # Получаем текущие метаданные
    current = get_object_metadata(s3_client, bucket, key)
    merged_metadata = {**current["metadata"], **new_metadata}

    # CopyObject с replace metadata
    s3_client.copy_object(
        CopySource={"Bucket": bucket, "Key": key},
        Bucket=bucket,
        Key=key,
        MetadataDirective="REPLACE",   # REPLACE: использовать новые метаданные
        Metadata={k: str(v) for k, v in merged_metadata.items()},
        StorageClass=current["storage_class"]
    )
    print(f"Метаданные обновлены: {key}")


# Пример: маркировка файла как «обработанного»
set_object_metadata(
    MINIO_CLIENT,
    bucket="raw-data",
    key="events/2024/01/15/events_001.parquet",
    new_metadata={
        "processed": "true",
        "processed_at": "2024-01-15T10:30:00Z",
        "pipeline_version": "2.3.1",
        "row_count": "1500000"
    }
)


def tag_object(s3_client, bucket: str, key: str, tags: dict) -> None:
    """
    Устанавливает теги на объект.
    Теги отличаются от метаданных:
    - Метаданные: хранятся внутри объекта, доступны через head_object
    - Теги: внешние атрибуты, не хранятся в объекте, доступны через get_object_tagging
    - Теги можно использовать в Lifecycle Rules (удалять по тегу)
    - Теги можно обновлять без copy-on-write
    - Ограничение: до 10 тегов на объект
    """
    tag_set = [{"Key": k, "Value": str(v)} for k, v in tags.items()]
    s3_client.put_object_tagging(
        Bucket=bucket,
        Key=key,
        Tagging={"TagSet": tag_set}
    )


def get_object_tags(s3_client, bucket: str, key: str) -> dict:
    """Получает теги объекта."""
    response = s3_client.get_object_tagging(Bucket=bucket, Key=key)
    return {tag["Key"]: tag["Value"] for tag in response.get("TagSet", [])}

Presigned URLs: временный доступ без credentials

Presigned URL - это временная подписанная ссылка, которая позволяет стороннему пользователю (без AWS credentials) выполнить конкретную S3-операцию в течение ограниченного времени.

from datetime import timedelta


def generate_presigned_get_url(
    s3_client,
    bucket: str,
    key: str,
    expires_in_seconds: int = 3600
) -> str:
    """
    Генерирует presigned URL для скачивания объекта.
    URL действителен expires_in_seconds секунд.
    Максимальный срок: 7 дней (604800 секунд) для AWS,
    MinIO поддерживает до 7 дней по умолчанию.

    Используется:
    - Выдача файлов партнёрам без создания им credentials
    - Временный доступ к отчётам/выгрузкам для внутренних пользователей
    - Webhook: пайплайн генерирует URL и отправляет в систему оповещений
    """
    url = s3_client.generate_presigned_url(
        ClientMethod="get_object",
        Params={"Bucket": bucket, "Key": key},
        ExpiresIn=expires_in_seconds
    )
    print(f"Presigned GET URL (действует {expires_in_seconds}с): {url[:80]}...")
    return url


def generate_presigned_put_url(
    s3_client,
    bucket: str,
    key: str,
    content_type: str = "application/octet-stream",
    expires_in_seconds: int = 3600
) -> str:
    """
    Генерирует presigned URL для загрузки объекта.
    Клиент выполняет PUT-запрос напрямую в S3, минуя ваш сервер.

    Используется в паттерне Direct Upload:
    Frontend запрашивает у Backend presigned PUT URL,
    Frontend загружает файл напрямую в S3 (без проксирования через Backend).
    Это экономит трафик на Backend-сервере.
    """
    url = s3_client.generate_presigned_url(
        ClientMethod="put_object",
        Params={
            "Bucket": bucket,
            "Key": key,
            "ContentType": content_type
        },
        ExpiresIn=expires_in_seconds
    )
    return url


def generate_presigned_post(
    s3_client,
    bucket: str,
    key_prefix: str,
    max_size_bytes: int = 100 * 1024 * 1024,   # 100 МБ
    expires_in_seconds: int = 3600
) -> dict:
    """
    Генерирует presigned POST форму с условиями загрузки.
    В отличие от presigned PUT, позволяет задавать условия:
    - max_size_bytes: ограничение на размер файла
    - allowed_prefixes: разрешённые пути
    - content_type_prefix: разрешённые MIME-типы

    Используется для пользовательской загрузки файлов через HTML-форму.
    """
    conditions = [
        ["content-length-range", 0, max_size_bytes],   # ограничение размера
        ["starts-with", "$key", key_prefix]             # ключ должен начинаться с префикса
    ]

    response = s3_client.generate_presigned_post(
        Bucket=bucket,
        Key=f"{key_prefix}${{filename}}",   # placeholder для имени файла
        Conditions=conditions,
        ExpiresIn=expires_in_seconds
    )

    # response["url"] - URL формы
    # response["fields"] - поля формы (включая подпись)
    return response


# Пример: выдача отчёта партнёру
def share_report_with_partner(
    s3_client,
    report_key: str,
    partner_email: str,
    valid_hours: int = 24
) -> str:
    """
    Генерирует временную ссылку на отчёт и логирует выдачу.
    """
    url = generate_presigned_get_url(
        s3_client,
        bucket="reports",
        key=report_key,
        expires_in_seconds=valid_hours * 3600
    )

    print(f"Ссылка выдана {partner_email}")
    print(f"Файл: {report_key}")
    print(f"Действует до: {valid_hours} часов")
    print(f"URL: {url}")

    return url

Multipart Upload: загрузка больших файлов

Для файлов больше 100 МБ рекомендуется использовать Multipart Upload (MPU). Boto3 управляет MPU автоматически через upload_file(), но иногда нужен явный контроль.

import os
import math
from concurrent.futures import ThreadPoolExecutor, as_completed


def multipart_upload(
    s3_client,
    local_file: str,
    bucket: str,
    key: str,
    part_size_mb: int = 100,
    max_workers: int = 4
) -> None:
    """
    Многопоточная загрузка большого файла через Multipart Upload.
    Каждая часть загружается параллельно - ускоряет загрузку при широком канале.

    part_size_mb: размер каждой части. AWS минимум 5 МБ (кроме последней части).
    max_workers: число параллельных потоков загрузки.

    Шаги:
    1. create_multipart_upload - инициализация, получаем UploadId
    2. upload_part × N - параллельная загрузка частей
    3. complete_multipart_upload - финализация, S3 собирает части в объект
    """
    part_size = part_size_mb * 1024 * 1024
    file_size = os.path.getsize(local_file)
    num_parts = math.ceil(file_size / part_size)

    print(f"MPU: {local_file} ({file_size / 1024**3:.2f} ГБ), {num_parts} частей")

    # Шаг 1: создаём multipart upload
    mpu = s3_client.create_multipart_upload(Bucket=bucket, Key=key)
    upload_id = mpu["UploadId"]
    print(f"UploadId: {upload_id}")

    parts = []

    def upload_part(part_number: int) -> dict:
        """Загружает одну часть файла."""
        start = (part_number - 1) * part_size
        end = min(start + part_size, file_size)
        chunk_size = end - start

        with open(local_file, "rb") as f:
            f.seek(start)
            data = f.read(chunk_size)

        response = s3_client.upload_part(
            Bucket=bucket,
            Key=key,
            UploadId=upload_id,
            PartNumber=part_number,
            Body=data
        )

        return {"PartNumber": part_number, "ETag": response["ETag"]}

    # Шаг 2: параллельная загрузка частей
    try:
        with ThreadPoolExecutor(max_workers=max_workers) as executor:
            futures = {
                executor.submit(upload_part, part_num): part_num
                for part_num in range(1, num_parts + 1)
            }

            for future in as_completed(futures):
                part_result = future.result()
                parts.append(part_result)
                print(f"Часть {part_result['PartNumber']}/{num_parts} загружена")

        # Сортируем части по номеру (порядок важен для complete)
        parts.sort(key=lambda p: p["PartNumber"])

        # Шаг 3: финализация
        s3_client.complete_multipart_upload(
            Bucket=bucket,
            Key=key,
            UploadId=upload_id,
            MultipartUpload={"Parts": parts}
        )
        print(f"MPU завершён: s3://{bucket}/{key}")

    except Exception as e:
        # При ошибке: отменяем MPU, иначе незавершённые части будут тарифицироваться
        print(f"Ошибка MPU: {e}. Отменяем upload...")
        s3_client.abort_multipart_upload(
            Bucket=bucket,
            Key=key,
            UploadId=upload_id
        )
        raise


def cleanup_incomplete_uploads(s3_client, bucket: str, older_than_days: int = 3) -> int:
    """
    Находит и удаляет незавершённые multipart uploads.
    Незавершённые MPU продолжают тарифицироваться на S3!

    Рекомендуется запускать ежедневно или настроить Lifecycle Rule.
    """
    from datetime import datetime, timezone, timedelta

    cutoff = datetime.now(timezone.utc) - timedelta(days=older_than_days)

    response = s3_client.list_multipart_uploads(Bucket=bucket)
    uploads = response.get("Uploads", [])

    old_uploads = [u for u in uploads if u["Initiated"] < cutoff]

    print(f"Незавершённых MPU: {len(uploads)}, старше {older_than_days} дней: {len(old_uploads)}")

    aborted = 0
    for upload in old_uploads:
        s3_client.abort_multipart_upload(
            Bucket=bucket,
            Key=upload["Key"],
            UploadId=upload["UploadId"]
        )
        aborted += 1
        print(f"Отменён MPU: {upload['Key']} (начат {upload['Initiated']})")

    return aborted

Lifecycle Rules: автоматическое управление жизненным циклом

S3 Lifecycle Rules позволяют автоматически переводить объекты между классами хранения или удалять их по заданным правилам - без кода.

import json


def set_lifecycle_policy(
    s3_client,
    bucket: str,
    rules: list[dict]
) -> None:
    """
    Устанавливает Lifecycle Policy на бакет.

    Пример rules:
    [
        {
            "ID": "move-old-events-to-glacier",
            "Status": "Enabled",
            "Filter": {"Prefix": "events/"},
            "Transitions": [
                {"Days": 30, "StorageClass": "STANDARD_IA"},
                {"Days": 90, "StorageClass": "GLACIER_IR"}
            ],
            "Expiration": {"Days": 365}
        }
    ]
    """
    s3_client.put_bucket_lifecycle_configuration(
        Bucket=bucket,
        LifecycleConfiguration={"Rules": rules}
    )
    print(f"Lifecycle Policy установлена на {bucket}: {len(rules)} правил")


def create_data_lifecycle_rules() -> list[dict]:
    """
    Типовой набор Lifecycle Rules для аналитической платформы:
    1. Горячие данные (events/): STANDARD → STANDARD_IA через 30 дней → удаление через 2 года
    2. Промежуточные данные (tmp/): удаление через 7 дней
    3. Незавершённые MPU: отмена через 3 дня
    """
    return [
        {
            "ID": "hot-events-tiering",
            "Status": "Enabled",
            "Filter": {"Prefix": "events/"},
            "Transitions": [
                {"Days": 30, "StorageClass": "STANDARD_IA"},
                {"Days": 90, "StorageClass": "GLACIER_IR"}
            ],
            "Expiration": {"Days": 730}   # 2 года
        },
        {
            "ID": "cleanup-tmp-files",
            "Status": "Enabled",
            "Filter": {"Prefix": "tmp/"},
            "Expiration": {"Days": 7}
        },
        {
            "ID": "cleanup-spark-staging",
            "Status": "Enabled",
            "Filter": {"Prefix": "_spark_staging/"},
            "Expiration": {"Days": 1}
        },
        {
            "ID": "abort-incomplete-multipart",
            "Status": "Enabled",
            "Filter": {"Prefix": ""},   # применяется ко всему бакету
            "AbortIncompleteMultipartUpload": {"DaysAfterInitiation": 3}
        }
    ]


# Устанавливаем lifecycle policy
set_lifecycle_policy(
    MINIO_CLIENT,
    bucket="analytics",
    rules=create_data_lifecycle_rules()
)

Интеграция Boto3 с Spark-пайплайнами

В реальных Data Engineering пайплайнах Boto3 и Spark работают вместе: Boto3 выполняет административные операции до и после Spark-задач.

from pyspark.sql import SparkSession
import boto3


class DataPipelineOrchestrator:
    """
    Оркестратор пайплайна обработки данных.
    Использует Boto3 для проверки входных данных и очистки,
    Spark - для трансформации.
    """

    def __init__(
        self,
        s3_client,
        spark: SparkSession,
        raw_bucket: str,
        processed_bucket: str
    ):
        self.s3 = s3_client
        self.spark = spark
        self.raw_bucket = raw_bucket
        self.processed_bucket = processed_bucket

    def check_input_files_exist(self, prefix: str, min_count: int = 1) -> bool:
        """
        Проверяет наличие входных файлов перед запуском Spark Job.
        Boto3 намного быстрее Spark для простой проверки существования файлов.
        Spark запускается 30-60 секунд только для старта JVM.
        Boto3 делает один head_object за ~50 мс.
        """
        try:
            response = self.s3.list_objects_v2(
                Bucket=self.raw_bucket,
                Prefix=prefix,
                MaxKeys=min_count
            )
            found = response.get("KeyCount", 0)
            print(f"Найдено файлов в {prefix}: {found} (минимум: {min_count})")
            return found >= min_count
        except ClientError:
            return False

    def get_unprocessed_files(self, prefix: str) -> list[str]:
        """
        Находит файлы, которые ещё не обработаны.
        Проверяет тег 'processed' на каждом объекте.
        """
        all_objects = list_objects_paginated(self.s3, self.raw_bucket, prefix)
        unprocessed = []

        for obj in all_objects:
            tags = get_object_tags(self.s3, self.raw_bucket, obj["Key"])
            if tags.get("processed") != "true":
                unprocessed.append(obj["Key"])

        print(f"Необработанных файлов: {len(unprocessed)}/{len(all_objects)}")
        return unprocessed

    def run_spark_processing(self, input_keys: list[str], output_prefix: str) -> int:
        """
        Запускает Spark для обработки файлов.
        Возвращает число обработанных строк.
        """
        # Строим пути S3 для чтения
        input_paths = [f"s3a://{self.raw_bucket}/{k}" for k in input_keys]

        df = self.spark.read.parquet(*input_paths)
        row_count = df.count()

        # Трансформации
        result = df.filter(df["event_type"].isNotNull()) \
                   .dropDuplicates(["event_id"]) \
                   .repartition(10)   # оптимизируем число выходных файлов

        output_path = f"s3a://{self.processed_bucket}/{output_prefix}"
        result.write.mode("overwrite").parquet(output_path)

        return row_count

    def mark_as_processed(self, keys: list[str], pipeline_version: str) -> None:
        """
        После успешной обработки помечает входные файлы тегом 'processed'.
        Это идемпотентная операция - повторный запуск не обработает их снова.
        """
        from datetime import datetime

        for key in keys:
            tag_object(self.s3, self.raw_bucket, key, {
                "processed": "true",
                "processed_at": datetime.utcnow().isoformat(),
                "pipeline_version": pipeline_version
            })

    def cleanup_temp_files(self, temp_prefix: str) -> None:
        """Удаляет временные файлы Spark после завершения задачи."""
        count = delete_prefix(
            self.s3, self.processed_bucket, temp_prefix, dry_run=False
        )
        print(f"Удалено временных файлов: {count}")

    def run_full_pipeline(self, date: str, pipeline_version: str = "1.0.0") -> dict:
        """
        Полный pipeline:
        1. Проверяем входные файлы
        2. Находим необработанные
        3. Запускаем Spark
        4. Помечаем как обработанные
        5. Очищаем временные файлы
        """
        prefix = f"events/{date}/"
        temp_prefix = f"events/{date}/_spark_staging/"

        print(f"\n{'='*50}")
        print(f"Pipeline запущен: {date}")

        # Шаг 1: проверка входных данных
        if not self.check_input_files_exist(prefix):
            print("Нет входных файлов - пропускаем")
            return {"status": "skipped", "date": date}

        # Шаг 2: поиск необработанных файлов
        unprocessed = self.get_unprocessed_files(prefix)
        if not unprocessed:
            print("Все файлы уже обработаны - пропускаем")
            return {"status": "already_processed", "date": date}

        # Шаг 3: Spark обработка
        row_count = self.run_spark_processing(
            input_keys=unprocessed,
            output_prefix=f"processed/events/{date}/"
        )

        # Шаг 4: пометка как обработанных
        self.mark_as_processed(unprocessed, pipeline_version)

        # Шаг 5: очистка
        self.cleanup_temp_files(temp_prefix)

        result = {
            "status": "success",
            "date": date,
            "files_processed": len(unprocessed),
            "rows_processed": row_count
        }
        print(f"Pipeline завершён: {result}")
        return result

Мониторинг и аудит через Boto3

from collections import defaultdict


def audit_bucket_structure(
    s3_client,
    bucket: str,
    prefix: str = ""
) -> dict:
    """
    Полный аудит структуры бакета:
    - Распределение файлов по размеру
    - Статистика по классам хранения
    - Топ-10 самых крупных файлов
    - Оценка ежемесячных расходов
    """
    # Тарифы AWS S3 (приблизительные, us-east-1)
    storage_prices = {
        "STANDARD":     0.023,   # $/ГБ·мес
        "STANDARD_IA":  0.0125,
        "ONEZONE_IA":   0.01,
        "GLACIER_IR":   0.004,
        "GLACIER":      0.0036,
        "DEEP_ARCHIVE": 0.00099,
        "UNKNOWN":      0.023    # если класс не указан
    }

    all_objects = list_objects_paginated(s3_client, bucket, prefix)

    # Распределение по размеру
    size_buckets = defaultdict(int)
    storage_class_stats = defaultdict(lambda: {"count": 0, "size": 0})
    total_size = 0

    for obj in all_objects:
        size = obj["Size"]
        storage_class = obj.get("StorageClass", "STANDARD")

        total_size += size
        storage_class_stats[storage_class]["count"] += 1
        storage_class_stats[storage_class]["size"] += size

        # Классификация по размеру
        if size < 1024:                  size_buckets["< 1 КБ"] += 1
        elif size < 100 * 1024:          size_buckets["1 КБ – 100 КБ"] += 1
        elif size < 10 * 1024**2:        size_buckets["100 КБ – 10 МБ"] += 1
        elif size < 128 * 1024**2:       size_buckets["10 МБ – 128 МБ"] += 1
        elif size < 512 * 1024**2:       size_buckets["128 МБ – 512 МБ"] += 1
        else:                            size_buckets["> 512 МБ"] += 1

    # Топ-10 крупнейших файлов
    top_files = sorted(all_objects, key=lambda o: o["Size"], reverse=True)[:10]

    # Оценка ежемесячных расходов
    monthly_cost = sum(
        stats["size"] / 1024**3 * storage_prices.get(cls, 0.023)
        for cls, stats in storage_class_stats.items()
    )

    report = {
        "bucket": bucket,
        "prefix": prefix,
        "total_objects": len(all_objects),
        "total_size_gb": total_size / 1024**3,
        "size_distribution": dict(size_buckets),
        "storage_class_stats": {
            cls: {
                "count": s["count"],
                "size_gb": s["size"] / 1024**3,
                "price_per_gb": storage_prices.get(cls, 0.023)
            }
            for cls, s in storage_class_stats.items()
        },
        "estimated_monthly_cost_usd": monthly_cost,
        "top_10_files": [
            {"key": f["Key"], "size_mb": f["Size"] / 1024**2}
            for f in top_files
        ]
    }

    # Вывод
    print(f"\n{'='*60}")
    print(f"Аудит бакета: s3://{bucket}/{prefix}")
    print(f"Объектов: {report['total_objects']:,}")
    print(f"Объём: {report['total_size_gb']:.2f} ГБ")
    print(f"Ежемесячная стоимость хранения: ${monthly_cost:.2f}")
    print("\nРаспределение по размеру:")
    for range_name, count in sorted(size_buckets.items()):
        bar = "█" * min(count * 20 // max(size_buckets.values(), 1), 20)
        print(f"  {range_name:20s}: {count:6,} {bar}")

    return report

Anti-Patterns и типичные ошибки

Anti-Pattern 1: Не использовать Paginator

# НЕПРАВИЛЬНО: list_objects_v2 возвращает максимум 1000 объектов
# Для бакета с 50 000 объектами вы получите только первые 1000!
response = s3_client.list_objects_v2(Bucket="analytics", Prefix="events/")
objects = response.get("Contents", [])   # максимум 1000!

# ПРАВИЛЬНО: всегда используйте Paginator
paginator = s3_client.get_paginator("list_objects_v2")
all_objects = []
for page in paginator.paginate(Bucket="analytics", Prefix="events/"):
    all_objects.extend(page.get("Contents", []))

Anti-Pattern 2: Не обрабатывать незавершённые MPU

# НЕПРАВИЛЬНО: если upload упадёт, части останутся на S3 и будут тарифицироваться
# Незавершённый MPU 10 ГБ = 10 ГБ × $0.023 = $0.23/мес навсегда
s3_client.upload_file("huge_file.parquet", "bucket", "key")
# При сбое: части остались, но объект не создан

# ПРАВИЛЬНО: настройте Lifecycle Rule для отмены старых MPU
set_lifecycle_policy(s3_client, "bucket", [{
    "ID": "abort-stale-mpu",
    "Status": "Enabled",
    "Filter": {"Prefix": ""},
    "AbortIncompleteMultipartUpload": {"DaysAfterInitiation": 3}
}])

Anti-Pattern 3: Не проверять ошибки при batch delete

# НЕПРАВИЛЬНО: delete_objects не бросает исключение при частичных ошибках!
# Если некоторые объекты не удалились (нет прав, нет объекта) - вы не узнаете
s3_client.delete_objects(Bucket="bucket", Delete={"Objects": [...]})

# ПРАВИЛЬНО: всегда проверять Errors в ответе
response = s3_client.delete_objects(
    Bucket="bucket",
    Delete={"Objects": [...], "Quiet": False}  # Quiet=False возвращает Deleted список
)
if response.get("Errors"):
    for err in response["Errors"]:
        print(f"Не удалось удалить {err['Key']}: {err['Code']} - {err['Message']}")
    raise RuntimeError(f"Ошибки при удалении: {len(response['Errors'])} объектов")

Anti-Pattern 4: Хранить credentials в коде

# НЕПРАВИЛЬНО: credentials в коде -> попадают в git, логи, трейсы
s3_client = boto3.client("s3",
    aws_access_key_id="AKIAIOSFODNN7EXAMPLE",
    aws_secret_access_key="wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY"
)

# ПРАВИЛЬНО: переменные окружения или secrets manager
import os
s3_client = boto3.client("s3",
    aws_access_key_id=os.environ["S3_ACCESS_KEY"],
    aws_secret_access_key=os.environ["S3_SECRET_KEY"],
    endpoint_url=os.environ.get("S3_ENDPOINT_URL")
)

# Ещё лучше для AWS: IAM Role (тогда credentials вообще не нужны в коде)
s3_client = boto3.client("s3", region_name="eu-central-1")
# Boto3 сам возьмёт credentials из EC2/EKS Instance Profile

Production кейс: автоматизация intake-пайплайна

Ситуация. Партнёры ежедневно загружают CSV-файлы в папку s3://landing-bucket/incoming/. Нужно:

  1. Проверять появление новых файлов каждые 15 минут
  2. Валидировать файлы (размер, формат)
  3. Запускать Spark Job для конвертации в Parquet
  4. Перемещать обработанные файлы в архив
  5. Уведомлять партнёра о завершении через Presigned URL готового Parquet
import time
from datetime import datetime
import requests   # для вызова webhook партнёра


def intake_pipeline(s3_client, spark: SparkSession) -> None:
    """
    Полный intake pipeline для партнёрских данных.
    """
    incoming_prefix = "incoming/"
    archive_prefix = "archive/"
    processed_prefix = "processed/"
    bucket = "landing-bucket"

    # Шаг 1: Найти необработанные CSV файлы
    all_objects = list_objects_paginated(s3_client, bucket, incoming_prefix)
    csv_files = [
        obj for obj in all_objects
        if obj["Key"].endswith(".csv")
        and get_object_tags(s3_client, bucket, obj["Key"]).get("processed") != "true"
    ]

    if not csv_files:
        print("Новых файлов нет")
        return

    print(f"Найдено CSV файлов для обработки: {len(csv_files)}")

    for obj in csv_files:
        key = obj["Key"]
        size_mb = obj["Size"] / 1024**2

        print(f"\nОбрабатываем: {key} ({size_mb:.2f} МБ)")

        # Шаг 2: Валидация
        if size_mb > 500:
            print(f"ПРЕДУПРЕЖДЕНИЕ: файл слишком большой ({size_mb:.0f} МБ)")

        if size_mb < 0.001:
            print(f"ПРОПУСК: пустой файл")
            continue

        # Шаг 3: Spark конвертация CSV → Parquet
        input_path = f"s3a://{bucket}/{key}"
        output_key = key.replace("incoming/", processed_prefix).replace(".csv", ".parquet")
        output_path = f"s3a://{bucket}/{output_key}"

        df = spark.read \
            .option("header", "true") \
            .option("inferSchema", "true") \
            .csv(input_path)

        df.coalesce(1).write.mode("overwrite").parquet(output_path)
        print(f"Конвертирован в Parquet: {output_key}")

        # Шаг 4: Перемещение в архив
        archive_key = key.replace("incoming/", archive_prefix)
        move_object(s3_client, bucket, key, bucket, archive_key)
        print(f"Файл заархивирован: {archive_key}")

        # Шаг 5: Presigned URL для партнёра
        url = generate_presigned_get_url(
            s3_client,
            bucket=bucket,
            key=output_key,
            expires_in_seconds=24 * 3600   # 24 часа
        )

        # Уведомление партнёра (webhook)
        webhook_payload = {
            "status": "completed",
            "original_file": key,
            "parquet_download_url": url,
            "row_count": df.count(),
            "processed_at": datetime.utcnow().isoformat()
        }
        print(f"Webhook отправлен партнёру с URL: {url[:60]}...")


# Запускаем каждые 15 минут через Airflow или cron
intake_pipeline(MINIO_CLIENT, spark)

Этот паттерн демонстрирует, как Boto3 и Spark дополняют друг друга: Boto3 управляет файлами и доступом (листинг, перемещение, теги, presigned URLs), Spark выполняет тяжёлую трансформацию данных. Попытка использовать только Spark для таких операций была бы избыточна и медленна: запуск Spark-кластера ради проверки «есть ли новые файлы» - оverkill, когда Boto3 делает это за 50 мс через один API-вызов.