AI Architect Trainer Відкрити інтерактивний трек

ГоловнаАрхітектура корпоративної аналітики

Фабрика даних і розподілені обчислення

Модель виконання MapReduce і Spark, партиціонування та перекіс даних, справжня ціна shuffle, кешування й балансування навантаження — інженерна фізика під кожною корпоративною платформою даних.

Востаннє переглянуто: 2026-09-04 · In English

У цьому напрямі

Перший шар корпоративної платформи даних — це не про вибір фреймворку, а про фізику. Мережа дорожча за диск, диск дорожчий за пам'ять, а стадія завершується за часом найповільнішої задачі, а не середньої. Цей модуль розбирає три важелі, якими архітектор реально керує: скільки байтів читається, скільки з них перетинає мережу і наскільки рівномірно вони розкладені. Усе, що стоїть вище — семантичний шар, рушій політик, агентні міркування — успадковує затримки та режими відмов, визначені саме тут.

Чому обчислення їде до даних: модель виконання MapReduce і Spark

E1.1

Базова асиметрія. Скомпільований job важить кілобайти. Набір даних важить терабайти. Тому платформа не тягне дані до коду — вона розсилає код на вузли, де ці байти вже лежать на локальних дисках. Це принцип локальності, з якого виріс MapReduce і який Spark успадкував без змін. Планувальник навіть ранжує розміщення задач: PROCESS_LOCAL (дані в тій самій JVM) → NODE_LOCAL (той самий вузол) → RACK_LOCALANY, і кілька сотень мілісекунд чекає на кращий варіант, перш ніж погодитися на гірший.

Одиниця паралелізму — партиція, а не рядок. Набір даних ділиться на партиції; на кожну партицію припадає рівно одна задача. Кількість партицій задає верхню межу паралелізму: 1 ТБ у 12 партиціях — це 12 задач, і кластер на 480 ядер стоятиме на 97% без роботи, скільки б пам'яті ви не додали.

Ліниве обчислення й межі стадій. Трансформації нічого не виконують — вони будують DAG. Дія (count, write, collect) запускає планувальник, який розрізає DAG на стадії за типом залежностей:

Чому це визначає архітектуру. Наступна стадія не почнеться, доки не завершиться попередня — цілком, до останньої задачі. Тому інженерних важелів рівно три: читати менше (колонковий формат, проєкція, предикати вниз до сховища), перемішувати менше (агрегація до join, broadcast малої сторони) і розкладати рівномірно (урок E1.2).

Що змінилося після класичного MapReduce. MapReduce матеріалізував результат кожної фази в розподілену ФС, тож ітеративні алгоритми платили повний цикл I/O на кожній ітерації. Spark тримає проміжний результат у пам'яті виконавця, а втрачену партицію не реплікує, а переобчислює за лінеажем. Це не «швидше», а інший компроміс: дешевше повторне використання, дорожча відмова наприкінці довгого DAG. Саме на цій властивості тримаються семантичні шари корпоративних платформ — те, що в Palantir Foundry описують як онтологію над фабрикою даних, є похідними наборами, які мусять уміти перебудуватися з першоджерела.

# Stage boundaries are cut at WIDE dependencies. Read the plan, not the docs.
orders = spark.read.parquet('s3a://lake/orders')        # 1.2 TB, 9,600 files

# NARROW: map-side only, no bytes cross the network, stays inside Stage 0
recent = (orders
          .filter(orders.order_ts >= '2026-01-01')
          .select('order_id', 'customer_id', 'total'))

# WIDE: groupBy forces an Exchange -> Stage 0 must fully finish first
per_customer = recent.groupBy('customer_id').sum('total')

per_customer.explain(mode='formatted')
# == Physical Plan ==
# * HashAggregate (final)                              <- Stage 1
# +- Exchange hashpartitioning(customer_id, 200)       <- SHUFFLE = stage cut
#    +- * HashAggregate (partial)                      <- Stage 0, runs on the data
#       +- * Project [order_id, customer_id, total]    <- column projection pushed
#          +- * FileScan parquet
#                PushedFilters: [GreaterThanOrEqual(order_ts, 2026-01-01)]

# Partition count = task count = ceiling on parallelism
print(recent.rdd.getNumPartitions())   # 3,100 -> 3,100 tasks on 480 cores

На практиці

Оператор зв'язку обробляв 4,1 ТБ CDR щодня. Job читав 9 600 Parquet-файлів, але аналітик поставив repartition(24) одразу після читання — «щоб на виході було 24 акуратні файли». Кластер на 480 ядер виконував важку стадію 24 задачами: 5% утилізації, 3 год 40 хв. Прибрали ранній repartition (залишилося 3 100 партицій по ~180 МБ), а coalesce(24) перенесли на крок перед записом. Час — 26 хвилин, вихід побайтово той самий.

