Erasure Coding: Reed-Solomon коды, overhead vs репликация, hot/cold данные

Полный разбор Erasure Coding в HDFS 3.x: математика Reed-Solomon, сравнение с репликацией 3x по дискам и CPU, разрушение Data Locality для Spark, аппаратное ускорение Intel ISA-L, архитектурный паттерн Hot/Cold Tiering, CLI настройка политик и сравнительный бенчмарк Spark.

storage platform

1. Экономический и физический тупик классической репликации 3x

Hadoop изначально строился вокруг одной простой идеи обеспечения надёжности: хранить три копии каждого блока данных на разных физических серверах. Эта схема работала превосходно - при потере любого сервера или даже двух серверов одновременно данные оставались доступны. Но у этой простоты была цена, которая с ростом объёмов данных стала экономически неприемлемой.

Анатомия overhead'а репликации 3x

Рассмотрим конкретные числа. Крупная финансовая компания хранит в Data Lake транзакционные данные за 10 лет: 2 PB полезных данных. При репликации 3x это означает:

  • Дисковое пространство: 2 PB × 3 = 6 PB физической ёмкости
  • Сетевой трафик при записи: каждый записанный байт передаётся трижды
  • Стоимость серверов: серверы с дисками ёмкостью ~300 TB каждый → нужно 20 серверов только для хранения

Из 6 PB физического пространства 4 PB - это pure overhead. Деньги на серверы, электричество, охлаждение, замену дисков уходят не на полезные данные, а на поддержание двух лишних копий.

Жизненный цикл данных: почему репликация 3x неэффективна для «холодных» слоёв

Ключевое наблюдение: не все данные в Data Lake одинаково важны для быстрого доступа. Данные проходят через жизненный цикл:

Схема жизненного цикла показывает, что большая часть данных - COLD и FROZEN - читается редко, но занимает большую часть физического пространства. Именно для этих слоёв репликация 3x расточительна: зачем хранить три копии данных, которые читаются раз в квартал?

На практике в крупном Data Lake распределение данных выглядит примерно так:

Слой Доля объёма Частота доступа Оптимальная стратегия
HOT (< 1 года) 15% Ежедневно Репликация 3x
WARM (1-2 года) 20% Еженедельно Репликация 2x или EC
COLD (2-5 лет) 35% Ежеквартально EC RS-6-3
FROZEN (5+ лет) 30% Ежегодно EC RS-10-4 или Glacier

Только на переводе 65% объёма данных с репликации 3x на EC можно сэкономить около 43% общей дисковой ёмкости кластера.


2. Математика и физика Erasure Coding: коды Рида-Соломона

Erasure Coding - это математический способ обеспечить избыточность хранения без полного дублирования данных. Вместо «храни три одинаковые копии» идея такая: «преобразуй данные математически так, чтобы любые K из N блоков позволяли восстановить исходное сообщение».

Интуиция за кодами Рида-Соломона

Представим простой пример. Возьмём два числа: a = 3 и b = 5. Это «данные». Теперь вычислим «контрольный блок»: c = a + b = 8. Храним три значения: 3, 5, 8.

Если потеряли a (3): восстанавливаем a = c - b = 8 - 5 = 3. Если потеряли b (5): b = c - a = 8 - 3 = 5. Если потеряли c (8): c = a + b = 3 + 5 = 8.

Мы можем восстановить любое одно из трёх значений. При этом «overhead» составляет 33% (одно пarity-число на два data-числа). Реальные коды Рида-Соломона работают с полями Галуа (GF(2^8)) вместо обычной арифметики, но математическая интуиция та же: из любых K блоков можно восстановить все исходные данные.

Схема RS-k-m: параметры и их смысл

В HDFS Erasure Coding определяется схемой RS-k-m:

  • k - количество блоков данных (Data Blocks)
  • m - количество блоков чётности (Parity Blocks)
  • Всего блоков: k + m
  • Максимальное число допустимых потерь: m (любые m блоков из k+m)

На схеме видно, как 600 MB файл делится на 6 блоков по 100 MB, затем кодировщик Reed-Solomon вычисляет 3 parity-блока. Итого хранится 900 MB - при том что полезных данных 600 MB. Overhead - ровно 50%, против 200% у репликации 3x.

Популярные схемы EC в HDFS и их характеристики

HDFS 3.x поставляется со следующими предустановленными схемами:

Схема Data Parity Всего Overhead Допустимые потери Мин. DataNode Применение
RS-3-2-1024k 3 2 5 67% 2 5 Маленькие кластеры
RS-6-3-1024k 6 3 9 50% 3 9 Стандарт для Cold
RS-10-4-1024k 10 4 14 40% 4 14 Крупные кластеры, архив
RS-LEGACY-6-3-1024k 6 3 9 50% 3 9 Совместимость с Hadoop 2.x
XOR-2-1-1024k 2 1 3 50% 1 3 Простое XOR, минимальный CPU

Число 1024k в названии - это размер ячейки (cell size): каждый блок данных нарезается на ячейки по 1024 KB (1 MB). Ячейки чередуются между DataNode (striping) - подробнее об этом ниже.

Striping: как данные нарезаются на ячейки

В отличие от классической репликации (где блок целиком лежит на одном диске), EC использует striping - чередование ячеек между несколькими DataNode. Это принципиально меняет физическую структуру хранения:

