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.
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 должен:
- Запросить cell1 с DN1 по сети
- Запросить cell2 с DN2 по сети
- Запросить cell3 с DN3 по сети
- Запросить cell4 с DN4 по сети
- Запросить cell5 с DN5 по сети
- Запросить cell6 с DN6 по сети
- Собрать все 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:
- Executor читает все доступные data-ячейки: cell1, cell2, cell4, cell5, cell6 (5 из 6)
- Executor читает один из parity-блоков: p1 (с DN7)
- Декодер Reed-Solomon вычисляет cell3:
cell3 = RS_decode(cell1, cell2, cell4, cell5, cell6, p1) - Только теперь 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-пайплайнов.