Головна › Архітектура корпоративної аналітики
Фабрика даних і розподілені обчислення
Модель виконання MapReduce і Spark, партиціонування та перекіс даних, справжня ціна shuffle, кешування й балансування навантаження — інженерна фізика під кожною корпоративною платформою даних.
У цьому напрямі
Перший шар корпоративної платформи даних — це не про вибір фреймворку, а про фізику. Мережа дорожча за диск, диск дорожчий за пам'ять, а стадія завершується за часом найповільнішої задачі, а не середньої. Цей модуль розбирає три важелі, якими архітектор реально керує: скільки байтів читається, скільки з них перетинає мережу і наскільки рівномірно вони розкладені. Усе, що стоїть вище — семантичний шар, рушій політик, агентні міркування — успадковує затримки та режими відмов, визначені саме тут.
Чому обчислення їде до даних: модель виконання MapReduce і Spark
Базова асиметрія. Скомпільований job важить кілобайти. Набір даних важить терабайти. Тому платформа не тягне дані до коду — вона розсилає код на вузли, де ці байти вже лежать на локальних дисках. Це принцип локальності, з якого виріс MapReduce і який Spark успадкував без змін. Планувальник навіть ранжує розміщення задач: PROCESS_LOCAL (дані в тій самій JVM) → NODE_LOCAL (той самий вузол) → RACK_LOCAL → ANY, і кілька сотень мілісекунд чекає на кращий варіант, перш ніж погодитися на гірший.
Одиниця паралелізму — партиція, а не рядок. Набір даних ділиться на партиції; на кожну партицію припадає рівно одна задача. Кількість партицій задає верхню межу паралелізму: 1 ТБ у 12 партиціях — це 12 задач, і кластер на 480 ядер стоятиме на 97% без роботи, скільки б пам'яті ви не додали.
Ліниве обчислення й межі стадій. Трансформації нічого не виконують — вони будують DAG. Дія (count, write, collect) запускає планувальник, який розрізає DAG на стадії за типом залежностей:
- Вузька залежність: кожна вихідна партиція читає рівно одну вхідну —
filter,select,withColumn. Дані не перетинають мережу, а кілька операцій зливаються в один прохід над рядком. - Широка залежність: вихідна партиція читає багато вхідних —
groupBy,join,distinct,repartition. Це shuffle: запис на диск, передача мережею та обов'язкова межа стадії.
Чому це визначає архітектуру. Наступна стадія не почнеться, доки не завершиться попередня — цілком, до останньої задачі. Тому інженерних важелів рівно три: читати менше (колонковий формат, проєкція, предикати вниз до сховища), перемішувати менше (агрегація до 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 задач чекають на одну
Перекіс — це не збій, це властивість реальних ключів. Shuffle розкладає рядки за hash(key) % N. Хеш рівномірний щодо ключів, а не щодо рядків: якщо 34% таблиці мають один і той самий customer_id, усі ці рядки за визначенням опиняться в одній партиції. Збільшення кількості партицій нічого не змінює — гарячий ключ хешується туди ж, просто поруч з ним з'явиться більше порожніх задач.
Симптом і його читання. У Spark UI дивіться на стадію, а не на job: медіана задачі 30 с, максимум 50 хв. Далі одне порівняння, яке відрізняє два різні діагнози:
- Повільна задача має такі самі вхідні байти й рядки, як медіана → це не перекіс, це хворий вузол (деградований диск, сусід по CPU, GC). Лікується спекулятивним виконанням.
- Повільна задача має у 40 разів більше рядків → це перекіс ключа. Спекуляція не допоможе: копія тієї самої партиції на іншому вузлі оброблятиме ті самі 61 млн рядків.
Три робочі засоби, у порядку вартості.
- Adaptive Query Execution. Планувальник дивиться на реальну статистику на межі стадії й розбиває партицію, що перевищує медіану у
skewedPartitionFactorразів, на підпартиції. Безкоштовно з погляду коду, але працює лише для sort-merge join і не рятує від одного мега-ключа. - Ізоляція виродженого ключа. Найчастіший перекіс у корпоративних даних — це не «популярний клієнт», а сурогат:
NULL,-1,UNKNOWN,guest. Такі рядки за визначенням ні з чим не з'єднуються осмислено — їх виносять в окрему гілку і додаютьunion. - Салтінг. До гарячого ключа дописується випадковий суфікс
0..S-1, а мала сторона реплікується на всі S значень солі черезexplode. Ключ штучно розщеплюється на S партицій. Ціна: мала сторона зростає в S разів — тому S підбирають, а не ставлять 1000.
Партиціонування на диску — інше рішення. Не плутайте партиції виконання з каталогами сховища. Каталогове партиціонування за 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 і економіка кешування
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 на боці, де лежать терабайти. Це найдешевша перемога в шарі. Її ціна — пам'ять драйвера й виконавців, тому поріг підіймають свідомо, а не «на максимум».
Кешування виграє рівно в одному випадку. Коли набір читається більше одного разу і переобчислення дорожче за пам'ять, яку він займе. Правила, що відрізняють користь від шкоди:
- Кеш конкурує за ту саму пам'ять, що й виконання. 340 ГБ кешу на 288 ГБ доступної пам'яті — це не кеш, це гарантований spill плюс витіснення.
persist()лінивий: без дії, що його матеріалізує, перше ж використання все одно порахує все заново.unpersist()обов'язковий, коли набір більше не потрібен — інакше він тримає пам'ять до кінця job.- Кеш не обрізає лінеаж. Проти довгого DAG з дорогою відмовою працює
checkpoint(), який матеріалізує дані в надійне сховище й розриває ланцюг залежностей.
Балансування — це рівномірність, а не кількість. Додавання виконавців допомагає лише тоді, коли є достатньо партицій, щоб їх зайняти, і жодна з них не є на порядок важчою за решту. Інакше нові вузли просто чекають разом зі старими.
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 на кілька гігабайтів.Джерела, з яких виведено напрям
- Spark - The Definitive Guide - Big data processing made simple
- Martin-Kleppmann---Designing-Data-Intensive-Applications -O’Reilly-Media-(2017)
- A Review on Big Data Analytics Framework
- Handling Data Skew in MapReduce Cluster by Using P
- A study of skew in mapreduce application
- A Balanced Solution for the Partition ba
- Characterizing and Fixing Silent Data Loss in Spark on AWS
- DeltaLakeOptimizationTechniquesforScalableLakehouseArchitectures
Пройти інтерактивно
У кожного напряму є питання, картки з інтервальним повторенням і облік прогресу. Для них потрібен акаунт — безкоштовний і на одну хвилину.
Відкрити інтерактивний трек Створити безкоштовний акаунт