Эта разница в физической структуре хранения - ключ к пониманию как преимуществ EC (50% overhead), так и его главного недостатка (потеря Data Locality для Spark). Схема наглядно показывает: при репликации каждый блок полностью находится на одном сервере, и Spark может читать его локально. При EC каждый логический блок «размазан» по 9 серверам, и для его чтения нужно собрать ячейки со всех 9.


3. Великий удар по производительности Spark: разрушение Data Locality

Это самый важный раздел урока с точки зрения практики Data Engineering. Erasure Coding - это не просто «другой способ хранения», это фундаментальное изменение того, как Spark читает данные, и если этот факт игнорировать, производительность ETL-пайплайнов на EC-данных будет значительно хуже ожидаемой.

Почему NODE_LOCAL невозможен для EC-данных

Вспомним как работает Data Locality. При чтении HDFS-файла Spark запрашивает у NameNode расположение блоков, затем пытается запустить задачи на тех же серверах, где лежат блоки (NODE_LOCAL уровень).

При репликации 3x это работает отлично: блок 128 MB лежит целиком на Server1 (и ещё на Server2 и Server3 как реплики). Spark запускает задачу на Server1 - данные читаются с локального диска, сеть не нагружается.

При EC RS-6-3 один логический блок файла разбит на 6 ячеек данных по 1 MB каждая, которые распределены по 6 разным серверам. Чтобы прочитать один логический блок, Executor должен:

  1. Запросить cell1 с DN1 по сети
  2. Запросить cell2 с DN2 по сети
  3. Запросить cell3 с DN3 по сети
  4. Запросить cell4 с DN4 по сети
  5. Запросить cell5 с DN5 по сети
  6. Запросить cell6 с DN6 по сети
  7. Собрать все 6 ячеек вместе в памяти

Независимо от того, на каком сервере запущен Executor, минимум 5 из 6 ячеек будут читаться по сети. NODE_LOCAL физически невозможен. Максимально достижимый уровень локальности - RACK_LOCAL, и то только если все 6 Data-серверов находятся в одной rack (что само по себе нарушало бы отказоустойчивость).

Сетевой шторм: ALL-to-ALL I/O при чтении EC-данных

Рассмотрим что происходит когда Spark Stage с 200 задачами читает EC-данные:

Диаграмма показывает, что один Executor делает 5 сетевых запросов только для чтения одного логического блока. При 200 параллельных задачах это превращается в 1000 одновременных сетевых соединений только для одного Stage. Это и есть «сетевой шторм» - резкий рост межсерверного трафика даже при чтении без Shuffle.

На практике это выражается в:

  • Рост сетевого трафика в 5–6 раз по сравнению с репликацией при той же аналитической задаче
  • Увеличение времени выполнения Stage'а в 1.5–3 раза на реальных нагрузках
  • Насыщение сетевых интерфейсов DataNode при параллельном чтении большого числа Executor'ов

Деградация при восстановлении: Reconstruction Penalty

Самый тяжёлый сценарий - чтение EC-данных когда один из DataNode недоступен. В этом случае HDFS Client Executor'а не может получить одну из ячеек (например, DN3 упал и cell3 недоступна). Тогда запускается EC Reconstruction:

  1. Executor читает все доступные data-ячейки: cell1, cell2, cell4, cell5, cell6 (5 из 6)
  2. Executor читает один из parity-блоков: p1 (с DN7)
  3. Декодер Reed-Solomon вычисляет cell3: cell3 = RS_decode(cell1, cell2, cell4, cell5, cell6, p1)
  4. Только теперь Executor имеет полный набор данных

Этот процесс требует:

  • Дополнительного сетевого запроса (parity-блок с DN7)
  • CPU-интенсивных вычислений Reed-Solomon decode (умножение матриц в GF(2^8))
  • Дополнительной памяти для хранения всех ячеек до завершения декодирования

В условиях активного кластера где DataNode'ы временно перегружены или сеть имеет пакетные потери - Reconstruction Penalty превращается в регулярное явление, которое суммарно может замедлить Spark-джобу в 2–4 раза.

# Эффект EC на Spark-задачи можно увидеть в метриках:
# При нормальной работе (все DN доступны):
# - Task Duration: ~30 сек
# - Input Bytes Read: 128 MB (один блок)
# - Bytes Read From Remote: 107 MB (5 из 6 ячеек по сети)

# При деградации (1 DN недоступен, идёт reconstruction):
# - Task Duration: ~90 сек (в 3 раза дольше!)
# - Input Bytes Read: 128 MB (тот же объём данных)
# - Bytes Read From Remote: 128 MB (все 6 ячеек по сети + parity)
# - Extra CPU Time: +15 сек (RS decode)

4. Оптимизация вычислений: Intel ISA-L и аппаратное ускорение

Математика Reed-Solomon красива теоретически, но требовательна к CPU практически. Кодирование и декодирование матриц в поле Галуа GF(2^8) - это интенсивные операции XOR и умножения, которые при выполнении на стандартном Java-коде потребляют значительные ресурсы процессора.

Почему чистый Java RS-decode - это проблема

Стандартная Java-реализация RS-кодирования в Hadoop выполняет операции с байтами в цикле. Для блока 128 MB при схеме RS-6-3 это означает:

  • Матрица кодирования 3×6 (parity_count × data_count)
  • Для каждого байта из 134 217 728 байт (128 MB): 18 операций умножения и сложения в GF(2^8)
  • Итого: ~2.4 миллиарда операций для кодирования одного блока

