Пакетная обработка: MapReduce, Spark, партиционирование
Пакетная обработка (batch processing) — это вычисление над ограниченным набором данных: у него есть начало и конец, и мы знаем оба. «Выручка за вчера», «пересобрать витрину пользователей за 90 дней», «сматчить логи с профилями» — всё это батч. Противоположность — потоковая обработка, где данные не заканчиваются никогда; ей посвящена следующая статья трека.
Не спешите считать батч устаревшим. Большая часть данных в компаниях до сих пор считается батчами, потому что батч дёшев, детерминирован и его можно пересчитать заново. Последнее — главная суперсила: нашли баг в логике — просто перезапустили джобу; стриминговый пайплайн пришлось бы чинить вместе с накопленным состоянием.
Эта статья — про то, как батч физически устроен: почему он распараллеливается, где теряется производительность, и почему почти все проблемы в проде сводятся к одному слову — shuffle. Предполагается, что вы прочитали https://courses.digitable.life/post/data-engineering/01-etl-vs-elt/ (откуда берутся данные) и https://courses.digitable.life/post/data-engineering/02-data-modeling/ (в какую форму мы их приводим).
1. Первый принцип: почему нужен распределённый батч
Конкретные числа. Один NVMe читает ~2 ГБ/с: терабайт логов вы прочитаете за ~8 минут, даже не начав обрабатывать; петабайт — за 5.8 суток. Оптимизация однопоточного кода тут не поможет — упирается физика носителя, а не CPU. Единственный выход — читать с N дисков одновременно. Отсюда вся конструкция: данные заранее нарезаны на блоки по многим машинам (HDFS-блоки 128 МБ, объекты S3, файлы Parquet); вычисление отправляется к данным (передвинуть 5 КБ кода дешевле, чем 128 МБ данных); каждая машина обрабатывает свой кусок независимо, и только потом результаты сводятся.
Всё, что считается независимо по кускам, масштабируется почти линейно. Всё, что требует сведения данных с разных машин, упирается в сеть. Эта граница — главный водораздел батча.
Запомните этот рисунок. Вся дальнейшая оптимизация батча — это «как сделать так, чтобы правая колонка случалась реже, обрабатывала меньше байт и не перекашивалась».
2. MapReduce: модель, а не только фреймворк
MapReduce описан в статье Google 2004 года (Dean & Ghemawat). Сам движок Hadoop MapReduce сегодня почти не используют, но модель осталась основой всего: Spark, Flink, Presto, BigQuery внутри делают ту же последовательность фаз.
2.1 Контракт модели
Пользователь пишет две чистые функции — map(k1, v1) -> list[(k2, v2)] и reduce(k2, list[v2]) -> list[(k3, v3)]. Всё остальное — параллелизм, перемещение данных, перезапуск упавших задач — делает фреймворк. Ключевое ограничение и одновременно источник силы: reduce видит все значения одного ключа и только их. Значит, фреймворк обязан собрать все записи с одинаковым k2 на одну машину. Это и есть shuffle.
2.2 Полный конвейер фаз
Фаз на самом деле шесть, и понимать надо все:
парсинг, эмиссия пар k,v"] M --> C["2 · combine
локальная предагрегация
необязательна, но резко снижает сеть"] C --> P["3 · partition
hash k mod R → номер редьюсера"] P --> SP["4 · spill + sort
запись отсортированных сегментов
на локальный диск"] SP --> SH["5 · shuffle / fetch
редьюсер тянет свои сегменты
со всех мапперов по сети"] SH --> MS["6 · merge-sort + reduce
слияние отсортированных потоков,
вызов reduce на каждый ключ"] MS --> OUT["Выходной файл part-r-0000N"]
Combine — это reduce, применённый локально до отправки по сети. Для суммы, счётчика, min/max он корректен, потому что операция ассоциативна и коммутативна. Для «среднего» — некорректен: avg(avg(a,b), c) ≠ avg(a,b,c). Правильный приём — считать пару (sum, count) и делить в самом конце. Общее правило: чтобы агрегат распараллеливался, он должен быть моноидом — иметь ассоциативную операцию слияния и нейтральный элемент.
Partition решает, какому редьюсеру достанется ключ; по умолчанию hash(key) % R. Отсюда две классические беды: перекос (один ключ = 40% данных) и то, что число редьюсеров R жёстко задаёт число выходных файлов.
Sort перед reduce обязателен: слияние отсортированных потоков позволяет отдавать редьюсеру группы ключей потоково, не держа всю группу в памяти. Именно поэтому Hadoop мог посчитать группу, которая не влезает в RAM, а наивный dict в Python — не может.
2.3 Как shuffle выглядит физически
Две важные вещи. Во-первых, между стадиями стоит барьер: reduce не может начаться, пока не завершился последний map-таск (иначе редьюсер не знает, что получил все значения ключа). Во-вторых, промежуточные данные материализуются на диск — это делает джобу устойчивой (упавший редьюсер перечитает shuffle-файлы), но и медленной.
2.4 Word count: псевдокод и рабочая реализация
map(document_id, text):
for word in tokenize(text): emit(word, 1)
combine(word, counts): emit(word, sum(counts)) # тот же reduce, локально
reduce(word, counts): emit(word, sum(counts))
Ниже — честная реализация фаз на чистом Python. Распределённости нет, но есть все шесть фаз, включая spill на диск и merge-sort:
import heapq, itertools, json, os, re, tempfile
from collections import defaultdict
TOKEN_RE = re.compile(r"[a-zа-яё]+")
def map_phase(lines):
"""Фаза map: строка -> поток пар (слово, 1)."""
for line in lines:
for word in TOKEN_RE.findall(line.lower()):
yield word, 1
def spill_sorted(pairs, spill_dir, memory_limit=100_000):
"""Фазы combine + spill: копим в памяти, при переполнении предагрегируем,
сортируем по ключу и сбрасываем сегмент на диск.
Память: O(memory_limit), а НЕ O(число уникальных ключей)."""
paths, buffer = [], defaultdict(int)
def flush():
if not buffer:
return
fd, path = tempfile.mkstemp(dir=spill_dir, suffix=".seg")
with os.fdopen(fd, "w", encoding="utf-8") as f:
for key in sorted(buffer): # сортировка сегмента: O(m log m)
f.write(json.dumps([key, buffer[key]]) + "\n")
buffer.clear()
paths.append(path)
for key, value in pairs:
buffer[key] += value # combine: локальная предагрегация
if len(buffer) >= memory_limit:
flush()
flush()
return paths
def merge_and_reduce(paths, num_reducers=2):
"""Фазы shuffle + merge-sort + reduce. heapq.merge сливает уже отсортированные
сегменты потоково: память O(число сегментов), а не O(объём данных)."""
def read(path):
with open(path, encoding="utf-8") as f:
for line in f:
yield tuple(json.loads(line))
merged = heapq.merge(*(read(p) for p in paths), key=lambda kv: kv[0])
results = [dict() for _ in range(num_reducers)]
for key, group in itertools.groupby(merged, key=lambda kv: kv[0]):
target = hash(key) % num_reducers # фаза partition
results[target][key] = sum(v for _, v in group) # фаза reduce
return results
if __name__ == "__main__":
corpus = ["мама мыла раму", "раму мыла мама снова", "снова снова раму"]
with tempfile.TemporaryDirectory() as tmp:
segments = spill_sorted(map_phase(corpus), tmp, memory_limit=2)
for i, part in enumerate(merge_and_reduce(segments)):
print(f"part-r-{i:05d}: {dict(sorted(part.items()))}")
Ключевой урок кода: память ограничена сортировкой и слиянием, а не размером датасета. Наивный Counter(all_words) на терабайте логов упал бы по OOM. Так работает external sort — и так же работает sortMergeJoin в Spark.
2.5 Сложность
Пусть N — общий объём данных, P — число параллельных задач, M = N/P — размер блока на задачу, k — доля данных, выживающая после combine.
| Фаза | Время (wall clock) | Сеть | Диск |
|---|---|---|---|
| map | O(M) |
0 | чтение N |
| combine | O(M) |
0 | 0 |
| spill + sort | O(M log M) |
0 | запись k·N |
| shuffle | O(k·N / (P · bandwidth)) |
k·N |
чтение k·N |
| reduce (merge) | O(M log P) |
0 | запись результата |
Итого время ≈ O(N/P · log(N/P)) + O(k·N / bandwidth). Второе слагаемое не делится на P: сколько байт нужно перекинуть по сети, столько и придётся, сколько машин ни добавь — сеть общая. Отсюда практическое правило: добавление узлов помогает, пока джоба CPU-bound, и перестаёт помогать, как только она стала shuffle-bound. Это закон Амдала, применённый к I/O. По памяти на задачу — O(размер буфера + число сливаемых сегментов), константа, не зависящая от N.
2.6 Почему ушли от Hadoop MapReduce
Каждая пара map+reduce материализуется в HDFS: итеративный алгоритм (PageRank, k-means) из 20 итераций пишет и читает датасет 20 раз, и репликация ×3 съедает всё. Оптимизатора нет — порядок джойнов и проталкивание фильтров были на совести инженера, вручную, на Java. Модель выражения бедная: один SELECT ... JOIN ... GROUP BY превращался в 3–5 связанных джоб (отсюда надстройки Hive и Pig). Spark родился как ответ на первый пункт — держать промежуточные данные в памяти между стадиями (Zaharia et al., RDD, NSDI 2012).
3. Spark: DAG вместо цепочки джоб
3.1 RDD, lineage и ленивость
Базовая абстракция Spark — RDD: неизменяемая распределённая коллекция, разбитая на партиции. Ключевая идея — lineage (линия происхождения): RDD помнит не данные, а рецепт, по которому получен из родителей. Упал узел — Spark пересчитает только потерянные партиции по рецепту, а не восстановит из реплики. Это отказоустойчивость без репликации промежуточных данных.
Отсюда же ленивость: трансформации (map, filter, join) только достраивают граф, вычисление запускает лишь action (count, collect, write). Пока action не вызван, оптимизатор видит весь граф целиком и может его переписать.
3.2 Узкие и широкие зависимости — главное различие в Spark
- Narrow (узкая): каждая партиция-родитель питает ровно одну партицию-потомка —
map,filter,union,mapPartitions. Такие операции сливаются в один проход (pipelining), данные не покидают исполнителя. - Wide (широкая): партиция-потомок зависит от многих родителей —
groupByKey,reduceByKey,join,distinct,repartition. Требует shuffle.
Граница стадии (stage) = широкая зависимость. Всё между двумя shuffle выполняется как один конвейер без материализации. Поэтому в Spark UI вы видите «Stage 3 of 7»: семь стадий — это шесть shuffle.
Parquet"] --> F1["filter dt = вчера"] F1 --> M1["select user_id, amount"] end subgraph ST2["Stage 2 · без shuffle"] R2["read users
Parquet"] --> F2["filter is_active"] end M1 -->|"shuffle: hash user_id"| J["Stage 3 · sort-merge join"] F2 -->|"shuffle: hash user_id"| J J --> AGG["groupBy country"] AGG -->|"shuffle: hash country"| W["Stage 4 · агрегация + запись"] classDef wide fill:#b4562f,stroke:#8a3d1f,color:#fff class J,W wide
3.3 Кто с кем разговаривает во время джобы
выбор стратегии join по статистикам C->>S: физический план S->>S: нарезка на стадии по границам shuffle S->>E1: задачи стадии 1 (партиции 0..3) S->>E2: задачи стадии 1 (партиции 4..7) E1-->>S: shuffle-файлы записаны, метрики E2-->>S: shuffle-файлы записаны, метрики Note over S: барьер: стадия 1 завершена S->>S: AQE читает фактические размеры
и корректирует план S->>E1: задачи стадии 2 (fetch + reduce) E1->>E2: fetch shuffle-блоков по сети E2->>E1: fetch shuffle-блоков по сети E1-->>U: результат / статус записи
Важная деталь — AQE (Adaptive Query Execution) из Spark 3.0. Оптимизатор изначально работает по статистике, которая часто врёт. AQE после каждой стадии смотрит на фактические размеры shuffle-блоков и на лету схлопывает лишние партиции, переключает sort-merge join на broadcast, разбивает перекошенные партиции. Это самый дешёвый прирост производительности из существующих:
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
3.4 RDD, DataFrame или SQL
Пишите на DataFrame/SQL. Catalyst проталкивает фильтры в Parquet, обрезает колонки, меняет порядок джойнов и генерирует байткод (whole-stage codegen), а Tungsten хранит данные в бинарном off-heap формате без боксинга. На RDD оптимизатор не видит ничего, кроме непрозрачной лямбды, — разница на реальных джобах 2–10×. RDD оправдан только для кастомных алгоритмов с нестандартным партиционированием.
Отдельный анти-паттерн — Python UDF: каждая строка сериализуется из JVM в Python-процесс и обратно. Если UDF неизбежен, берите векторизованный pandas_udf (Arrow, батчами) — обычно в 3–20 раз быстрее.
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import DoubleType
import pandas as pd
spark = SparkSession.builder.appName("daily-revenue").getOrCreate()
events = (
spark.read.parquet("s3://lake/events/")
.where(F.col("dt") == "2026-07-15") # pushdown до уровня каталогов партиций
.select("user_id", "amount", "currency") # обрезка колонок на уровне Parquet
)
# Маленький справочник курсов — идеальный кандидат на broadcast join:
# он рассылается на каждый executor целиком, shuffle большой таблицы не нужен.
rates = spark.read.parquet("s3://lake/dim_currency_rates/dt=2026-07-15")
enriched = events.join(F.broadcast(rates), on="currency", how="left")
@F.pandas_udf(DoubleType()) # векторизованный UDF: работает над Series
def apply_fee(amount: pd.Series) -> pd.Series:
return amount * 0.975
daily = (
enriched
.withColumn("net", apply_fee(F.col("amount") * F.col("rate")))
.groupBy("user_id") # единственный неизбежный shuffle
.agg(F.sum("net").alias("revenue"), F.count("*").alias("tx_count"))
)
4. Shuffle: где сгорает бюджет
4.1 groupByKey вместо reduceByKey
rdd.groupByKey().mapValues(sum) # ПЛОХО: по сети едут ВСЕ значения,
# редьюсер держит всю группу в памяти
rdd.reduceByKey(lambda a, b: a + b) # ХОРОШО: combine на маппере, по сети —
# одно число на ключ на маппер
Разница может быть в сотни раз по объёму сети. В DataFrame API groupBy().agg() делает частичную агрегацию сам — ещё один довод в его пользу.
4.2 Перекос данных (data skew)
Самая коварная проблема батча: 199 задач завершились за 20 секунд, а одна работает 40 минут, потому что на неё пришёлся ключ user_id = NULL или country = 'RU' с 60% трафика. Диагностика: в Spark UI на вкладке Stages сравните Min / Median / Max по Shuffle Read Size и Duration. Если max/median > 10 — у вас перекос.
Три лекарства по убыванию предпочтительности: включить AQE skew join (Spark сам разрежет крупные партиции — решает большинство случаев); broadcast join, если одна сторона мала (shuffle исчезает полностью); salting, если не помогло ничего:
SALT = 64 # число «долей», на которые режем горячий ключ
big_salted = big.withColumn("salt", (F.rand() * SALT).cast("int"))
# Маленькую сторону размножаем на SALT копий, чтобы каждая доля нашла пару
small_exploded = small.withColumn(
"salt", F.explode(F.array([F.lit(i) for i in range(SALT)]))
)
joined = big_salted.join(small_exploded, on=["user_id", "salt"]).drop("salt")
Цена: маленькая таблица раздувается в SALT раз. Поэтому salting применяют точечно — только к списку заранее известных горячих ключей, а остальное джойнят обычным способом.
Отдельно про NULL: он хэшируется в одно значение и создаёт гигантскую партицию, хотя семантически NULL != NULL и такие строки всё равно отвалятся. WHERE key IS NOT NULL до джойна иногда ускоряет джобу вдвое одной строкой.
4.3 Стратегии join и их стоимость
| Стратегия | Условие применимости | Shuffle | Стоимость |
|---|---|---|---|
| Broadcast hash join | одна сторона < autoBroadcastJoinThreshold (по умолчанию 10 МБ) |
нет | O(N + S·E), E — число executor’ов |
| Sort-merge join | обе стороны большие, ключ сортируемый | обеих сторон | O((N+M) log((N+M)/P)) |
| Shuffle hash join | одна сторона влезает в память executor’а после shuffle | обеих сторон | O(N + M), но рискует OOM |
| Broadcast nested loop | нет условия равенства (<, BETWEEN) |
нет | O(N · M) — почти всегда катастрофа |
Порог броадкаста (spark.sql.autoBroadcastJoinThreshold) стоит поднимать осознанно — например до 100–256 МБ на жирных executor’ах. Но броадкаст материализуется сначала в драйвере, потом в каждом executor’е: слишком большой порог валит драйвер по OOM. Если в df.explain("formatted") вы видите BroadcastNestedLoopJoin — почти наверняка потерялось условие равенства и джоба не закончится никогда.
5. Партиционирование: два разных смысла одного слова
Источник вечной путаницы, который надо развести раз и навсегда.
- Партиции вычисления (Spark partitions) — куски данных на исполнителях во время джобы. Управляются
spark.sql.shuffle.partitions,repartition(),coalesce(). Живут минуты. - Партиции хранения (table partitions) — физическая раскладка таблицы по каталогам:
dt=2026-07-15/country=RU/. Живут годами.
Первое влияет на скорость текущей джобы. Второе — на скорость всех будущих запросов, и потому важнее.
5.1 Партиционирование хранения и partition pruning
Смысл прост: если фильтр запроса совпадает с ключом партиционирования, движок вообще не открывает лишние файлы. Это не «на 20% быстрее» — это на порядки меньше прочитанных байт, а в облаке ещё и прямая экономия денег (BigQuery и Athena тарифицируются по объёму сканирования).
Правила выбора ключа:
- Ключ совпадает с типичным фильтром. В 95% случаев это дата события (
dt,event_date). - Целевой размер партиции — 128 МБ … 1 ГБ. Меньше — «проблема мелких файлов», больше — нет параллелизма.
- Максимум 1–2 уровня вложенности.
dt/country/city/deviceдаёт комбинаторный взрыв каталогов, и листинг метаданных в S3 начинает занимать больше времени, чем чтение данных. - Никогда не партиционируйте по высококардинальному полю —
user_id,order_id, timestamp с точностью до секунды. Миллион каталогов по 2 КБ — гарантированная смерть пайплайна.
Проблема мелких файлов заслуживает отдельного абзаца. Каждый файл в S3/HDFS — отдельный сетевой запрос, отдельная запись в метасторе, отдельный таск. 100 000 файлов по 1 МБ читаются в десятки раз медленнее, чем 1 000 файлов по 100 МБ при том же объёме. Лечение: coalesce() перед записью, регулярный compaction-джоб, OPTIMIZE в Delta Lake или rewrite_data_files в Iceberg.
5.2 Бакетирование и кластеризация
Партиционирование режет по значению; бакетирование — по хэшу на фиксированное число корзин. Две таблицы, забакетированные по одному ключу с одинаковым числом бакетов, джойнятся без shuffle: Spark знает, что нужные строки уже лежат в соответствующих бакетах.
-- Обе таблицы бакетируем по user_id: последующий join между ними обходится без shuffle
CREATE TABLE events_bucketed (
user_id BIGINT,
amount DECIMAL(18, 2),
dt DATE
)
USING parquet
PARTITIONED BY (dt)
CLUSTERED BY (user_id) INTO 256 BUCKETS;
Trade-off: число бакетов фиксируется на этапе записи, менять его больно, а сама запись дорожает (нужен shuffle при записи). Оправдано, когда одна и та же пара таблиц джойнится десятки раз в день. Современная альтернатива — сортировка / Z-ordering внутри партиции (Delta OPTIMIZE ... ZORDER BY, Iceberg sort orders, ClickHouse ORDER BY): она не убирает shuffle, но резко улучшает min/max-статистику row group’ов, и движок пропускает файлы целиком по predicate pushdown. Подробнее о форматах — https://courses.digitable.life/post/data-engineering/06-storage-and-formats/.
5.3 Партиции вычисления: сколько их должно быть
Дефолт spark.sql.shuffle.partitions = 200 придуман в 2014 году и не подходит почти никому. Практическое правило: число партиций ≈ (объём данных после shuffle) / (128–256 МБ), округлённое вверх до кратного числу ядер кластера (executors × cores), чтобы не было «хвоста» из недогруженной последней волны задач. При 2 ТБ shuffle это ~10 000 партиций, а не 200. С включённым AQE можно смело ставить заведомо большое значение — лишнее Spark схлопнет сам.
И помните разницу: repartition(n) — полный shuffle, партиции равного размера, можно увеличивать и уменьшать; coalesce(n) — без shuffle, только объединяет соседние партиции, только уменьшает и может схлопнуть параллелизм всей стадии (coalesce(1) перед тяжёлым вычислением заставит его выполняться в один поток).
6. Инкрементальность, backfill и идемпотентность
Наивный батч пересчитывает всё каждый раз — это работает до первого терабайта и первого счёта от облака. Продакшн-батч почти всегда инкрементальный: обрабатывает только новые партиции. Жизненный цикл одной партиции витрины:
_SUCCESS источника Готова_к_расчёту --> Считается: планировщик выдал слот Считается --> Записана_во_временный: запись в _tmp/dt=... Записана_во_временный --> Опубликована: атомарная подмена
+ маркер _SUCCESS Считается --> Упала: OOM / битые данные Упала --> Готова_к_расчёту: retry с экспонентой Упала --> Карантин: превышен лимит попыток Опубликована --> Готова_к_расчёту: backfill после фикса логики Опубликована --> [*] Карантин --> [*]
Схему делают рабочей три требования.
1. Идемпотентность. Повторный запуск за ту же дату даёт тот же результат, а не дубли. Реализуется через overwrite партиции, а не append:
# Перезаписываем ТОЛЬКО те партиции, что реально есть в датафрейме, остальные не трогаем.
# Без этой настройки mode("overwrite") сотрёт ВСЮ таблицу — одна из самых дорогих
# ошибок в индустрии. Проверьте свои джобы прямо сейчас.
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
daily.write.mode("overwrite").partitionBy("dt").parquet("s3://lake/marts/daily_revenue/")
2. Атомарная публикация. Читатель не должен увидеть половину записанной партиции. Классика — писать во временный каталог и переименовывать, оставляя маркер _SUCCESS. Но в S3 переименование не атомарно (это copy+delete), поэтому современный ответ — табличные форматы Delta/Iceberg/Hudi с транзакционным журналом.
3. Управление водяным знаком. Джоба должна знать, до какого момента данные уже обработаны, и уметь пересчитать «хвост» из поздно пришедших событий. Типичный компромисс: каждый день пересчитывать не вчерашнюю партицию, а окно [сегодня−3; сегодня−1].
Про сенсоры, ретраи и backfill в Airflow — https://courses.digitable.life/post/data-engineering/05-orchestration/.
7. Когда батч, когда стрим, когда просто SQL
Отдельно скажем честно: если ваших данных меньше 100 ГБ — Spark вам, скорее всего, не нужен. DuckDB, ClickHouse или обычный PostgreSQL посчитают то же самое быстрее и без кластера: современный сервер с 256 ГБ RAM берёт такие объёмы влёгкую, а Spark тратит десятки секунд только на подъём executor’ов и платит за сериализацию. Классика на эту тему — «Scalability! But at what COST?» (HotOS 2015): многие распределённые системы проигрывают однопоточной программе на ноутбуке.
8. Типичные ошибки
collect()на большом датафрейме — тянет всё в память драйвера. Используйтеwrite,limit(n).collect()илиtoLocalIterator().mode("overwrite")безpartitionOverwriteMode=dynamic— стирает всю таблицу.count()ради проверки внутри пайплайна — каждыйcount()это полный проход; три отладочных = тройная стоимость джобы.- Нет
cache()там, где датафрейм используется дважды — из-за ленивости Spark пересчитает всю ветку графа. Обратная крайность —cache()всего подряд забивает память и вызывает spill. - Партиционирование по высококардинальному полю — миллион каталогов, метастор умирает.
- Мелкие файлы без compaction — через полгода пайплайн замедлился в 5 раз без изменений в коде.
- Игнорирование перекоса. «Джоба стала медленной» почти всегда = один ключ разросся; смотрите max/median.
- Python UDF в горячем пути вместо встроенных функций или
pandas_udf. - Джоба, не идемпотентная по времени запуска — использует
CURRENT_DATEвместо переданной даты выполнения. Backfill за прошлый месяц пересчитает всё как «сегодня» и тихо испортит данные. - Слепая вера в дефолтные 200 shuffle-партиций на многотерабайтном датасете.
SELECT *из широкой Parquet-таблицы — колоночный формат позволяет читать 3 колонки из 300; не отказывайтесь от бесплатного ускорения в 100×.- Отсутствие алертов по стоимости — в облаке ушедшая вразнос джоба стоит реальных денег.
9. Как это выглядит в проде
Sizing. Стартовая эвристика: executor на 4–5 ядер и 16–32 ГБ памяти. Больше 5 ядер — деградирует пропускная способность HDFS/S3-клиента; меньше 4 — накладные расходы JVM не окупаются. Оставляйте ~10% на spark.executor.memoryOverhead, иначе YARN/K8s будет убивать контейнеры без внятного сообщения.
Наблюдаемость. Минимум метрик на джобу: время выполнения, прочитанные байты, shuffle read/write, число задач с retry, max/median времени задачи (детектор перекоса), стоимость в деньгах. Включите Spark History Server — без event log разбирать инцидент будет нечем. Ключевая привычка перед выкаткой: смотреть df.explain("formatted") и искать Exchange (это shuffle) и BroadcastNestedLoopJoin (это беда).
Стоимость. Батч отлично живёт на спотовых/прерываемых инстансах: lineage позволяет пересчитать потерянные партиции, типичная экономия 60–80%. Драйвер при этом держат на обычной ноде.
Слоистость. В проде батч почти всегда организован медальонно: bronze (сырые данные как есть, append-only) → silver (очищенные, типизированные, дедуплицированные) → gold (витрины под конкретные вопросы бизнеса). Пересчитать можно с любого слоя вниз — это главная страховка от ошибок в логике. Про контракты между слоями — https://courses.digitable.life/post/data-engineering/07-data-quality-and-governance/.
Тестирование. Батч-джоба тестируется как обычный код: чистые функции трансформации отдельно от чтения/записи, юнит-тесты на маленьких датафреймах в локальном Spark, плюс проверки данных на выходе (dbt test, Great Expectations, Deequ). Отсутствие тестов у пайплайна — не «специфика данных», а обычный технический долг.
10. Мини-итог
- Батч работает с ограниченными данными, и его можно пересчитать — это преимущество перед стримингом, а не недостаток.
- MapReduce задал модель
map → combine → partition → sort → shuffle → reduce; все современные движки внутри делают то же самое, просто быстрее и с оптимизатором. - Операции делятся на узкие (масштабируются линейно, бесплатны) и широкие (требуют shuffle, упираются в сеть). Оптимизация батча = уменьшение числа и объёма shuffle.
- Время ≈
O(N/P · log(N/P))на вычисление плюсO(k·N / bandwidth)на сеть; второе слагаемое не ускоряется добавлением узлов. - Партиционирование хранения важнее партиционирования вычисления: ключ — по типичному фильтру, партиция — 128 МБ … 1 ГБ, вложенность — не больше двух уровней.
- Перекос, мелкие файлы и
groupByKey— три главных убийцы производительности; идемпотентность и атомарная публикация — условие того, что backfill не испортит данные. - И самое неудобное правило: до сотни гигабайт распределённый движок чаще вредит, чем помогает.
Источники
- Jeffrey Dean, Sanjay Ghemawat. MapReduce: Simplified Data Processing on Large Clusters, OSDI 2004.
- Matei Zaharia et al. Resilient Distributed Datasets, NSDI 2012.
- Michael Armbrust et al. Spark SQL: Relational Data Processing in Spark, SIGMOD 2015 — про Catalyst.
- Frank McSherry et al. Scalability! But at what COST?, HotOS 2015.
- Martin Kleppmann. Designing Data-Intensive Applications, главы 10–11 — dataintensive.net; Holden Karau, Rachel Warren. High Performance Spark, O’Reilly — прикладной разбор shuffle и перекоса.
- Spark docs: Performance Tuning и Adaptive Query Execution; Apache Iceberg: Partitioning — скрытое партиционирование и эволюция схемы разбиения.
Что дальше
Батч упирается в фундаментальное ограничение: результат появляется только после того, как окно данных закрылось. Если бизнесу нужен ответ через секунды после события, нужна другая модель вычисления — с неограниченным потоком, окнами по времени события и явными гарантиями доставки.
Дальше: Потоковая обработка: Kafka, Flink, окна и семантика доставки.