Антипатерн

Викликати count() після кожної трансформації «щоб перевірити». Кожен виклик — окрема дія, яка проганяє весь DAG з нуля: п'ять перевірок означають п'ять повних сканувань джерела, а не п'ять дешевих перевірок.

Партиціонування і перекіс даних: чому 199 задач чекають на одну

E1.2

Перекіс — це не збій, це властивість реальних ключів. Shuffle розкладає рядки за hash(key) % N. Хеш рівномірний щодо ключів, а не щодо рядків: якщо 34% таблиці мають один і той самий customer_id, усі ці рядки за визначенням опиняться в одній партиції. Збільшення кількості партицій нічого не змінює — гарячий ключ хешується туди ж, просто поруч з ним з'явиться більше порожніх задач.

Симптом і його читання. У Spark UI дивіться на стадію, а не на job: медіана задачі 30 с, максимум 50 хв. Далі одне порівняння, яке відрізняє два різні діагнози:

Три робочі засоби, у порядку вартості.

Партиціонування на диску — інше рішення. Не плутайте партиції виконання з каталогами сховища. Каталогове партиціонування за event_date дозволяє планувальнику відкинути цілі теки ще до читання. Партиціонування за ключем високої кардинальності (user_id) породжує мільйони дрібних файлів і вбиває продуктивність листингу.

from pyspark.sql import functions as F

# 1) DIAGNOSE before you fix. What share of the table is the heaviest key?
(orders.groupBy('customer_id').count()
       .orderBy(F.desc('count')).show(5, False))
# +-----------+----------+
# |customer_id|count     |
# +-----------+----------+
# |-1         |784215330 |   <- 34% of 2.3B rows: the 'guest checkout' surrogate
# |4417019    |  1204880 |
# +-----------+----------+

# 2) Let the engine handle the moderate tail
spark.conf.set('spark.sql.adaptive.enabled', 'true')
spark.conf.set('spark.sql.adaptive.skewJoin.enabled', 'true')
spark.conf.set('spark.sql.adaptive.skewJoin.skewedPartitionFactor', '5')

# 3) Isolate the degenerate key - it joins to nothing anyway
real   = orders.filter(F.col('customer_id') != -1).join(dim, 'customer_id')
guests = orders.filter(F.col('customer_id') == -1).withColumn('segment', F.lit('GUEST'))
result = real.unionByName(guests, allowMissingColumns=True)

# 4) Salt what is still hot AFTER isolation. Both sides must agree on the salt.
S = 64
big   = orders.withColumn('salt', (F.rand() * S).cast('int'))
small = dim.withColumn('salt', F.explode(F.array([F.lit(i) for i in range(S)])))
joined = big.join(small, ['customer_id', 'salt'])   # small side is now 64x - budget for it

На практиці

Рітейлер з'єднував 2,3 млрд замовлень з довідником клієнтів. 199 із 200 задач стадії завершувалися за 40 с, одна працювала 71 хв. Причина: гостьові замовлення писалися з customer_id = -1 — 784 млн рядків, 34% таблиці, один хеш. Ізоляція цієї гілки в окремий union плюс сіль S=64 на залишковий хвіст дали 4 хв 10 с. Жодного нового вузла не додали.

Антипатерн

Реагувати на straggler підняттям spark.sql.shuffle.partitions з 200 до 4000. Гарячий ключ хешується в ту саму партицію незалежно від їх кількості: ви отримали 3 800 порожніх задач, зайві накладні витрати на планування — і той самий 71-хвилинний хвіст.

Ціна shuffle і економіка кешування

E1.3

Shuffle — єдина операція, що платить обидві ціни. Вона записує весь проміжний набір на локальний диск кожного виконавця (shuffle write), а потім тягне його мережею до інших виконавців (shuffle read). Усе інше — читання, фільтрація, проєкція — платить лише одну. Тому бюджет оптимізації в цьому шарі витрачається на shuffle, а не на «швидший код».

Розмір партиції після shuffle — головний числовий важіль. Цільовий орієнтир: 128–200 МБ на партицію. Дефолт spark.sql.shuffle.partitions = 200 писався для наборів у десятки гігабайтів; на 8 ТБ він дає ~40 ГБ на задачу, що гарантує spill і найчастіше OOM. Занадто дрібні партиції — інша крайність: 100 000 задач по 2 МБ, де планування задачі коштує дорожче за саму роботу.

