Data Engineering и ETL Пакетная обработка: MapReduce, Spark, партиционирование
0%

Пакетная обработка: MapReduce, Spark, партиционирование

Пакетная обработка: 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 Полный конвейер фаз

Фаз на самом деле шесть, и понимать надо все:

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 выглядит физически

Фазы map, shuffle и reduce: 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.

3.3 Кто с кем разговаривает во время джобы

Важная деталь — 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

Раскладка таблицы по каталогам dt и country и отсечение лишних партиций при чтении

Смысл прост: если фильтр запроса совпадает с ключом партиционирования, движок вообще не открывает лишние файлы. Это не «на 20% быстрее» — это на порядки меньше прочитанных байт, а в облаке ещё и прямая экономия денег (BigQuery и Athena тарифицируются по объёму сканирования).

Правила выбора ключа:

  1. Ключ совпадает с типичным фильтром. В 95% случаев это дата события (dt, event_date).
  2. Целевой размер партиции — 128 МБ … 1 ГБ. Меньше — «проблема мелких файлов», больше — нет параллелизма.
  3. Максимум 1–2 уровня вложенности. dt/country/city/device даёт комбинаторный взрыв каталогов, и листинг метаданных в S3 начинает занимать больше времени, чем чтение данных.
  4. Никогда не партиционируйте по высококардинальному полю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 и идемпотентность

Наивный батч пересчитывает всё каждый раз — это работает до первого терабайта и первого счёта от облака. Продакшн-батч почти всегда инкрементальный: обрабатывает только новые партиции. Жизненный цикл одной партиции витрины:

Схему делают рабочей три требования.

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. Типичные ошибки

  1. collect() на большом датафрейме — тянет всё в память драйвера. Используйте write, limit(n).collect() или toLocalIterator().
  2. mode("overwrite") без partitionOverwriteMode=dynamic — стирает всю таблицу.
  3. count() ради проверки внутри пайплайна — каждый count() это полный проход; три отладочных = тройная стоимость джобы.
  4. Нет cache() там, где датафрейм используется дважды — из-за ленивости Spark пересчитает всю ветку графа. Обратная крайность — cache() всего подряд забивает память и вызывает spill.
  5. Партиционирование по высококардинальному полю — миллион каталогов, метастор умирает.
  6. Мелкие файлы без compaction — через полгода пайплайн замедлился в 5 раз без изменений в коде.
  7. Игнорирование перекоса. «Джоба стала медленной» почти всегда = один ключ разросся; смотрите max/median.
  8. Python UDF в горячем пути вместо встроенных функций или pandas_udf.
  9. Джоба, не идемпотентная по времени запуска — использует CURRENT_DATE вместо переданной даты выполнения. Backfill за прошлый месяц пересчитает всё как «сегодня» и тихо испортит данные.
  10. Слепая вера в дефолтные 200 shuffle-партиций на многотерабайтном датасете.
  11. SELECT * из широкой Parquet-таблицы — колоночный формат позволяет читать 3 колонки из 300; не отказывайтесь от бесплатного ускорения в 100×.
  12. Отсутствие алертов по стоимости — в облаке ушедшая вразнос джоба стоит реальных денег.

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 не испортит данные.
  • И самое неудобное правило: до сотни гигабайт распределённый движок чаще вредит, чем помогает.

Источники

Что дальше

Батч упирается в фундаментальное ограничение: результат появляется только после того, как окно данных закрылось. Если бизнесу нужен ответ через секунды после события, нужна другая модель вычисления — с неограниченным потоком, окнами по времени события и явными гарантиями доставки.

Дальше: Потоковая обработка: Kafka, Flink, окна и семантика доставки.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов