MinIO: установка, настройка bucket policy и интеграция со Spark

Глубокий разбор MinIO как фундамента Open-Source Data Lakehouse: архитектура, Erasure Coding, Bucket Policy и IAM, полная интеграция со Spark через S3A, Delta Lake и Iceberg поверх MinIO.

storage

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
  • В Resource ARN указывать ДВА ресурса: бакет (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