Spill — це не збій, це рахунок. Коли дані задачі не вміщуються в її частку пам'яті, вони пишуться на диск і читаються назад. Метрики Spill (Memory) / Spill (Disk) у Spark UI — прямий вимір того, наскільки ви промахнулися з розміром партиції.

Broadcast join прибирає shuffle цілком. Якщо одна сторона менша за autoBroadcastJoinThreshold (типово 10 МБ, на практиці підіймають до 32–256 МБ), вона надсилається кожному виконавцю один раз, і велика сторона з'єднується локально — нульовий shuffle на боці, де лежать терабайти. Це найдешевша перемога в шарі. Її ціна — пам'ять драйвера й виконавців, тому поріг підіймають свідомо, а не «на максимум».

Кешування виграє рівно в одному випадку. Коли набір читається більше одного разу і переобчислення дорожче за пам'ять, яку він займе. Правила, що відрізняють користь від шкоди:

Балансування — це рівномірність, а не кількість. Додавання виконавців допомагає лише тоді, коли є достатньо партицій, щоб їх зайняти, і жодна з них не є на порядок важчою за решту. Інакше нові вузли просто чекають разом зі старими.

from pyspark import StorageLevel
from pyspark.sql import functions as F

# 1) Size the shuffle for the DATA, not for the default.
#    8 TB / 200 partitions = 40 GB per task -> spill, then OOM.
#    8 TB / 49,152 partitions ~= 170 MB per task -> inside the target band.
spark.conf.set('spark.sql.shuffle.partitions', 49152)

# 2) Broadcast the small side: the 18 TB fact table is never shuffled.
spark.conf.set('spark.sql.autoBroadcastJoinThreshold', 64 * 1024 * 1024)   # 64 MB
enriched = fact.join(F.broadcast(dim_sku), 'sku_id')                       # dim_sku = 41 MB

# 3) Cache ONLY what is genuinely read more than once, and materialise it once.
sessions = (enriched
            .filter(F.col('event_type').isin('view', 'add_to_cart', 'purchase'))
            .repartition('user_id')
            .persist(StorageLevel.MEMORY_AND_DISK))
sessions.count()                     # single materialisation pass

funnel  = sessions.groupBy('user_id').agg(F.collect_list('event_type').alias('path'))
revenue = sessions.filter(F.col('event_type') == 'purchase').groupBy('sku_id').sum('total')

funnel.write.mode('overwrite').parquet('s3a://lake/marts/funnel')
revenue.write.mode('overwrite').parquet('s3a://lake/marts/revenue')
sessions.unpersist()                 # give the memory back before the next stage

# 4) Long DAG with an expensive late failure? checkpoint(), not cache() -
#    only checkpoint truncates lineage.
spark.sparkContext.setCheckpointDir('s3a://lake/_checkpoints')
stable = enriched.checkpoint(eager=True)

На практиці

Фінансова платформа поставила .cache() на 12 DataFrame «про всяк випадок». Сумарний обсяг кешу — 340 ГБ на 288 ГБ пам'яті виконавців: почалося витіснення, а разом з ним 2,1 ТБ spill на диск. Job став у 2,4 раза повільнішим за версію взагалі без кешу. Залишили persist на одному наборі, який справді читався шість разів, додали count() для матеріалізації та unpersist() після останнього використання: 41 хв замість 96.

Антипатерн

Лишати spark.sql.shuffle.partitions на дефолтних 200 для набору в терабайти, а потім «лікувати» OOM збільшенням пам'яті виконавця. Ви не зменшили обсяг даних на задачу — ви лише зробили кожен виконавець дорожчим і відсунули той самий spill на кілька гігабайтів.

Джерела, з яких виведено напрям

  1. Spark - The Definitive Guide - Big data processing made simple
  2. Martin-Kleppmann---Designing-Data-Intensive-Applications -O’Reilly-Media-(2017)
  3. A Review on Big Data Analytics Framework
  4. Handling Data Skew in MapReduce Cluster by Using P
  5. A study of skew in mapreduce application
  6. A Balanced Solution for the Partition ba
  7. Characterizing and Fixing Silent Data Loss in Spark on AWS
  8. DeltaLakeOptimizationTechniquesforScalableLakehouseArchitectures

Пройти інтерактивно

У кожного напряму є питання, картки з інтервальним повторенням і облік прогресу. Для них потрібен акаунт — безкоштовний і на одну хвилину.

Відкрити інтерактивний трек Створити безкоштовний акаунт

Далі в цьому треку