На современном CPU без векторизации (scalar mode): ~200 ms на блок. На кластере с сотнями одновременных операций кодирования/декодирования CPU становится узким местом.

Intel ISA-L: векторизованное вычисление через SIMD

Intel Storage Acceleration Library (ISA-L) - это открытая C++ библиотека, реализующая алгоритмы Erasure Coding с использованием SIMD инструкций процессора:

  • SSE4.2: обрабатывает 16 байт за одну инструкцию
  • AVX2: обрабатывает 32 байта за одну инструкцию
  • AVX-512: обрабатывает 64 байта за одну инструкцию

Практический эффект:

Реализация Скорость кодирования (1 core) Ускорение
Java (scalar) ~500 MB/s 1x
C без SIMD ~1.2 GB/s 2.4x
С SSE4.2 ~4 GB/s 8x
С AVX2 ~8 GB/s 16x
С AVX-512 ~16 GB/s 32x

При AVX-512 кодирование блока 128 MB занимает ~8 ms вместо ~200 ms на Java. Это меняет EC из «CPU-bottleneck» в «network-bottleneck» операцию.

Как ISA-L интегрируется с Hadoop

ISA-L интегрируется через ту же нативную библиотеку libhadoop.so, о которой мы говорили в уроке про Short-Circuit Reads. Если Hadoop собран с поддержкой ISA-L, libhadoop.so автоматически использует векторизованные реализации при доступности соответствующих инструкций CPU.

Схема показывает путь вызова: если libhadoop.so доступна, Spark использует нативные ISA-L реализации с автоматическим выбором наиболее мощных инструкций для конкретного CPU. Если библиотека недоступна - fallback на медленную Java-реализацию.

Проверка и активация ISA-L

# 1. Проверить поддержку AVX-512 на процессоре
grep -m1 "avx512" /proc/cpuinfo
# Если вывод есть - AVX-512 поддерживается

# 2. Проверить наличие ISA-L библиотеки в Hadoop
find /usr/lib/hadoop /opt/hadoop -name "libhadoop.so*" 2>/dev/null

# 3. Проверить что ISA-L скомпилирован в libhadoop.so
nm -D /usr/lib/hadoop/lib/native/libhadoop.so | grep -i "isal\|isa_l"
# Должны быть символы: erasure_code_cal_matrix, gf_gen_rs_matrix и т.д.

# 4. Проверить через Hadoop CLI что EC использует нативный codec
hadoop checknative -a
# Ожидаемый вывод:
# native-hadoop: true /usr/lib/hadoop/lib/native/libhadoop.so.1.0.0
# isa-l: true /usr/lib/hadoop/lib/native/libhadoop.so.1.0.0
# java zlib: false (не нужно для EC)

# 5. Если ISA-L не встроен - установить и пересобрать Hadoop
# Ubuntu/Debian:
sudo apt-get install libisal-dev
# CentOS/RHEL:
sudo yum install isa-l-devel

# Пересборка Hadoop с ISA-L поддержкой:
mvn package -Pdist,native -DskipTests \
    -Disal.prefix=/usr \
    -Disal.lib=/usr/lib64 2>&1 | tail -20

