Boto3: Python SDK для S3-операций - листинг, загрузка, Presigned URLs и управление метаданными
Глубокий разбор Boto3 как инструмента администрирования S3-совместимого хранилища из Python: архитектура клиента и ресурса, операции с объектами и бакетами, Presigned URLs для временного доступа, управление метаданными, жизненный цикл объектов и интеграция с Spark-пайплайнами.
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/. Нужно:
- Проверять появление новых файлов каждые 15 минут
- Валидировать файлы (размер, формат)
- Запускать Spark Job для конвертации в Parquet
- Перемещать обработанные файлы в архив
- Уведомлять партнёра о завершении через 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-вызов.