Spark vs Hadoop MapReduce: in-memory вычисления и Lazy Evaluation
Почему Spark в 10–100 раз быстрее MapReduce, как работают ленивые вычисления и как Catalyst оптимизирует запросы до начала выполнения.
В предыдущем уроке мы разобрали анатомию кластера: Driver строит DAG, Task Scheduler распределяет задачи по Executor-ам, данные живут в Partition-ах. Теперь разберём почему эта архитектура в 10–100 раз быстрее MapReduce и как механизм Lazy Evaluation позволяет Spark оптимизировать план ещё до первой строки вычислений.
Наследие MapReduce: диск как узкое горлышко¶
MapReduce появился в Google в 2004 году как ответ на вопрос «как обработать петабайт данных на тысячах ненадёжных машин?». Решение элегантное: разбей задачу на Map (обработай каждую запись независимо) и Reduce (агрегируй результаты по ключу), запускай везде параллельно. Apache Hadoop реализовал эту идею как open-source, и к 2010 году стал стандартом для Big Data.
Проблема не в самой модели - проблема в обязательном дисковом I/O.
Каждая MapReduce-задача обязана записать промежуточные данные на HDFS перед тем, как передать их следующему этапу. Это не баг - это гарантия отказоустойчивости: если узел упал посреди Reduce, следующий Job читает уже записанный результат, а не теряет его. За надёжность платят скоростью.
Для одного SQL-запроса с двумя JOIN-ами Hive компилировал в три отдельных MapReduce Job-а - каждый со своим чтением HDFS, записью промежуточных файлов и повторным чтением. Суммарно: 6 операций с распределённым диском на один запрос.
Итеративные алгоритмы: катастрофа¶
Машинное обучение работает итеративно: алгоритм многократно проходит по данным, корректируя веса. Логистическая регрессия - 100 итераций, K-Means - 50 итераций, PageRank - 30 итераций.
В MapReduce каждая итерация = отдельный Job = чтение HDFS + запись HDFS. 100 итераций = 200 операций с HDFS на холодном диске.
Бенчмарк из оригинальной статьи о Spark (Zaharia et al., Berkeley AMPLab, 2012):
| Задача | Hadoop MapReduce | Apache Spark |
|---|---|---|
| Логистическая регрессия (100 итераций) | ~110 сек / итерация | ~0.9 сек / итерация |
| K-Means кластеризация | ~26 мин | ~80 сек |
Разница не в качестве кода - разница в том, что Spark держит данные в памяти между итерациями.
Spark: промежуточные данные в памяти¶
Spark решает проблему на уровне архитектуры: промежуточные данные хранятся в памяти Executor-ов, а не записываются на диск после каждого шага.
Скорость разных типов хранения:
| Тип | Пропускная способность | Latency |
|---|---|---|
| DDR4 RAM | ~50 GB/s | ~100 нс |
| NVMe SSD | ~3–7 GB/s | ~100 мкс |
| SATA SSD | ~0.5 GB/s | ~200 мкс |
| HDD (HDFS) | ~0.1–0.2 GB/s | ~10 мс |
RAM быстрее HDD в 100–500 раз по latency и в 25–50 раз по пропускной способности. Именно это скрывается за формулой «10–100× быстрее».
Для итеративного алгоритма разница наглядна:
Если RAM не хватает - Spark сбрасывает (spill) часть данных на локальный диск Executor-а (не на HDFS), но это исключение. При достаточном объёме памяти итеративный алгоритм читает источник один раз, а дальнейшие итерации работают полностью в RAM.
Когда разница меньше¶
Важная оговорка: «10–100×» - это итеративные сценарии и сложные многошаговые запросы.
| Сценарий | Практический выигрыш |
|---|---|
| Итеративный ML (100 итераций) | 50–100× |
| Сложный ETL: JOIN × 3, 10+ трансформаций | 5–20× |
| Одиночный SELECT с GROUP BY | 2–5× |
| Одиночное полное сканирование большого файла | ~1–2× |
Если задача - один проход по данным с записью на выход, Spark быстрее за счёт лучшего параллелизма и оптимизатора Catalyst, но не за счёт in-memory магии.
DAG: произвольный граф вместо Map → Reduce¶
MapReduce ограничен ровно двумя фазами. Сложный запрос требует нескольких последовательных Job-ов с HDFS между ними.
Spark строит DAG (Directed Acyclic Graph) - произвольный граф зависимостей произвольной глубины. Весь SQL-запрос с тремя JOIN-ами и несколькими агрегациями превращается в один DAG, который оптимизируется как единое целое.
Из урока 2 мы знаем: DAG разбивается на Stage-ы по границам shuffle. Внутри каждого Stage все операции выполняются в один проход по данным без промежуточных записей. На локальный диск пишутся только файлы shuffle - то, что нужно для передачи данных между Stage-ами. HDFS при этом не трогается.
Lazy Evaluation: откладываем всё до последнего¶
Это ключевая концепция урока. Spark не выполняет трансформации сразу. Когда вы пишете df.filter(...) или df.groupBy(...), Spark только записывает эту операцию в план - ничего не происходит.
Выполнение начинается только при вызове Action - операции, требующей реального результата.
# Ни одна из этих строк ничего не вычисляет.
# Spark читает схему Parquet-файла, но не данные.
df = spark.read.parquet("s3a://data/events/")
# Каждая строка ниже добавляет узел в DAG - не более.
filtered = df.filter(df.event_type == "purchase")
enriched = filtered.withColumn(
"amount_usd", col("amount") / col("exchange_rate")
)
aggregated = enriched.groupBy("user_id") \
.agg(sum("amount_usd").alias("total"))
# Вот здесь Spark начинает работать:
# читает Parquet, выполняет filter → withColumn → groupBy → agg → запись
aggregated.write.mode("overwrite").parquet("s3a://result/user_totals/")
Почему это важно? До Action у Spark есть весь план. Оптимизатор Catalyst видит все операции сразу и может перестроить их порядок, устранить лишние вычисления и передать условия фильтрации прямо в Parquet Reader - до того, как первый Task попадёт на Executor.
Если бы Spark выполнял каждую трансформацию немедленно, эти оптимизации были бы невозможны: к моменту write() про filter() уже было бы «забыто».
Трансформации и Actions - полный список¶
| Тип | Что делает | Примеры |
|---|---|---|
| Transformation | Добавляет операцию в DAG, возвращает новый DataFrame | filter(), select(), groupBy(), join(), withColumn(), union(), repartition(), distinct() |
| Action | Запускает выполнение DAG, возвращает результат | count(), collect(), show(), first(), take(N), write.*(), foreach(), toPandas() |
Практическое следствие: не бойтесь длинных цепочек трансформаций. Стоимость - не в их количестве, а в объёме данных и наличии shuffle. Две последовательных filter() - не два прохода, а один после оптимизации.
Catalyst оптимизирует ваш план до выполнения¶
Когда вы вызываете Action, Catalyst проходит несколько стадий оптимизации (детально разбираем в модуле об оптимизации). Два результата, которые ощущаются в производительности прямо сейчас:
Predicate Pushdown¶
Если в вашем коде есть filter(), Catalyst передаёт условие фильтрации вниз к хранилищу - Parquet Reader читает только строки, удовлетворяющие условию. Не поднимает весь файл в память, чтобы потом отфильтровать.
Parquet хранит статистику (min/max значения) для каждой группы строк (row group). Если min_date > '2024-01' - весь row group пропускается без чтения. Это не ошибка оптимизации - это работа Catalyst + формат Parquet вместе.
Column Pruning¶
SELECT user_id, amount из Parquet-файла с 80 колонками читает физически только 2 колонки. Catalyst передаёт список нужных колонок в Parquet Reader - остальные 78 не читаются вообще. В Avro (строковый формат) так не работает: придётся читать всю строку целиком.
Как проверить: explain()¶
explain() показывает физический план до выполнения. Вызывайте перед запуском на production:
df = spark.read.parquet("s3a://data/events/")
result = df.filter(col("event_type") == "purchase") \
.select("user_id", "amount")
result.explain(mode="formatted")
Вывод (упрощённый):
== Physical Plan ==
*(1) Project [user_id#10L, amount#11]
+- *(1) Filter (isnotnull(event_type#12) AND (event_type#12 = purchase))
+- FileScan parquet [user_id#10L, amount#11, event_type#12]
Batched: true,
DataFilters: [isnotnull(event_type#12), (event_type#12 = purchase)],
Format: Parquet,
PushedFilters: [IsNotNull(event_type), EqualTo(event_type,purchase)],
ReadSchema: struct<user_id:bigint,amount:double,event_type:string>
Что читать:
PushedFilters- условия переданы в Parquet Reader. Если пусто - pushdown не сработал (проверьте тип колонки и формат).ReadSchema- только 3 колонки из исходных 80. Column pruning работает.FileScan- чтение с диска / S3, всегда в самом низу плана.
Для сложных запросов с join-ами используйте:
result.explain(mode="extended") # показывает все стадии: Analyzed → Optimized → Physical
Кэширование: когда данные нужны повторно¶
Lazy Evaluation не всегда плюс. Если один и тот же DataFrame используется в нескольких Action-ах, Spark пересчитывает его с нуля каждый раз: перечитывает источник, повторяет все трансформации.
# ❌ Плохо: df_cleaned пересчитывается дважды, S3 читается дважды
df_cleaned = raw.filter(col("status") == "active") \
.dropDuplicates(["user_id"])
row_count = df_cleaned.count() # Action 1: читаем S3, считаем
df_cleaned.write.parquet("s3a://out/") # Action 2: читаем S3 снова
Решение - persist():
# ✅ Хорошо: читаем источник один раз, результат остаётся в RAM
from pyspark import StorageLevel
df_cleaned = raw.filter(col("status") == "active") \
.dropDuplicates(["user_id"]) \
.persist(StorageLevel.MEMORY_AND_DISK)
row_count = df_cleaned.count() # Action 1: читаем S3, результат → RAM
df_cleaned.write.parquet("s3a://out/") # Action 2: читаем из RAM
df_cleaned.unpersist() # Явно освобождаем память
persist() срабатывает при первом Action - это тоже Lazy Evaluation. Второй Action уже читает из памяти.
cache() - алиас для persist(StorageLevel.MEMORY_AND_DISK_DESER). Разница в StorageLevel:
| Уровень | Где | Когда использовать |
|---|---|---|
MEMORY_AND_DISK |
RAM, переполнение → диск | По умолчанию. Защита от OOM |
MEMORY_ONLY |
Только RAM | DataFrame небольшой, нужна максимальная скорость |
DISK_ONLY |
Только диск | RAM ограничена, но пересчёт дорогой (сложный join) |
Когда кэшировать:
- DataFrame используется в 2+ Action-ах
- Вычисление дорогое: несколько join-ов, дедупликация, сложная агрегация
- Итеративный алгоритм (MLlib, Graph-алгоритмы)
Когда не кэшировать:
- DataFrame читается один раз - только лишний overhead на сериализацию
- Данные не помещаются в RAM → постоянный spill, выигрыш нулевой
- Простой проход: filter → select → write без повторного использования
Всегда вызывайте unpersist() явно: Spark не знает, когда вы закончили использовать кэш. Накопленные кэши вытесняют рабочую память и замедляют shuffle.
Антипаттерн: collect() убивает Driver¶
collect() переносит все данные с Executor-ов в память Driver-процесса. Driver - один JVM-процесс с фиксированным spark.driver.memory (по умолчанию 1–2 GB).
# ❌ Катастрофа: 10 млн строк × 10 колонок → Driver OOM → приложение падает
all_rows = df.collect()
# ✅ Вместо этого - пишите результат в хранилище
df.write.mode("overwrite").parquet("s3a://result/")
# ✅ Для отладки - берите ограниченную выборку
sample = df.limit(100).collect() # только 100 строк безопасно
# ✅ Для агрегатных метрик - collect() безопасен на маленьком результате
stats = df.agg(
count("*").alias("total"),
sum("amount").alias("sum_amount"),
avg("amount").alias("avg_amount")
).collect()
# Результат - список из одной Row с тремя числами
print(stats[0]["total"], stats[0]["sum_amount"])
Правило: collect() допустим только если вы точно знаете, что результат мал - после limit(), на агрегированном результате с малым числом ключей, или на groupBy() с заведомо ограниченным числом групп.
При падении по OOM в Spark UI увидите: ExecutorLostFailure или java.lang.OutOfMemoryError: Java heap space на Driver-узле.
Практика: наблюдаем Lazy Evaluation в Spark UI¶
Запустите код и откройте Spark UI (http://localhost:4040):
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as _sum, count
spark = SparkSession.builder \
.appName("LazyEvalDemo") \
.master("local[*]") \
.getOrCreate()
# Датасет в памяти для локального теста
data = [(i, i % 10, float(i * 1.5)) for i in range(1_000_000)]
df = spark.createDataFrame(data, ["id", "category", "value"])
# ── Трансформации: в Spark UI нет ни одного Job ──
step1 = df.filter(col("value") > 100.0)
step2 = step1.withColumn("value_doubled", col("value") * 2)
step3 = step2.filter(col("category") < 5)
step4 = step3.groupBy("category").agg(
count("*").alias("cnt"),
_sum("value").alias("total")
)
# Проверяем план ДО выполнения
step4.explain(mode="formatted")
# ── Action: появится один Job в Spark UI ──
result = step4.collect()
print(f"Результат: {len(result)} групп")
# В Spark UI → Jobs: один Job, два Stage
# Stage 1: filter + withColumn + filter (всё narrow, в один проход)
# Stage 2: groupBy + agg (wide - нужен shuffle)
Ключевое наблюдение: четыре отдельных шага (filter, withColumn, filter, groupBy) Spark выполнил в два Stage - потому что первые три narrow-трансформации склеились в один проход. Это Lazy Evaluation в действии: Spark знал план заранее и объединил операции.
Итог: Spark vs MapReduce¶
| Аспект | MapReduce | Spark |
|---|---|---|
| Промежуточные данные | HDFS (диск) | RAM + spill на локальный диск |
| Граф зависимостей | Map → Reduce (2 фазы) | DAG произвольной глубины |
| Сложный запрос (JOIN × 3) | 3 Job-а, 6 записей HDFS | 1 DAG, 1–2 Stage |
| Итеративный ML (100 итер.) | 200 операций HDFS | 1 чтение + 99 итераций в RAM |
| Оптимизация | Нет | Catalyst: pushdown, pruning, AQE |
| Python API | Hadoop Streaming (неудобно) | PySpark - полноценный API |
MapReduce сегодня практически не используется в новых проектах. Если в компании есть Hadoop-кластер, Spark запускается поверх него через YARN - аппаратная инфраструктура остаётся, вычислительный движок меняется.
Чек-лист¶
- [ ] Трансформации не выполняются - каждая добавляет узел в DAG
- [ ] Action запускает весь DAG целиком
- [ ]
explain()перед запуском на production - убедись, что видишьPushedFilters - [ ]
persist()только если DataFrame используется в 2+ Action-ах - [ ]
unpersist()явно после использования кэша - [ ]
collect()только на заведомо маленьком результате (послеlimit()или на агрегате) - [ ] Для записи результатов -
write.parquet(), неcollect()+ python-файл
В следующем уроке разберём RDD - базовую абстракцию Spark, из которой вырос DataFrame API, и поймём, почему в 2024 году она почти не используется напрямую, но знать её нужно.