Передача настроек ISA-L в SparkSession (для Executor'ов):

spark = SparkSession.builder \
    # Путь к нативным библиотекам Hadoop (включая libhadoop.so с ISA-L)
    .config("spark.executor.extraJavaOptions",
            "-Djava.library.path=/usr/lib/hadoop/lib/native") \
    .config("spark.executorEnv.LD_LIBRARY_PATH",
            "/usr/lib/hadoop/lib/native") \
    # Принудительно использовать нативный EC codec
    .config("spark.hadoop.io.erasurecode.codec.rs.rawcoder",
            "org.apache.hadoop.io.erasurecode.rawcoder.NativeRSRawErasureCoderFactory") \
    .getOrCreate()

5. Архитектурный паттерн: Hot/Cold Data Tiering

Зная про overhead EC и его влияние на Spark, правильная стратегия - не выбирать между «EC везде» или «репликация везде», а применять каждый метод там, где он даёт максимальный результат.

Проектирование зон Data Lakehouse

Схема показывает практическую организацию Data Lakehouse: горячие данные остаются на репликации для максимальной Spark-производительности, холодные данные переходят на EC для экономии дискового пространства. Переходы автоматизируются через Airflow.

Настройка EC Policy через HDFS CLI

HDFS 3.x предоставляет удобный CLI для управления EC политиками:

# 1. Просмотр доступных EC политик
hdfs ec -listPolicies

# Пример вывода:
# ErasureCodingPolicy=[Name=RS-6-3-1024k, Schema=[ECSchema=[Codec=rs, numDataUnits=6,
#   numParityUnits=3]], CellSize=1048576, State=ENABLED]
# ErasureCodingPolicy=[Name=RS-3-2-1024k, Schema=[ECSchema=[Codec=rs, numDataUnits=3,
#   numParityUnits=2]], CellSize=1048576, State=ENABLED]
# ErasureCodingPolicy=[Name=RS-10-4-1024k, Schema=[ECSchema=[Codec=rs, numDataUnits=10,
#   numParityUnits=4]], CellSize=1048576, State=ENABLED]

# 2. Просмотр текущей политики директории
hdfs ec -getPolicy -path /user/spark/data/bronze/current/
# Directory: /user/spark/data/bronze/current/
# ErasureCodingPolicy=null (нет EC, использует репликацию)

hdfs ec -getPolicy -path /user/spark/data/bronze/archive_2y/
# ErasureCodingPolicy=RS-6-3-1024k

# 3. Установка EC политики на директорию
# ВАЖНО: политика применяется к НОВЫМ файлам!
# Уже существующие файлы не конвертируются автоматически!
hdfs ec -setPolicy -path /user/spark/data/bronze/archive_2y -policy RS-6-3-1024k

# 4. Удаление EC политики (возврат к репликации для новых файлов)
hdfs ec -unsetPolicy -path /user/spark/data/bronze/archive_2y

# 5. Конвертация СУЩЕСТВУЮЩИХ реплицированных файлов в EC
# hdfs distcp с автоматическим применением EC политики целевой директории:
hdfs distcp \
  -Ddfs.replication=1 \
  hdfs://cluster/user/spark/data/bronze/2022/ \
  hdfs://cluster/user/spark/data/cold/bronze/2022/
# Целевая директория должна иметь установленную EC политику заранее!

# 6. Проверка статуса блоков конкретного файла
hdfs fsck /user/spark/data/cold/bronze/2022/part-001.parquet \
  -files -blocks -locations
# Вывод покажет EC блоки и их расположение по DataNode

Автоматизация через Airflow: DAG перевода холодных данных в EC

# dags/hdfs_ec_migration_dag.py
# Airflow DAG: автоматический перевод данных старше N месяцев в EC

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.apache.hdfs.hooks.hdfs import HDFSHook
import subprocess
import logging


logger = logging.getLogger(__name__)

DEFAULT_ARGS = {
    "owner": "data-platform",
    "depends_on_past": False,
    "retries": 2,
    "retry_delay": timedelta(minutes=30),
    "email_on_failure": True,
}

# Конфигурация: какие директории переводить в EC и когда
EC_MIGRATION_CONFIG = [
    {
        "source_base": "/user/spark/data/bronze",
        "cold_base": "/user/spark/data/cold/bronze",
        "ec_policy": "RS-6-3-1024k",
        "age_months": 24,        # данные старше 24 месяцев → EC
        "description": "Bronze raw layer archival",
    },
    {
        "source_base": "/user/spark/data/silver",
        "cold_base": "/user/spark/data/cold/silver",
        "ec_policy": "RS-6-3-1024k",
        "age_months": 18,        # Silver - переводим чуть быстрее
        "description": "Silver cleaned layer archival",
    },
]


def get_partitions_older_than(hdfs_path: str, months: int) -> list[str]:
    """
    Возвращает список партиций (поддиректорий date=YYYY-MM-DD)
    старше указанного числа месяцев.
    Предполагает партиционирование по дате вида /path/date=2022-01-15/.
    """
    cutoff_date = datetime.now() - timedelta(days=months * 30)

    result = subprocess.run(
        ["hdfs", "dfs", "-ls", f"{hdfs_path}/"],
        capture_output=True, text=True
    )

    old_partitions = []
    for line in result.stdout.splitlines():
        # Парсим строку вида: drwxr-xr-x  ... /user/spark/data/bronze/date=2022-01-15
        parts = line.split()
        if len(parts) < 8:
            continue
        path = parts[-1]
        if "date=" not in path:
            continue

        date_str = path.split("date=")[-1].rstrip("/")
        try:
            partition_date = datetime.strptime(date_str[:10], "%Y-%m-%d")
            if partition_date < cutoff_date:
                old_partitions.append(path)
        except ValueError:
            continue

    logger.info(f"Найдено {len(old_partitions)} партиций старше {months} месяцев в {hdfs_path}")
    return old_partitions


def setup_cold_directory(cold_path: str, ec_policy: str) -> None:
    """
    Создаёт целевую директорию для холодных данных и устанавливает EC политику.
    EC политика должна быть установлена ДО копирования данных!
    """
    # Создаём директорию если не существует
    subprocess.run(["hdfs", "dfs", "-mkdir", "-p", cold_path], check=True)

    # Устанавливаем EC политику
    subprocess.run(
        ["hdfs", "ec", "-setPolicy", "-path", cold_path, "-policy", ec_policy],
        check=True
    )
    logger.info(f"EC политика {ec_policy} установлена на {cold_path}")


def migrate_partition_to_ec(source_path: str, cold_base: str, ec_policy: str) -> bool:
    """
    Мигрирует одну партицию из репликации в EC через hdfs distcp.

    Алгоритм:
    1. Создаём целевую директорию с EC политикой
    2. Копируем данные через distcp (distcp читает старые файлы, пишет новые с EC)
    3. Верифицируем что целевые файлы корректны (checksums)
    4. Удаляем исходные файлы
    5. Создаём symlink или обновляем Hive метастор (если используется)
    """
    partition_name = source_path.split("/")[-1]
    target_path = f"{cold_base}/{partition_name}"

    logger.info(f"Начинаем миграцию: {source_path}{target_path}")

    # Шаг 1: настройка целевой директории
    setup_cold_directory(cold_base, ec_policy)

    # Шаг 2: копирование через distcp
    # -p: сохранять атрибуты файлов (timestamps, permissions)
    # -update: пропускать уже скопированные файлы (идемпотентность)
    # -delete: удалять файлы в целевой директории которых нет в источнике
    distcp_result = subprocess.run([
        "hadoop", "distcp",
        "-p",
        "-update",
        source_path,
        target_path,
    ], capture_output=True, text=True)

    if distcp_result.returncode != 0:
        logger.error(f"distcp завершился с ошибкой: {distcp_result.stderr}")
        return False

    # Шаг 3: верификация checksums
    verify_result = subprocess.run([
        "hdfs", "dfs", "-expunge"  # убираем из корзины
    ])

    fsck_result = subprocess.run([
        "hdfs", "fsck", target_path,
        "-files", "-blocks", "-locations"
    ], capture_output=True, text=True)

    if "Status: HEALTHY" not in fsck_result.stdout and "HEALTHY" not in fsck_result.stdout:
        logger.warning(f"fsck предупреждение для {target_path}")

    # Шаг 4: удаление исходных файлов (только после успешного копирования!)
    delete_result = subprocess.run([
        "hdfs", "dfs", "-rm", "-r", "-skipTrash", source_path
    ], capture_output=True, text=True)

    if delete_result.returncode != 0:
        logger.error(f"Ошибка удаления источника: {delete_result.stderr}")
        return False

    logger.info(f"Миграция завершена: {partition_name}")
    return True


def run_ec_migration(**context) -> None:
    """
    Основная функция DAG: сканирует директории и мигрирует старые партиции.
    """
    for config in EC_MIGRATION_CONFIG:
        logger.info(f"Обрабатываем: {config['description']}")

        old_partitions = get_partitions_older_than(
            config["source_base"],
            config["age_months"]
        )

        if not old_partitions:
            logger.info(f"Нет партиций для миграции в {config['source_base']}")
            continue

        logger.info(f"Найдено {len(old_partitions)} партиций для миграции в EC")

        success_count = 0
        error_count = 0

        for partition in old_partitions:
            try:
                success = migrate_partition_to_ec(
                    partition,
                    config["cold_base"],
                    config["ec_policy"]
                )
                if success:
                    success_count += 1
                else:
                    error_count += 1
            except Exception as e:
                logger.error(f"Ошибка при миграции {partition}: {e}")
                error_count += 1

        logger.info(
            f"Итого: {success_count} успешно, {error_count} ошибок "
            f"из {len(old_partitions)} партиций"
        )


with DAG(
    dag_id="hdfs_ec_cold_migration",
    default_args=DEFAULT_ARGS,
    description="Автоматический перевод архивных данных HDFS в Erasure Coding",
    schedule_interval="0 2 1 * *",  # Каждый первый день месяца в 2:00
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=["hdfs", "ec", "archival"],
) as dag:

    migrate_task = PythonOperator(
        task_id="migrate_cold_partitions_to_ec",
        python_callable=run_ec_migration,
        doc_md="""
        Сканирует Bronze и Silver директории HDFS.
        Переводит партиции старше N месяцев в Erasure Coding.
        Экономия: ~50% дискового пространства на архивных данных.
        """,
    )

6. Мониторинг и аудит EC-пайплайнов через Spark UI и метрики

Работа со смешанным хранилищем (часть данных в репликации, часть в EC) требует понимания того, какие задачи читают EC-данные и как это влияет на производительность.

Идентификация EC-чтения в логах Spark

При чтении EC-файлов HDFS использует специальный класс DFSStripedInputStream вместо стандартного DFSInputStream. Это можно увидеть в DEBUG-логах:

# Включение детальных логов чтения HDFS
spark-submit \
  --conf "spark.driver.extraJavaOptions=-Dlog4j.logger.org.apache.hadoop.hdfs=DEBUG" \
  --conf "spark.executor.extraJavaOptions=-Dlog4j.logger.org.apache.hadoop.hdfs=DEBUG" \
  your_job.py
# Лог при чтении EC-файла (DFSStripedInputStream):
DEBUG DFSStripedInputStream: opening striped input stream for file
  /user/spark/data/cold/bronze/2022/part-001.parquet
DEBUG StripedBlockUtil: file is striped with policy RS-6-3-1024k:
  numDataUnits=6, numParityUnits=3, cellSize=1048576
DEBUG StripedBlockUtil: stripe 0: block BP-123.../blk_001@DN1, blk_002@DN2...

# Лог при чтении обычного реплицированного файла (DFSInputStream):
DEBUG DFSInputStream: opening input stream for file
  /user/spark/data/bronze/current/part-001.parquet
DEBUG BlockReaderLocal: creating short-circuit reader for block blk_001
  (NODE_LOCAL)

По наличию DFSStripedInputStream vs DFSInputStream в логах можно однозначно определить, читаются ли данные из EC или из репликации.

Косвенные признаки EC-чтения в Spark UI

Если уровень логирования не включён, на EC-данные указывают косвенные признаки в Spark UI:

Мониторинг EC через HDFS метрики

# ec_monitoring.py
# Мониторинг состояния Erasure Coding на HDFS кластере

import requests
import json
from typing import Optional


def get_ec_stats(namenode_host: str, port: int = 9870) -> dict:
    """
    Получает статистику EC операций с NameNode через JMX.

    Ключевые метрики:
    - ECReconstructionBytesRead: байт прочитано для EC реконструкции
    - ECReconstructionWritten: байт записано при EC реконструкции
    - ECReconstructionBlocks: число блоков на реконструкции
    - ECDecodingBlocksCount: текущее число блоков в процессе декодирования
    """
    url = f"http://{namenode_host}:{port}/jmx?qry=Hadoop:service=NameNode,name=NameNodeInfo"

    try:
        response = requests.get(url, timeout=10)
        data = response.json()

        for bean in data.get("beans", []):
            if "NameNodeInfo" not in bean.get("name", ""):
                continue

            return {
                "ec_policy_stats": bean.get("ErasureCodingPolicies", ""),
                "under_replicated_blocks": bean.get("UnderReplicatedBlocks", 0),
                "pending_deletions": bean.get("PendingDeletionBlocks", 0),
            }
    except Exception as e:
        return {"error": str(e)}


def get_datanode_ec_metrics(datanode_host: str, port: int = 9864) -> dict:
    """
    Получает EC-специфичные метрики DataNode:
    - ECReconstructionBytesRead: байт прочитано для восстановления EC блоков
    - ECReconstructionWritten: байт записано в ходе EC реконструкции
    - ECDecodingMillis: время затраченное на EC декодирование
    """
    url = (f"http://{datanode_host}:{port}/jmx"
           f"?qry=Hadoop:service=DataNode,name=DataNodeActivity")

    try:
        response = requests.get(url, timeout=5)
        data = response.json()

        for bean in data.get("beans", []):
            if "DataNodeActivity" not in bean.get("name", ""):
                continue

            ec_bytes_read = bean.get("ECReconstructionBytesRead", 0)
            ec_bytes_written = bean.get("ECReconstructionWritten", 0)
            total_bytes_read = bean.get("BytesRead", 1)  # избегаем деление на 0

            ec_reconstruction_pct = (
                ec_bytes_read / total_bytes_read * 100
                if total_bytes_read > 0 else 0
            )

            return {
                "hostname": datanode_host,
                "ec_reconstruction_bytes_read": ec_bytes_read,
                "ec_reconstruction_bytes_written": ec_bytes_written,
                "ec_reconstruction_pct_of_total_reads": ec_reconstruction_pct,
                "ec_decoding_millis": bean.get("ECDecodingMillis", 0),
            }
    except Exception as e:
        return {"hostname": datanode_host, "error": str(e)}


def audit_ec_health(datanodes: list[str], threshold_pct: float = 20.0) -> None:
    """
    Аудит EC health: предупреждает если доля EC реконструкции
    превышает порог (признак частых отказов DataNode).

    При высоком EC reconstruction:
    - DataNode часто недоступны (проблемы с железом или сетью)
    - Spark-задачи тратят лишние CPU и сетевые ресурсы на декодирование
    - Нужно проверить состояние кластера
    """
    print("\n" + "=" * 60)
    print("АУДИТ ERASURE CODING МЕТРИК")
    print("=" * 60)

    problem_nodes = []

    for dn in datanodes:
        metrics = get_datanode_ec_metrics(dn)

        if "error" in metrics:
            print(f"  {dn}: ОШИБКА - {metrics['error']}")
            continue

        pct = metrics["ec_reconstruction_pct_of_total_reads"]
        status = "✅" if pct < threshold_pct else "⚠️ "

        print(f"  {status} {dn}: EC reconstruction = {pct:.1f}% "
              f"({metrics['ec_reconstruction_bytes_read'] / 1024**3:.1f} GB)")

        if pct >= threshold_pct:
            problem_nodes.append(dn)

    if problem_nodes:
        print(f"\n⚠️  Высокая EC reconstruction нагрузка на {len(problem_nodes)} нодах!")
        print("Возможные причины: нестабильные DataNode, сетевые проблемы, ")
        print("высокая конкуренция за ресурсы при Spark чтении EC данных.")
    else:
        print("\n✅ EC метрики в норме")

7. Практика: настройка EC политик и сравнительный бенчмарк Spark Job

Подготовка стенда: создание HDFS директорий с разными политиками

# Шаг 1: Создать директории для эксперимента
hdfs dfs -mkdir -p /benchmark/replication-3x
hdfs dfs -mkdir -p /benchmark/ec-rs-6-3

# Шаг 2: Проверить текущие политики (пока обе используют репликацию)
hdfs ec -getPolicy -path /benchmark/replication-3x
# ErasureCodingPolicy=null (Default: 3x replication)

hdfs ec -getPolicy -path /benchmark/ec-rs-6-3
# ErasureCodingPolicy=null

# Шаг 3: Установить EC политику на вторую директорию
# ВАЖНО: для RS-6-3 нужно минимум 9 DataNode!
# Проверяем количество доступных DataNode:
hdfs dfsadmin -report | grep "Live datanodes"
# Live datanodes (9)  ← нужно минимум 9 для RS-6-3

hdfs ec -setPolicy -path /benchmark/ec-rs-6-3 -policy RS-6-3-1024k

# Шаг 4: Верифицируем
hdfs ec -getPolicy -path /benchmark/ec-rs-6-3
# ErasureCodingPolicy=RS-6-3-1024k

# Шаг 5: Запомним свободное место ДО записи
hdfs dfs -df -h /benchmark/

Сравнительный бенчмарк PySpark

# benchmark_ec_vs_replication.py
# Сравнение производительности Spark чтения для EC vs репликация данных

import time
import subprocess
from dataclasses import dataclass
from pyspark.sql import SparkSession
from pyspark.sql import functions as F


@dataclass
class BenchmarkResult:
    storage_type: str      # "replication-3x" или "ec-rs-6-3"
    write_time_sec: float  # время записи датасета
    disk_usage_gb: float   # занятое место после записи
    read_time_sec: float   # время аналитического чтения
    agg_time_sec: float    # время агрегации


def create_spark() -> SparkSession:
    return (
        SparkSession.builder
        .master("yarn")
        .appName("ec-vs-replication-benchmark")
        .config("spark.executor.instances", "9")   # по числу DataNode
        .config("spark.executor.cores", "4")
        .config("spark.executor.memory", "8g")
        # Нативные библиотеки для ISA-L EC decode
        .config("spark.executor.extraJavaOptions",
                "-Djava.library.path=/usr/lib/hadoop/lib/native")
        .config("spark.executorEnv.LD_LIBRARY_PATH",
                "/usr/lib/hadoop/lib/native")
        # Для EC данных locality.wait = 0 - NODE_LOCAL всё равно недостижим
        .config("spark.locality.wait", "0")
        .config("spark.sql.adaptive.enabled", "true")
        .getOrCreate()
    )


def get_hdfs_disk_usage_gb(path: str) -> float:
    """Возвращает реальный объём, занятый данными на HDFS (включая parity блоки)."""
    result = subprocess.run(
        ["hdfs", "dfs", "-du", "-s", "-h", path],
        capture_output=True, text=True
    )
    # Формат вывода: "1.5 G  4.5 G  /benchmark/replication-3x"
    # Первое число = размер данных, второе = реальное занятое место
    parts = result.stdout.split()
    if len(parts) >= 3:
        size_str = parts[1]  # второе число = DiskSpaceConsumed (с репликами/parity)
        multiplier = {"K": 1e-6, "M": 1e-3, "G": 1.0, "T": 1e3}.get(size_str[-1], 1.0)
        try:
            return float(size_str[:-1]) * multiplier
        except ValueError:
            return 0.0
    return 0.0


def generate_and_write_dataset(
    spark: SparkSession, output_path: str, n_rows: int = 10_000_000
) -> float:
    """
    Генерирует и записывает тестовый датасет.
    Возвращает время записи в секундах.
    """
    print(f"Генерация датасета ({n_rows:,} строк) → {output_path}")
    start = time.time()

    df = spark.range(n_rows).select(
        F.col("id").alias("transaction_id"),
        (F.rand() * 500_000).cast("long").alias("user_id"),
        (F.rand() * 50_000).cast("long").alias("merchant_id"),
        (F.rand() * 10_000.0).alias("amount"),
        F.array(F.lit("RUB"), F.lit("USD"), F.lit("EUR")).getItem(
            (F.rand() * 3).cast("int")
        ).alias("currency"),
        F.array(
            F.lit("purchase"), F.lit("refund"), F.lit("transfer"), F.lit("fee")
        ).getItem((F.rand() * 4).cast("int")).alias("tx_type"),
        F.date_add(F.lit("2022-01-01"), (F.rand() * 730).cast("int")).alias("tx_date"),
        F.current_timestamp().alias("created_at"),
    )

    # repartition(90) = 10 файлов на DataNode, равномерное распределение
    df.repartition(90, "tx_date").write \
        .mode("overwrite") \
        .partitionBy("tx_date") \
        .parquet(output_path)

    elapsed = time.time() - start
    print(f"Записано за {elapsed:.1f} секунд")
    return elapsed


def run_analytical_query(spark: SparkSession, data_path: str) -> float:
    """
    Аналитический запрос: агрегация транзакций по дате, мерчанту и типу.
    Полное сканирование всего датасета с многоуровневой агрегацией.
    """
    start = time.time()

    df = spark.read.parquet(data_path)

    result = (
        df
        .filter(F.col("tx_type") == "purchase")
        .groupBy(
            F.date_trunc("month", "tx_date").alias("month"),
            "merchant_id",
            "currency"
        )
        .agg(
            F.count("*").alias("tx_count"),
            F.sum("amount").alias("total_amount"),
            F.avg("amount").alias("avg_amount"),
            F.countDistinct("user_id").alias("unique_users"),
        )
        .orderBy(F.desc("total_amount"))
    )

    row_count = result.count()
    elapsed = time.time() - start
    print(f"Результат: {row_count:,} строк за {elapsed:.1f} секунд")
    return elapsed


def run_full_benchmark() -> None:
    spark = create_spark()

    REPLICATION_PATH = "hdfs://mycluster/benchmark/replication-3x"
    EC_PATH = "hdfs://mycluster/benchmark/ec-rs-6-3"
    N_ROWS = 10_000_000  # ~1 GB данных

    results = {}

    for storage_type, path in [("Репликация 3x", REPLICATION_PATH),
                                ("EC RS-6-3", EC_PATH)]:
        print(f"\n{'='*60}")
        print(f"Тест: {storage_type}")
        print(f"Путь: {path}")
        print(f"{'='*60}")

        # Запись
        write_time = generate_and_write_dataset(spark, path, N_ROWS)

        # Занятое место
        disk_gb = get_hdfs_disk_usage_gb(path)
        print(f"Занято на HDFS: {disk_gb:.1f} GB")

        # Прогрев (warm-up)
        print("\nПрогрев...")
        spark.read.parquet(path).count()

        # Чтение с агрегацией (3 запуска)
        read_times = []
        for i in range(1, 4):
            t = run_analytical_query(spark, path)
            read_times.append(t)
            print(f"  Запуск {i}: {t:.1f} сек")

        avg_read = sum(read_times) / len(read_times)

        results[storage_type] = {
            "write_time_sec": write_time,
            "disk_gb": disk_gb,
            "avg_read_time_sec": avg_read,
        }

        time.sleep(10)

    # Вывод сравнения
    print("\n" + "=" * 60)
    print("ИТОГОВОЕ СРАВНЕНИЕ")
    print("=" * 60)

    rep = results["Репликация 3x"]
    ec = results["EC RS-6-3"]

    disk_savings_pct = (1 - ec["disk_gb"] / rep["disk_gb"]) * 100
    read_slowdown = ec["avg_read_time_sec"] / rep["avg_read_time_sec"]

    print(f"{'Метрика':<35} {'Репликация 3x':>15} {'EC RS-6-3':>15}")
    print("-" * 65)
    print(f"{'Время записи (сек)':<35} {rep['write_time_sec']:>14.1f}s "
          f"{ec['write_time_sec']:>14.1f}s")
    print(f"{'Занято на HDFS (GB)':<35} {rep['disk_gb']:>14.1f}  "
          f"{ec['disk_gb']:>14.1f}")
    print(f"{'Среднее время чтения (сек)':<35} {rep['avg_read_time_sec']:>14.1f}s "
          f"{ec['avg_read_time_sec']:>14.1f}s")
    print()
    print(f"Экономия дискового пространства: {disk_savings_pct:.0f}%")
    print(f"Замедление чтения: {read_slowdown:.2f}x")
    print()
    print("ВЫВОД:")
    if disk_savings_pct > 40 and read_slowdown < 3:
        print("✅ EC RS-6-3 рекомендуется для холодных/архивных данных.")
        print(f"   50% экономии дисков при {read_slowdown:.1f}x замедлении чтения —")
        print("   приемлемый компромисс для редко читаемых данных.")
    else:
        print("⚠️  Результаты зависят от конфигурации кластера.")

    spark.stop()


if __name__ == "__main__":
    run_full_benchmark()

Проверка физического расположения EC блоков

# После записи: проверяем как распределены EC блоки по DataNode
hdfs fsck /benchmark/ec-rs-6-3/tx_date=2022-01-01/ \
  -files -blocks -locations 2>&1 | head -50

# Пример вывода для EC файла:
# /benchmark/ec-rs-6-3/tx_date=2022-01-01/part-0000.parquet:
#   ErasureCoding policy=RS-6-3-1024k
#   Total size:    134217728 B
#   Total dirs:    0
#   Total files:   1
#   Total blocks (validated): 1
#   EC block group with missing blocks:0
#
#   Block group 0 of 1:
#     Data block 0: datanode1:50010 /data/hdfs/current/blk_001 cell=1048576
#     Data block 1: datanode2:50010 /data/hdfs/current/blk_002 cell=1048576
#     Data block 2: datanode3:50010 /data/hdfs/current/blk_003 cell=1048576
#     Data block 3: datanode4:50010 /data/hdfs/current/blk_004 cell=1048576
#     Data block 4: datanode5:50010 /data/hdfs/current/blk_005 cell=1048576
#     Data block 5: datanode6:50010 /data/hdfs/current/blk_006 cell=1048576
#     Parity block 0: datanode7:50010 /data/hdfs/current/blk_007 cell=1048576
#     Parity block 1: datanode8:50010 /data/hdfs/current/blk_008 cell=1048576
#     Parity block 2: datanode9:50010 /data/hdfs/current/blk_009 cell=1048576
#
# FSCK ended at ... with status HEALTHY

# Сравниваем с репликацией:
hdfs fsck /benchmark/replication-3x/tx_date=2022-01-01/ \
  -files -blocks -locations 2>&1 | head -30

# Пример вывода для реплицированного файла:
# Block 0 of 1:
#   Replica 0: datanode1:50010 /data/hdfs/current/blk_001
#   Replica 1: datanode3:50010 /data/hdfs/current/blk_001_r1
#   Replica 2: datanode7:50010 /data/hdfs/current/blk_001_r2
# (весь блок целиком на 3 разных DataNode)

Итоги: когда EC выгоден, а когда нет

Erasure Coding - это не серебряная пуля и не просто «более умная репликация». Это фундаментальный компромисс между дисковым пространством и вычислительными ресурсами + Data Locality.

EC выгоден, когда:

  • Данные читаются редко (раз в квартал / раз в год)
  • Стоимость дисков важнее скорости чтения
  • Данные большие и неизменяемые (исторические архивы, compliance данные)
  • Кластер достаточно велик (минимум 9 DataNode для RS-6-3)
  • Доступны SIMD-инструкции (Intel ISA-L) для ускорения EC decode

EC невыгоден (используйте репликацию), когда:

  • Данные читаются ежедневно Spark ETL-пайплайнами
  • Важна максимальная скорость (HOT слои Bronze/Silver)
  • Кластер маленький (< 9 DataNode для RS-6-3)
  • Файлы маленькие (< размера ячейки × k: для RS-6-3 < 6 MB - EC не эффективен)
  • Используется Short-Circuit Reads (EC полностью ломает Short-Circuit)

Универсальная рекомендация: Hot/Cold Tiering. Репликация 3x для активных данных последних 12–24 месяцев, EC RS-6-3 для архивных данных - даёт 30–40% экономии общего дискового пространства кластера при минимальном влиянии на производительность активных Spark-пайплайнов.