MinIO: установка, настройка bucket policy и интеграция со Spark
Глубокий разбор MinIO как фундамента Open-Source Data Lakehouse: архитектура, Erasure Coding, Bucket Policy и IAM, полная интеграция со Spark через S3A, Delta Lake и Iceberg поверх MinIO.
MinIO как фундамент Open-Source Data Lakehouse¶
MinIO - это высокопроизводительный, S3-совместимый объектный storage с открытым исходным кодом. Он стал де-факто стандартом для построения on-premise и Kubernetes-native Data Lakehouse в тех инфраструктурах, где нельзя или нецелесообразно использовать AWS S3 или аналогичные облачные сервисы.
Почему MinIO популярен, особенно в российском enterprise:
- Полная совместимость с S3 API. Весь стек, написанный под AWS S3 (Spark S3A Connector, Delta Lake, Apache Iceberg, Apache Hudi, boto3), работает с MinIO без изменения кода - достаточно поменять endpoint.
- Открытый исходный код. MinIO распространяется по лицензии GNU AGPL v3. Это означает независимость от иностранного cloud vendor, возможность развернуть на собственном железе и полный контроль над данными.
- Экстремальная производительность. MinIO позиционирует себя как самое быстрое объектное хранилище: в бенчмарках на NVMe-дисках достигает >325 GB/s на чтение и >165 GB/s на запись на кластере из 32 нод.
- Простота операционного обслуживания. Один бинарный файл без внешних зависимостей. Можно поднять за 5 минут локально или развернуть на Kubernetes через официальный Helm chart или Operator.
- Erasure Coding вместо репликации. MinIO не хранит N полных копий объекта (как HDFS с replication=3). Вместо этого он использует Erasure Coding: объект разбивается на шарды, что даёт отказоустойчивость при меньшем расходе дискового пространства.
Место MinIO в современном Data Stack¶
Схема показывает типичный современный Data Stack, где MinIO занимает роль Storage Layer. Все три уровня Medallion Architecture (Bronze → Silver → Gold) живут в MinIO-бакетах. Поверх объектов MinIO лежат транзакционные форматы - Delta Lake или Iceberg, которые добавляют ACID-семантику. Вычислительный Spark-кластер полностью stateless: его можно перезапустить, обновить или масштабировать, не трогая данные в MinIO.
Архитектурные режимы MinIO¶
MinIO поддерживает два принципиально разных режима работы. Понимание их различий критично для правильного выбора в зависимости от задачи.
Standalone (Single-Node Single-Drive / Single-Node Multi-Drive)¶
Standalone режим предназначен для разработки, тестирования и небольших production-нагрузок.
В самой простой форме MinIO запускается как один процесс на одном сервере с одним диском. Данные хранятся в обычной файловой системе без Erasure Coding. При падении сервера или диска - данные недоступны или потеряны.
Вариант Single-Node Multi-Drive (SNMD) чуть надёжнее: один сервер, несколько дисков. MinIO использует Erasure Coding на локальных дисках одного сервера. При потере одного диска - данные восстанавливаются. Но при падении сервера - потеря доступности.
Когда использовать: локальная разработка, CI/CD тестирование, небольшие non-critical данные.
Distributed (Multi-Node Multi-Drive)¶
Distributed режим - единственный production-ready вариант для хранения важных данных.
MinIO разворачивается на N серверах (нодах), каждый с M дисками. Минимальная конфигурация для Erasure Coding: 4 ноды. Данные распределяются между нодами через Erasure Coding. Кластер сохраняет работоспособность при потере N/2 нод или N×M/2 дисков (в зависимости от схемы EC).
Ключевые свойства:
- Горизонтальное масштабирование: добавление серверов увеличивает и ёмкость, и пропускную способность линейно
- Единое пространство имён: все бакеты видны через любую ноду кластера
- Load Balancing: запросы можно направлять на любую ноду - MinIO сам маршрутизирует к нужным данным
- Healing: при возвращении упавшей ноды MinIO автоматически восстанавливает пореплицированные шарды
Когда использовать: production Data Lake, shared storage для нескольких Spark-кластеров, любые данные с требованиями к durability.
Erasure Coding: почему MinIO не использует репликацию¶
Классический подход к отказоустойчивости в HDFS - хранить три копии каждого блока (replication factor = 3). Это просто и надёжно, но расточительно: 100 GB данных занимают 300 GB дискового пространства.
MinIO использует Reed-Solomon Erasure Coding - математический метод, позволяющий восстановить данные при потере части шардов. Объект разбивается на K data shards и M parity shards:
- Для восстановления достаточно любых
Kшардов изK+M - Можно потерять до
Mшардов без потери данных - Overhead: вместо 3× имеет
(K+M)/K- например, при EC:8+4 overhead = 12/8 = 1.5×
Что это означает для Spark: при чтении одного объекта из MinIO EC:8+4 данные физически читаются с 8 серверов одновременно. Пропускная способность одного GET-запроса - это суммарная пропускная способность 8 дисков, а не одного. Это встроенный параллелизм на уровне хранилища.
Установка и развёртывание MinIO¶
Standalone через Docker (для разработки)¶
Самый быстрый способ поднять MinIO для локальной разработки - Docker:
docker run -d \
--name minio-dev \
-p 9000:9000 \
-p 9001:9001 \
-e MINIO_ROOT_USER=minioadmin \
-e MINIO_ROOT_PASSWORD=minioadmin123 \
-v minio_data:/data \
minio/minio:RELEASE.2024-11-07T00-52-20Z \
server /data --console-address ":9001"
Порт 9000 - S3 API (то, к чему подключается Spark). Порт 9001 - Web Console (браузерный UI).
После запуска: http://localhost:9001 → логин minioadmin / minioadmin123.
Production-ready Docker Compose (Distributed, 4 ноды)¶
В реальных условиях MinIO разворачивают в distributed-режиме. Пример Docker Compose для 4-нодового кластера на одном хосте (для тестирования; в prod каждая нода - отдельный сервер):
# docker-compose.yml для MinIO Distributed (4 ноды)
version: "3.8"
x-minio-common: &minio-common
image: minio/minio:RELEASE.2024-11-07T00-52-20Z
command: server
http://minio-{1...4}:9000/data
--console-address ":9001"
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER:-minioadmin}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD:-minioadmin123}
MINIO_VOLUMES: "/data"
expose:
- "9000"
- "9001"
healthcheck:
test: ["CMD", "mc", "ready", "local"]
interval: 5s
timeout: 5s
retries: 5
services:
minio-1:
<<: *minio-common
hostname: minio-1
volumes:
- minio1_data:/data
minio-2:
<<: *minio-common
hostname: minio-2
volumes:
- minio2_data:/data
minio-3:
<<: *minio-common
hostname: minio-3
volumes:
- minio3_data:/data
minio-4:
<<: *minio-common
hostname: minio-4
volumes:
- minio4_data:/data
nginx:
image: nginx:1.25-alpine
hostname: nginx
volumes:
- ./nginx.conf:/etc/nginx/nginx.conf:ro
ports:
- "9000:9000" # S3 API через load balancer
- "9001:9001" # Console
depends_on:
- minio-1
- minio-2
- minio-3
- minio-4
volumes:
minio1_data:
minio2_data:
minio3_data:
minio4_data:
Конфигурация nginx для балансировки нагрузки между нодами MinIO:
# nginx.conf
worker_processes auto;
events {
worker_connections 1024;
}
http {
upstream minio_s3 {
least_conn;
server minio-1:9000;
server minio-2:9000;
server minio-3:9000;
server minio-4:9000;
}
upstream minio_console {
least_conn;
server minio-1:9001;
server minio-2:9001;
server minio-3:9001;
server minio-4:9001;
}
server {
listen 9000;
ignore_invalid_headers off;
client_max_body_size 0;
location / {
proxy_pass http://minio_s3;
proxy_set_header Host $http_host;
proxy_set_header X-Real-IP $remote_addr;
proxy_connect_timeout 300;
proxy_http_version 1.1;
chunked_transfer_encoding off;
}
}
server {
listen 9001;
location / {
proxy_pass http://minio_console;
proxy_set_header Host $http_host;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
}
}
}
Запуск: docker compose up -d
MinIO на Kubernetes (Operator)¶
Для Kubernetes-native развёртывания MinIO предоставляет официальный Operator:
# Установка MinIO Operator через Helm
helm repo add minio-operator https://operator.min.io
helm install operator minio-operator/operator \
--namespace minio-operator \
--create-namespace
# Создание MinIO Tenant (кластер объектного storage)
cat <<EOF | kubectl apply -f -
apiVersion: minio.min.io/v2
kind: Tenant
metadata:
name: minio-data-lake
namespace: minio
spec:
pools:
- servers: 4
volumesPerServer: 4
volumeClaimTemplate:
spec:
accessModes:
- ReadWriteOnce
resources:
requests:
storage: 1Ti
storageClassName: fast-nvme
image: "minio/minio:RELEASE.2024-11-07T00-52-20Z"
requestAutoCert: false
EOF
Operator автоматически создаёт StatefulSet, PersistentVolumeClaims, Service'ы и управляет lifecycle кластера. При обновлении версии MinIO - rolling update без downtime.
MinIO Client (mc): CLI для администрирования¶
mc - официальный CLI-инструмент для MinIO. Он поддерживает все операции: управление бакетами, объектами, пользователями, политиками, мониторинг.
# Установка mc
curl https://dl.min.io/client/mc/release/linux-amd64/mc \
-o /usr/local/bin/mc
chmod +x /usr/local/bin/mc
# Регистрация MinIO-сервера под псевдонимом "local"
mc alias set local http://localhost:9000 minioadmin minioadmin123
# Проверка подключения
mc admin info local
# Создание бакетов для Medallion Architecture
mc mb local/bronze --region us-east-1
mc mb local/silver --region us-east-1
mc mb local/gold --region us-east-1
# Просмотр структуры
mc ls local/
# Загрузка файла
mc cp ./data.parquet local/bronze/events/2024/01/data.parquet
# Рекурсивный листинг
mc ls --recursive local/bronze/
# Статистика использования
mc du local/bronze/
Безопасность: Bucket Policy и управление доступом¶
Почему нельзя использовать Root Account в Spark-пайплайнах¶
Root Account в MinIO (MINIO_ROOT_USER / MINIO_ROOT_PASSWORD) обладает абсолютными правами: создание/удаление бакетов, управление пользователями, изменение политик, полный доступ ко всем данным. Использование Root Credentials в Spark-джобах - критическая уязвимость:
- Эти credentials видны в Spark UI (вкладка Environment), в логах кластера, в переменных окружения Executor'ов
- Компрометация одного ETL-скрипта = полная компрометация всего хранилища
- Нарушает принцип наименьших привилегий (Principle of Least Privilege)
Правильная модель: создать отдельный Service Account с минимально необходимыми правами для конкретного пайплайна.
Структура Bucket Policy (JSON)¶
MinIO использует IAM-совместимый формат политик доступа - JSON-документ, идентичный AWS IAM Policies. Каждая политика состоит из:
- Version - версия языка политик (всегда
"2012-10-17") - Statement - массив правил доступа
- Effect -
"Allow"или"Deny" - Principal - кому применяется правило (
"*"= всем, или конкретный ARN) - Action - список разрешённых/запрещённых операций S3
- Resource - ARN ресурса (бакет или объекты внутри него)
Создание Service Account и политик через mc¶
# ─────────────────────────────────────────────────────────────
# Сценарий: два пайплайна, два пользователя, два бакета
# spark_ingest → может читать из raw-data, писать в bronze
# spark_transform → может читать bronze, писать silver и gold
# ─────────────────────────────────────────────────────────────
# Создаём бакеты
mc mb local/raw-data
mc mb local/bronze
mc mb local/silver
mc mb local/gold
# ─────────────────────────────────────────────────────────────
# Политика для spark_ingest:
# - ReadOnly на raw-data
# - ReadWrite на bronze
# ─────────────────────────────────────────────────────────────
cat > /tmp/spark-ingest-policy.json << 'EOF'
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:ListBucket",
"s3:GetBucketLocation"
],
"Resource": [
"arn:aws:s3:::raw-data",
"arn:aws:s3:::raw-data/*"
]
},
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:PutObject",
"s3:DeleteObject",
"s3:ListBucket",
"s3:GetBucketLocation",
"s3:AbortMultipartUpload",
"s3:ListMultipartUploadParts"
],
"Resource": [
"arn:aws:s3:::bronze",
"arn:aws:s3:::bronze/*"
]
}
]
}
EOF
# ─────────────────────────────────────────────────────────────
# Политика для spark_transform:
# - ReadOnly на bronze
# - ReadWrite на silver и gold
# ─────────────────────────────────────────────────────────────
cat > /tmp/spark-transform-policy.json << 'EOF'
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:ListBucket",
"s3:GetBucketLocation"
],
"Resource": [
"arn:aws:s3:::bronze",
"arn:aws:s3:::bronze/*"
]
},
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:PutObject",
"s3:DeleteObject",
"s3:ListBucket",
"s3:GetBucketLocation",
"s3:AbortMultipartUpload",
"s3:ListMultipartUploadParts"
],
"Resource": [
"arn:aws:s3:::silver",
"arn:aws:s3:::silver/*",
"arn:aws:s3:::gold",
"arn:aws:s3:::gold/*"
]
}
]
}
EOF
# Регистрируем политики в MinIO
mc admin policy create local spark-ingest-policy /tmp/spark-ingest-policy.json
mc admin policy create local spark-transform-policy /tmp/spark-transform-policy.json
# Создаём пользователей и привязываем политики
mc admin user add local spark_ingest "ingest_secret_key_32chars_minimum"
mc admin user add local spark_transform "transform_secret_key_32chars_min"
mc admin policy attach local spark-ingest-policy --user spark_ingest
mc admin policy attach local spark-transform-policy --user spark_transform
# Проверяем назначенные права
mc admin user info local spark_ingest
mc admin policy info local spark-ingest-policy
Почему DeleteObject обязателен для Spark¶
Многие инженеры удивляются, зачем Spark-воркеру нужно право на удаление объектов. Причин несколько:
1. mode("overwrite") - при перезаписи Parquet-директории Spark сначала записывает файлы во временную директорию (_temporary/), затем переносит их в финальную. При использовании дефолтного FileOutputCommitter «перенос» = copy + delete. Без s3:DeleteObject commit phase упадёт с 403.
2. Magic Committer - использует AbortMultipartUpload при сбоях. Без права на это незавершённые multipart uploads накапливаются и тарифицируются.
3. Delta Lake VACUUM - удаляет старые версии файлов из Delta-таблицы. Требует s3:DeleteObject.
4. Iceberg expire_snapshots - аналогично очищает устаревшие snapshots.
Правило: Spark-воркеру, который пишет данные, нужен полный набор: GetObject, PutObject, DeleteObject, ListBucket, AbortMultipartUpload, ListMultipartUploadParts.
Bucket Policy vs User Policy¶
MinIO поддерживает два типа политик:
User Policy (IAM-style) - привязывается к пользователю или группе. Действует глобально: один пользователь получает одну политику для всего хранилища. Управляется через mc admin policy.
Bucket Policy - привязывается к конкретному бакету. Управляет анонимным доступом (без credentials). Полезно для публичных бакетов или inter-service доступа без явной аутентификации.
# Bucket Policy: анонимный ReadOnly доступ к gold/ (для BI-инструментов)
cat > /tmp/gold-public-read.json << 'EOF'
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Principal": {"AWS": ["*"]},
"Action": ["s3:GetObject", "s3:ListBucket"],
"Resource": [
"arn:aws:s3:::gold",
"arn:aws:s3:::gold/*"
]
}
]
}
EOF
mc anonymous set-json /tmp/gold-public-read.json local/gold
Группы пользователей: RBAC в MinIO¶
При большом количестве пользователей удобнее управлять политиками через группы:
# Создание групп
mc admin group add local spark-writers spark_ingest spark_transform
mc admin group add local spark-readers analytics_user reporting_user
# Привязка политики к группе
mc admin policy attach local spark-transform-policy --group spark-writers
# Все пользователи группы spark-writers автоматически получают права
Интеграция Spark с MinIO: подробная настройка¶
Зависимости: правильная пара JAR¶
Ключевая проблема при интеграции Spark + MinIO - версионная совместимость JAR-файлов. Неправильная пара версий даёт NoSuchMethodError или ClassNotFoundException в runtime.
Матрица совместимости:
| Spark версия | Hadoop версия | hadoop-aws | aws-java-sdk-bundle |
|---|---|---|---|
| Spark 3.3.x | Hadoop 3.3.2 | 3.3.2 |
1.12.262 |
| Spark 3.4.x | Hadoop 3.3.4 | 3.3.4 |
1.12.367 |
| Spark 3.5.x | Hadoop 3.3.6 | 3.3.6 |
1.12.599 |
Проверить версию Hadoop в вашем Spark:
import pyspark
sc = pyspark.SparkContext.getOrCreate()
print(sc._jvm.org.apache.hadoop.util.VersionInfo.getVersion())
# Например: "3.3.4"
Полная конфигурация SparkSession для MinIO¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("MinIO-Integration") \
\
# ── Зависимости ────────────────────────────────────────────
.config(
"spark.jars.packages",
"org.apache.hadoop:hadoop-aws:3.3.4,"
"com.amazonaws:aws-java-sdk-bundle:1.12.367"
) \
\
# ── Подключение к MinIO ────────────────────────────────────
.config("spark.hadoop.fs.s3a.endpoint", "http://localhost:9000") \
.config("spark.hadoop.fs.s3a.access.key", "spark_ingest") \
.config("spark.hadoop.fs.s3a.secret.key", "ingest_secret_key_32chars_minimum") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
.config("spark.hadoop.fs.s3a.impl",
"org.apache.hadoop.fs.s3a.S3AFileSystem") \
.config("spark.hadoop.fs.s3a.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider") \
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false") \
\
# ── Производительность ──────────────────────────────────────
.config("spark.hadoop.fs.s3a.connection.maximum", "200") \
.config("spark.hadoop.fs.s3a.threads.max", "200") \
.config("spark.hadoop.fs.s3a.multipart.size", "67108864") \
.config("spark.hadoop.fs.s3a.multipart.threshold", "67108864") \
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "disk") \
.config("spark.hadoop.fs.s3a.readahead.range", "1048576") \
.config("spark.hadoop.fs.s3a.experimental.fadvise","random") \
\
# ── Magic Committer ────────────────────────────────────────
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.config("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
.config("spark.sql.parquet.output.committer.class",
"org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter") \
\
.getOrCreate()
Разбор ключевых параметров¶
fs.s3a.path.style.access=true - критичный параметр. MinIO не может создать DNS-запись вида bucket.minio.host для каждого бакета (в отличие от AWS S3, который управляет wildcard DNS *.s3.amazonaws.com). Без этого параметра S3A Connector попытается обратиться к bronze.localhost:9000 - и получит DNS resolution failed.
fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem - явная регистрация реализации FileSystem для схемы s3a://. Без этого параметра Spark не знает, какой класс использовать для работы с URI вида s3a://bucket/path. В кластерных дистрибутивах (EMR, CDH) этот класс уже зарегистрирован; в standalone Spark - нужно указать явно.
fs.s3a.aws.credentials.provider=SimpleAWSCredentialsProvider - явное указание провайдера credentials. Без него S3A перебирает цепочку провайдеров: Instance Profile → Environment Variables → ~/.aws/credentials → и только потом Simple (ключи из конфига). В корпоративных средах цепочка может случайно подхватить чужие credentials (например, от другого сервиса), что приведёт к загадочному 403.
fs.s3a.connection.ssl.enabled=false - отключение SSL для HTTP-эндпоинтов MinIO в dev/test. В production при HTTPS-MinIO это должно быть true.
Несколько MinIO-серверов в одной SparkSession¶
Иногда нужно читать из одного MinIO и писать в другой (например, разные environment или разные регионы). S3A поддерживает per-bucket конфигурацию:
spark = SparkSession.builder \
.appName("Multi-MinIO") \
\
# ── MinIO для dev (чтение) ─────────────────────────────────
.config("spark.hadoop.fs.s3a.bucket.raw-data.endpoint",
"http://minio-dev:9000") \
.config("spark.hadoop.fs.s3a.bucket.raw-data.access.key", "reader_key") \
.config("spark.hadoop.fs.s3a.bucket.raw-data.secret.key", "reader_secret") \
.config("spark.hadoop.fs.s3a.bucket.raw-data.path.style.access", "true") \
\
# ── MinIO для prod (запись) ────────────────────────────────
.config("spark.hadoop.fs.s3a.bucket.bronze.endpoint",
"http://minio-prod:9000") \
.config("spark.hadoop.fs.s3a.bucket.bronze.access.key", "writer_key") \
.config("spark.hadoop.fs.s3a.bucket.bronze.secret.key", "writer_secret") \
.config("spark.hadoop.fs.s3a.bucket.bronze.path.style.access", "true") \
\
# ── Общие настройки (fallback) ─────────────────────────────
.config("spark.hadoop.fs.s3a.impl",
"org.apache.hadoop.fs.s3a.S3AFileSystem") \
.getOrCreate()
# Теперь s3a://raw-data/ → minio-dev, s3a://bronze/ → minio-prod
df = spark.read.parquet("s3a://raw-data/events/")
df.write.parquet("s3a://bronze/events/")
Практика: сквозной ETL-пайплайн Bronze → Silver → Gold¶
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType
import time
# ─────────────────────────────────────────────────────────────
# 1. Инициализация SparkSession
# ─────────────────────────────────────────────────────────────
spark = SparkSession.builder \
.appName("Medallion-ETL-MinIO") \
.master("local[4]") \
.config("spark.jars.packages",
"org.apache.hadoop:hadoop-aws:3.3.4,"
"com.amazonaws:aws-java-sdk-bundle:1.12.367") \
.config("spark.hadoop.fs.s3a.endpoint", "http://localhost:9000") \
.config("spark.hadoop.fs.s3a.access.key", "spark_ingest") \
.config("spark.hadoop.fs.s3a.secret.key", "ingest_secret_key_32chars_minimum") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
.config("spark.hadoop.fs.s3a.impl",
"org.apache.hadoop.fs.s3a.S3AFileSystem") \
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false") \
.config("spark.hadoop.fs.s3a.fast.upload", "true") \
.config("spark.hadoop.fs.s3a.fast.upload.buffer", "disk") \
.config("spark.hadoop.fs.s3a.committer.name", "magic") \
.config("spark.hadoop.mapreduce.outputcommitter.factory.scheme.s3a",
"org.apache.hadoop.fs.s3a.commit.S3ACommitterFactory") \
.config("spark.sql.sources.commitProtocolClass",
"org.apache.spark.internal.io.cloud.PathOutputCommitProtocol") \
.config("spark.sql.parquet.output.committer.class",
"org.apache.spark.internal.io.cloud.BindingParquetOutputCommitter") \
.config("spark.sql.session.timeZone", "UTC") \
.getOrCreate()
# ─────────────────────────────────────────────────────────────
# 2. Генерируем синтетические сырые данные → Bronze
# ─────────────────────────────────────────────────────────────
print("=== Шаг 1: Ingestion → Bronze ===")
# Имитируем лог транзакций e-commerce
df_raw = spark.range(0, 1_000_000).select(
F.col("id").alias("transaction_id"),
F.date_sub(
F.current_date(),
(F.col("id") % 90).cast("int")
).alias("event_date"),
F.array(
F.lit("electronics"), F.lit("clothing"),
F.lit("food"), F.lit("books")
).getItem((F.col("id") % 4).cast("int")).alias("category"),
(F.rand() * 10000).alias("amount_raw"), # "сырое" float значение
F.when(F.col("id") % 50 == 0, None) # ~2% битых записей
.otherwise(F.col("id") % 1000).alias("user_id"),
F.lit("USD").alias("currency"),
# Намеренно добавляем мусор (некорректный json-лог)
F.when(F.col("id") % 200 == 0, F.lit("CORRUPTED_RECORD"))
.otherwise(F.concat(F.lit('{"ip":"'), (F.col("id") % 255).cast("string"), F.lit('"}')))
.alias("raw_log"),
)
start = time.time()
df_raw.coalesce(8).write \
.mode("overwrite") \
.partitionBy("event_date") \
.parquet("s3a://bronze/transactions/")
print(f"Bronze записан за {time.time() - start:.1f}s, "
f"{df_raw.rdd.getNumPartitions()} партиций")
# ─────────────────────────────────────────────────────────────
# 3. Чтение Bronze → очистка → Silver
# ─────────────────────────────────────────────────────────────
print("\n=== Шаг 2: Bronze → Silver (очистка) ===")
df_bronze = spark.read.parquet("s3a://bronze/transactions/")
print(f"Bronze: {df_bronze.count()} записей")
# Partition Pruning: читаем только последние 30 дней
from datetime import date, timedelta
cutoff = (date.today() - timedelta(days=30)).isoformat()
df_silver = (
df_bronze
# Отсекаем по дате (partition pruning на MinIO!)
.filter(F.col("event_date") >= cutoff)
# Убираем битые записи (NULL user_id)
.filter(F.col("user_id").isNotNull())
# Убираем CORRUPTED_RECORD из raw_log
.filter(~F.col("raw_log").startswith("CORRUPTED"))
# Нормализуем amount: округляем до 2 знаков
.withColumn("amount", F.round(F.col("amount_raw"), 2))
.drop("amount_raw")
# Добавляем технические колонки
.withColumn("processed_at", F.current_timestamp())
.withColumn("event_year", F.year("event_date"))
.withColumn("event_month", F.month("event_date"))
)
print(f"Silver (после фильтрации): {df_silver.count()} записей")
start = time.time()
df_silver.write \
.mode("overwrite") \
.partitionBy("event_year", "event_month") \
.parquet("s3a://silver/transactions/")
print(f"Silver записан за {time.time() - start:.1f}s")
# ─────────────────────────────────────────────────────────────
# 4. Агрегация Silver → Gold
# ─────────────────────────────────────────────────────────────
print("\n=== Шаг 3: Silver → Gold (агрегация) ===")
df_silver_read = spark.read.parquet("s3a://silver/transactions/")
df_gold = (
df_silver_read
.groupBy("event_year", "event_month", "category")
.agg(
F.count("*").alias("transaction_count"),
F.sum("amount").alias("total_revenue"),
F.avg("amount").alias("avg_transaction"),
F.countDistinct("user_id").alias("unique_users"),
)
.withColumn("revenue_per_user",
F.col("total_revenue") / F.col("unique_users"))
)
start = time.time()
df_gold.write \
.mode("overwrite") \
.partitionBy("event_year", "event_month") \
.parquet("s3a://gold/transactions_summary/")
print(f"Gold записан за {time.time() - start:.1f}s")
# ─────────────────────────────────────────────────────────────
# 5. Верификация результата
# ─────────────────────────────────────────────────────────────
df_result = spark.read.parquet("s3a://gold/transactions_summary/")
print(f"\nGold: {df_result.count()} агрегатных строк")
df_result.orderBy("total_revenue", ascending=False).show(5)
Проверка через mc¶
# Проверяем что данные реально записались в MinIO
mc ls local/bronze/transactions/ --recursive | head -20
mc du local/bronze/
mc du local/silver/
mc du local/gold/
# Проверяем структуру партиционирования
mc ls local/silver/transactions/
# output:
# [2024-01-15] DIR event_year=2024/
# ...
mc ls local/silver/transactions/event_year=2024/
# output:
# [2024-01-15] DIR event_month=1/
# [2024-01-15] DIR event_month=2/
# ...
Delta Lake поверх MinIO¶
Delta Lake и MinIO - идеальное сочетание для Open-Source Data Lakehouse. Delta Lake добавляет ACID-транзакции, schema evolution и time travel поверх Parquet-файлов в MinIO, полностью решая проблему rename (Delta Lake записывает файлы напрямую с UUID-именами, без _temporary/).
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from delta import configure_spark_with_delta_pip, DeltaTable
spark = configure_spark_with_delta_pip(
SparkSession.builder
.appName("Delta-on-MinIO")
.master("local[4]")
.config("spark.sql.extensions",
"io.delta.sql.DeltaSparkSessionExtension")
.config("spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog")
# MinIO конфигурация
.config("spark.hadoop.fs.s3a.endpoint", "http://localhost:9000")
.config("spark.hadoop.fs.s3a.access.key", "spark_ingest")
.config("spark.hadoop.fs.s3a.secret.key", "ingest_secret_key_32chars_minimum")
.config("spark.hadoop.fs.s3a.path.style.access", "true")
.config("spark.hadoop.fs.s3a.impl",
"org.apache.hadoop.fs.s3a.S3AFileSystem")
.config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false")
).getOrCreate()
DELTA_PATH = "s3a://bronze/delta/transactions/"
# ── Первоначальная запись ──────────────────────────────────────
df = spark.range(0, 100_000).select(
F.col("id").alias("transaction_id"),
(F.rand() * 1000).alias("amount"),
F.date_sub(F.current_date(), (F.col("id") % 30).cast("int")).alias("event_date"),
F.array(F.lit("A"), F.lit("B"), F.lit("C"))
.getItem((F.col("id") % 3).cast("int")).alias("category"),
)
df.write \
.format("delta") \
.mode("overwrite") \
.partitionBy("event_date") \
.save(DELTA_PATH)
print("Delta table создана")
# ── Инкрементальное добавление (APPEND) ────────────────────────
df_new = spark.range(100_000, 110_000).select(
F.col("id").alias("transaction_id"),
(F.rand() * 1000).alias("amount"),
F.current_date().alias("event_date"),
F.lit("D").alias("category"),
)
df_new.write \
.format("delta") \
.mode("append") \
.save(DELTA_PATH)
# ── MERGE (Upsert) - ACID-транзакция ──────────────────────────
updates = spark.createDataFrame([
(1, 9999.0, "UPDATED"),
(2, 8888.0, "UPDATED"),
(999_999, 7777.0, "NEW"), # новая строка
], ["transaction_id", "amount", "category"])
delta_table = DeltaTable.forPath(spark, DELTA_PATH)
delta_table.alias("target").merge(
updates.alias("source"),
"target.transaction_id = source.transaction_id"
).whenMatchedUpdate(set={
"amount": "source.amount",
"category": "source.category",
}).whenNotMatchedInsert(values={
"transaction_id": "source.transaction_id",
"amount": "source.amount",
"category": "source.category",
"event_date": F.current_date(),
}).execute()
# ── Time Travel: чтение исторической версии ────────────────────
df_v0 = spark.read.format("delta") \
.option("versionAsOf", 0) \
.load(DELTA_PATH)
print(f"Версия 0 (до merge): {df_v0.count()} строк")
df_current = spark.read.format("delta").load(DELTA_PATH)
print(f"Текущая версия: {df_current.count()} строк")
# ── История версий ─────────────────────────────────────────────
delta_table.history().select(
"version", "timestamp", "operation", "operationParameters"
).show(truncate=False)
# ── Очистка старых версий (VACUUM) ────────────────────────────
spark.conf.set("spark.databricks.delta.retentionDurationCheck.enabled", "false")
delta_table.vacuum(retentionHours=0) # в production: retentionHours=168 (7 дней)
# Смотрим что получилось в MinIO
# mc ls --recursive local/bronze/delta/transactions/_delta_log/
# output:
# [2024-01-15] 1.2 KiB _delta_log/00000000000000000000.json (initial write)
# [2024-01-15] 0.9 KiB _delta_log/00000000000000000001.json (append)
# [2024-01-15] 1.8 KiB _delta_log/00000000000000000002.json (merge)
# [2024-01-15] 0.5 KiB _delta_log/00000000000000000003.json (vacuum)
Почему Delta Lake не нужны специальные коммитеры: Delta Lake записывает каждый Parquet-файл с уникальным UUID-именем напрямую в финальную директорию. Транзакционная атомарность обеспечивается через _delta_log/ - JSON-файлы транзакций. Если запись прервалась - файлы данных могут существовать, но без записи в _delta_log они не видимы ни одному читателю. При следующем чтении/записи Delta автоматически очищает «осиротевшие» файлы.
Monitoring MinIO: Prometheus + Grafana¶
MinIO нативно экспортирует метрики в формате Prometheus через endpoint /minio/v2/metrics/cluster.
Включение метрик¶
# Генерация prometheus scrape token (Bearer token для auth)
mc admin prometheus generate local
# Пример вывода:
# scrape_configs:
# - job_name: minio-job
# bearer_token: eyJhbGciOiJIUzUxMiIsInR5cCI6IkpXVCJ9...
# metrics_path: /minio/v2/metrics/cluster
# scheme: http
# static_configs:
# - targets: ['localhost:9000']
Ключевые метрики для Data Engineer¶
# Использование пространства
minio_bucket_usage_total_bytes{bucket="bronze"} → объём данных в бакете
minio_bucket_objects_count{bucket="bronze"} → количество объектов
minio_cluster_capacity_usable_free_bytes → свободное место
# I/O запросы
minio_s3_requests_total{api="PutObject"} → количество PUT запросов
minio_s3_requests_total{api="GetObject"} → количество GET запросов
minio_s3_requests_errors_total{api="PutObject"} → ошибки
# Latency
minio_s3_requests_ttfb_seconds_distribution{api="GetObject", le="0.1"} → p100 < 100ms
# Erasure Coding health
minio_cluster_nodes_online_total → живые ноды
minio_cluster_drive_online_total → живые диски
minio_cluster_drive_offline_total → упавшие диски (должно быть 0)
Prometheus + Grafana через Docker Compose¶
# monitoring/docker-compose.yml
version: "3.8"
services:
prometheus:
image: prom/prometheus:v2.47.0
volumes:
- ./prometheus.yml:/etc/prometheus/prometheus.yml
ports:
- "9090:9090"
grafana:
image: grafana/grafana:10.1.0
ports:
- "3000:3000"
environment:
GF_SECURITY_ADMIN_PASSWORD: admin
volumes:
- grafana_data:/var/lib/grafana
volumes:
grafana_data:
# monitoring/prometheus.yml
global:
scrape_interval: 15s
scrape_configs:
- job_name: minio
bearer_token: "YOUR_BEARER_TOKEN_FROM_MC_ADMIN_PROMETHEUS"
metrics_path: /minio/v2/metrics/cluster
scheme: http
static_configs:
- targets: ["minio:9000"]
В Grafana импортировать официальный MinIO Dashboard (ID: 13502) из Grafana Dashboard Catalog.
Диагностика типичных ошибок Spark + MinIO¶
403 Forbidden: Access Denied¶
Самая частая ошибка. Причин несколько:
com.amazonaws.services.s3.model.AmazonS3Exception:
Access Denied (Service: Amazon S3; Status Code: 403; Error Code: AccessDenied)
Диагностика:
# 1. Проверяем права пользователя
mc admin user info local spark_ingest
# 2. Проверяем политику
mc admin policy info local spark-ingest-policy
# 3. Тестируем конкретную операцию
mc cp /tmp/test.txt local/bronze/test.txt # от имени spark_ingest
# 4. Частые причины:
# - Нет s3:ListBucket → Spark не может составить план чтения
# - Нет s3:DeleteObject → mode("overwrite") падает при commit
# - Нет s3:AbortMultipartUpload → Magic Committer не может rollback
# - Resource ARN неверный: "arn:aws:s3:::bronze" (бакет) vs "arn:aws:s3:::bronze/*" (объекты)
Важная деталь ARN: в Bucket Policy нужно указывать ДВА ресурса: бакет (для ListBucket) и объекты внутри него (для GetObject, PutObject и т.д.):
"Resource": [
"arn:aws:s3:::bronze", ← для ListBucket, GetBucketLocation
"arn:aws:s3:::bronze/*" ← для GetObject, PutObject, DeleteObject
]
Если указать только "arn:aws:s3:::bronze/*" - ListBucket упадёт с 403.
UnknownHostException: DNS не резолвит имя бакета¶
java.net.UnknownHostException: bronze.localhost
Причина: не установлен path.style.access=true. S3A пытается обратиться к bucket.endpoint вместо endpoint/bucket.
Решение:
.config("spark.hadoop.fs.s3a.path.style.access", "true")
SignatureDoesNotMatch¶
com.amazonaws.services.s3.model.AmazonS3Exception:
The request signature we calculated does not match the signature you provided
(Status Code: 403; Error Code: SignatureDoesNotMatch)
Причины:
- Неверный secret key (опечатка, скопировали с пробелом)
- Системное время сервера Spark и сервера MinIO расходятся более чем на 5 минут (AWS Signature V4 включает timestamp)
- Устаревший алгоритм подписи (для очень старых версий MinIO может потребоваться
S3SignerType)
Диагностика:
# Проверить синхронизацию времени
date && mc admin info local | grep -i time
# Проверить credentials вручную
mc ls --debug local/bronze/ 2>&1 | grep "Signature\|Authorization"
Connection refused / Connection reset¶
java.net.ConnectException: Connection refused to http://localhost:9000
Причина: MinIO не запущен, неверный endpoint, порт заблокирован.
Диагностика:
# Проверить что MinIO слушает
curl http://localhost:9000/minio/health/live
# Ожидаемый ответ: HTTP 200
# Проверить network доступность из контейнера Spark
docker exec spark-container curl http://minio:9000/minio/health/live
ClassNotFoundException: S3AFileSystem¶
java.lang.ClassNotFoundException: org.apache.hadoop.fs.s3a.S3AFileSystem
Причина: hadoop-aws JAR не в classpath.
Решение: добавить --packages org.apache.hadoop:hadoop-aws:3.3.4,... или в Docker-образ.
Дизайн структуры бакетов для Data Lakehouse¶
Правильная организация бакетов и путей в MinIO критически важна для partition pruning, управления доступом и операционного удобства.
Medallion Architecture в MinIO¶
bronze/
├── transactions/
│ ├── event_date=2024-01-01/
│ │ ├── part-00001-uuid.parquet
│ │ └── part-00002-uuid.parquet
│ └── event_date=2024-01-02/
├── events/
│ └── event_date=2024-01-01/
└── users/
silver/
├── transactions/
│ ├── event_year=2024/
│ │ ├── event_month=1/
│ │ │ └── part-00001-uuid.parquet
│ │ └── event_month=2/
│ └── event_year=2023/
└── users_enriched/
gold/
├── transactions_summary/
│ └── event_year=2024/event_month=1/
├── user_cohorts/
└── product_metrics/
Принципы проектирования:
- Отдельный бакет на слой (bronze/silver/gold), а не один бакет с префиксами. Это позволяет применять разные Bucket Policy: bronze - только
spark_ingestпишет, silver/gold - толькоspark_transform. - Партиционирование по дате для partition pruning:
event_date=YYYY-MM-DDна Bronze,event_year=YYYY/event_month=MMна Silver/Gold. - UUID в именах файлов - важно для S3-совместимых хранилищ: равномерное распределение нагрузки по prefix shards.
Почему НЕ стоит партиционировать слишком глубоко¶
Частая ошибка - чрезмерное партиционирование: year/month/day/hour/region/category/. При 4 уровнях иерархии и каждом уровне по 10 значений = 10 000 «директорий». Каждая из них - отдельный LIST-запрос к MinIO при планировании запроса. При миллионе партиций планирование занимает минуты.
Правило: для Batch-аналитики оптимально 1–2 уровня партиционирования. Для streaming (данные каждые 5 минут) - партиционирование по часу максимум, и обязательная компакция.
Антипаттерны MinIO + Spark¶
1. Root Credentials в production-пайплайнах. Root Account = full access. Компрометация credentials ETL-скрипта = доступ ко всем данным. Всегда создавайте Service Account с минимальными правами.
2. Публичные бакеты (anonymous read/write). Даже read-only публичный бакет - это утечка данных в открытый интернет. Для BI-инструментов внутри сети используйте ограниченные Service Accounts, а не публичный доступ.
3. Hardcoded credentials в SparkSession.builder. Эти конфиги видны в Spark UI (Environment tab) и в логах кластера. Используйте переменные окружения или Kubernetes Secrets:
import os
spark = SparkSession.builder \
.config("spark.hadoop.fs.s3a.access.key",
os.environ["MINIO_ACCESS_KEY"]) \
.config("spark.hadoop.fs.s3a.secret.key",
os.environ["MINIO_SECRET_KEY"]) \
.getOrCreate()
4. Миллионы мелких файлов. Каждый файл = объект в MinIO = запись в metadata service = отдельный HTTP-запрос при листинге. При 1 млн файлов в директории - listStatus() делает 1000 LIST-запросов к MinIO. Spark сидит и ждёт планирования. Всегда компактируйте: целевой размер файла 128 MB–1 GB.
5. Один бакет для всех данных. s3a://datalake/bronze/, s3a://datalake/silver/, s3a://datalake/gold/ в одном бакете означает, что нельзя применить разные Bucket Policy к разным слоям. Spark-воркер, пишущий в bronze, получает доступ ко всему бакету.
6. Не настроен Magic Committer. Дефолтный FileOutputCommitter использует rename (copy + delete). На MinIO это O(N) HTTP-запросов только для финализации записи. При 1000 файлов - минуты только на commit phase.
Итоговый чек-лист¶
- MinIO Standalone - только для dev/test; prod всегда Distributed с Erasure Coding
- Никогда не использовать Root Credentials в Spark-пайплайнах - создавать Service Account с минимальными правами
- Bucket Policy: для записи нужны
PutObject+DeleteObject+AbortMultipartUpload+ListBucket - В
ResourceARN указывать ДВА ресурса: бакет (arn:aws:s3:::bucket) И объекты (arn:aws:s3:::bucket/*) path.style.access=trueобязателен для MinIO - без него DNS ломаетсяimpl=org.apache.hadoop.fs.s3a.S3AFileSystemуказывать явно в standalone Spark- Magic Committer устраняет rename-проблему - включать всегда для Parquet без Delta/Iceberg
- Delta Lake и Iceberg не требуют специальных коммитеров - у них собственный commit protocol
- Отдельные бакеты на каждый слой Medallion (bronze/silver/gold) - для изоляции политик доступа
- Партиционирование: 1–2 уровня максимум для batch-аналитики
- Файлы 128 MB–1 GB - оптимальный размер; мелкие файлы убивают performance листинга в MinIO
- Мониторинг: Prometheus + Grafana Dashboard 13502 - следить за
drive_offline_totalи request errors