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

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

Lakehouse-сховище та інкрементальні подання

ACID поверх об'єктного сховища, формати таблиць Delta Lake та Iceberg, розкладка файлів і компакція, CDC та інкрементальне підтримання подань замість повного перерахунку.

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

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

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

ACID поверх об'єктного сховища: як журнал транзакцій робить із бакета таблицю

E2.1

Об'єктне сховище не є файловою системою. S3, GCS чи MinIO не мають атомарного перейменування каталогу, не мають блокувань, а лістинг префікса історично міг відставати від реальності. Класична схема «таблиця — це префікс у бакеті» ламається саме тут: читач, що стартував посеред запису, бачить половину нових файлів; збійний воркер лишає осиротілі частини; два паралельні записи мовчки перетирають один одного.

Формат таблиці (table format) — це шар метаданих, який робить набір Parquet-файлів транзакційним. Він відповідає рівно на одне питання: які саме файли складають таблицю у версії N. Далі все випливає з відповіді.

Що саме вважається конфліктом. Два дописування (append) у різні файли не конфліктують: обидва лише додають нові файли, і той, хто програв перегони, просто зсуває номер своєї версії. Конфліктують операції, які читають стан і залежать від прочитаного — MERGE, UPDATE, DELETE і компакція. Звідси найдешевший спосіб прибрати більшість конфліктів: звузити область операції предикатом на розділ. Якщо два писачі гарантовано не перетинаються за розділами, перевірка конфлікту пропускає обох, і жоден не чекає на інший.

Де це ламається на практиці. Коміт атомарний рівно настільки, наскільки атомарна операція існує на рівні сховища. Для Delta на S3 із кількох кластерів одночасно потрібен зовнішній координатор (DynamoDB log store) або умовний запис S3; для Iceberg — каталог, який справді підтримує compare-and-swap. Каталог «просто файл у бакеті» повертає вас рівно туди, звідки ви тікали: до тихого перетирання версій.

# What a commit actually is: one conditional PUT decides the winner.
# _delta_log/00000000000000000042.json either exists or it does not.

from deltalake import DeltaTable
from deltalake.exceptions import CommitFailedError

def merge_batch(table_uri, batch, attempts=5):
    for attempt in range(attempts):
        dt = DeltaTable(table_uri)
        pinned = dt.version()                 # optimistic: pin version N
        try:
            (dt.merge(source=batch,
                      predicate="t.order_id = s.order_id",
                      source_alias="s", target_alias="t")
               .when_matched_update_all()
               .when_not_matched_insert_all()
               .execute())                    # tries to create version N+1
            return dt.version()
        except CommitFailedError:
            # Someone else won N+1. Re-read the log, re-check the conflict
            # window, retry. Data files already staged are reused, not rewritten.
            print("lost race at v%d, retry %d" % (pinned, attempt + 1))
    raise RuntimeError("commit contention: %d failed attempts" % attempts)

# The audit trail is a table, not a log file:
#   DESCRIBE HISTORY sales.orders        -> version, operation, operationMetrics
#   RESTORE TABLE sales.orders TO VERSION AS OF 41   -- undo a bad load in seconds

На практиці

Платіжний процесор перевів 4,2 ТБ добових завантажень із Hive-таблиць на Delta. Інциденти «розділ записаний наполовину» впали з 11 за квартал до нуля, а відкат помилкового завантаження скоротився з 6 годин відновлення з бекапу до 90 секунд через RESTORE TO VERSION AS OF.

Антипатерн

Дописувати Parquet-файли напряму в каталог даних таблиці повз API формату («ми ж просто покладемо файл поруч»). Журнал про ці файли не знає: сканування їх ігнорує, звіт мовчки недораховує, а перша ж OPTIMIZE або VACUUM видаляє їх як сміття.

Розкладка файлів, кластеризація та компакція: чому запит відкриває 200 000 об'єктів

E2.2

Вартість запиту в lakehouse визначається не обсягом даних, а кількістю об'єктів, які треба відкрити. Кожен GET до об'єктного сховища — це 20–80 мс до першого байта плюс окремий раунд по метадані. Таблиця на 8 ГБ, розкладена у 200 000 файлів по 40 КБ, читається в десятки разів довше за ту саму таблицю в 60 файлах по 128 МБ, хоча байтів однаково.

Три важелі, і плутати їх не можна:

Компакція — це не прибирання, а частина SLA. Стрімінг із тригером на хвилину створює щонайменше 1 440 файлів на добу на кожен розділ. OPTIMIZE (Delta) чи rewrite_data_files (Iceberg) збирає їх у цільові 128 МБ – 1 ГБ. Але переписані файли не зникають одразу: старі знімки потрібні для time travel і для запитів, які вже виконуються. VACUUM та expire_snapshots звільняють місце лише після вікна ретенції — і саме тут архітектори руйнують собі відновлюваність, виставивши ретенцію 0 годин заради економії на сторіджі.

Copy-on-write проти merge-on-read. COW переписує весь файл заради одного оновленого рядка: читання швидке, запис дорогий, амплітуда запису на точкових апдейтах доходить до сотень разів. MOR пише окремі файли видалень і дельт: запис дешевий, читання платить за злиття, доки не відпрацює компакція. Вектори видалень у Delta — це керований компроміс: рядок позначається позиційно, а фізичне переписування відкладається до планової компакції.

-- Delta: bin-pack, then cluster on the columns the predicates actually use.
OPTIMIZE sales.orders
  WHERE order_date >= current_date() - INTERVAL 7 DAYS
  ZORDER BY (customer_id, status);

-- Retention is a recovery decision, not a cleanup chore.
-- 168h keeps a week of time travel; RETAIN 0 HOURS destroys it.
VACUUM sales.orders RETAIN 168 HOURS;

-- Iceberg: the same two jobs, as explicit procedures.
CALL catalog.system.rewrite_data_files(
  table      => 'sales.orders',
  strategy   => 'sort',
  sort_order => 'customer_id, status',
  options    => map('target-file-size-bytes', '536870912',  -- 512 MB
                    'min-input-files',        '32')
);
CALL catalog.system.expire_snapshots('sales.orders', TIMESTAMP '2026-08-25 00:00:00');

-- Hidden partitioning: the query filters on ts, the engine prunes on days(ts).
-- No derived date_str column for analysts to remember (or forget).
ALTER TABLE sales.orders ADD PARTITION FIELD days(ts);
ALTER TABLE sales.orders ADD PARTITION FIELD bucket(16, customer_id);

На практиці

Роздрібна мережа тримала 11 ТБ подій каси у 1,4 млн файлів після року стрімінгу з хвилинним тригером. Щоденна OPTIMIZE з ZORDER за (store_id, sku) звела таблицю до 9 300 файлів: p95 аналітичного запиту впав із 47 до 3,1 секунди, а рахунок за обчислення — на 38%.

Антипатерн

Розділяти таблицю за високо-кардинальним ключем (PARTITIONED BY (customer_id)) в надії пришвидшити точкові пошуки. Ви отримуєте мільйон каталогів по одному файлу, лістинг метаданих стає довшим за саме сканування, а точкового доступу все одно немає: розділення — це не індекс, для цього існує кластеризація.

CDC та інкрементальні подання: коли дельта дешевша за повний перерахунок

E2.3

Крок 1 — як ви взагалі дізнаєтеся про зміну. Опитування за високою міткою (WHERE updated_at > :watermark) — найдешевший і найпідступніший спосіб: воно не бачить фізичних видалень, губить рядок, оновлений двічі всередині інтервалу, і мовчки ламається щоразу, коли джерело оновлює запис в обхід тригера. Журнальний CDC (WAL або binlog через Debezium) віддає впорядкований потік INSERT/UPDATE/DELETE з LSN — це єдиний спосіб коректно поширити видалення далі за межі OLTP.

Крок 2 — приземлення потоку. Сирі зміни пишуться append-only. Поточний стан отримується MERGE-ом із обов'язковою дедуплікацією: для кожного ключа беремо запис із максимальним LSN у мікробатчі. Без цього два оновлення одного рядка в одному батчі дають недетермінований результат, а рушій законно відмовляється виконувати MERGE.

Крок 3 — підтримання агрегату. Тут ховається головне рішення модуля. Не всі агрегати можна оновити дельтою:

Крок 4 — коли інкрементальність узагалі виправдана. Груба, але робоча оцінка: інкрементальний шлях виграє, поки частка змінених розділів лишається малою — орієнтовно до 10–15% добового churn. Вище цього порогу вартість читання фіду змін, MERGE, переписування файлів і подальшої компакції перевищує вартість чесного CREATE OR REPLACE. Тому зріла архітектура тримає обидві гілки: інкрементальне оновлення щогодини і повний звірювальний перерахунок раз на добу або тиждень, який поглинає наслідки запізнілих і не по порядку прибулих подій.

І де тут NoSQL. Lakehouse дає сканування, join та ACID, але 100–300 мс на точковий пошук — це не те, що ставлять під продуктовий API. Ключ-значення сховище дає одиниці мілісекунд, але не має ані join, ані узгодженого знімка через сутності. Правильна відповідь — не «або-або»: lakehouse лишається системою запису, а KV-сховище є похідною сервісною проєкцією, яку перебудовує той самий потік змін. Дублювання даних тут — свідомо куплена латентність, а не помилка нормалізації.

-- 1. De-duplicate the micro-batch: one row per key, highest LSN wins.
CREATE OR REPLACE TEMP VIEW changes AS
SELECT * FROM (
  SELECT *, ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY lsn DESC) AS rn
  FROM bronze.orders_cdc
  WHERE _commit_version > (SELECT last_version FROM ops.watermarks WHERE job = 'orders_silver')
) WHERE rn = 1;

-- 2. Inserts, updates AND deletes land in one atomic commit.
MERGE INTO silver.orders t
USING changes s ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'D'       THEN DELETE
WHEN MATCHED AND s.lsn > t.lsn    THEN UPDATE SET *
WHEN NOT MATCHED AND s.op <> 'D'  THEN INSERT *;

-- 3. Self-maintainable aggregate: apply +new / -old straight from the feed.
MERGE INTO gold.customer_totals g
USING (
  SELECT customer_id,
         SUM(CASE WHEN _change_type IN ('insert', 'update_postimage') THEN amount
                  ELSE -amount END) AS delta_amount
  FROM table_changes('silver.orders', 812, 947)   -- versions since the last run
  GROUP BY customer_id
) d ON g.customer_id = d.customer_id
WHEN MATCHED     THEN UPDATE SET g.total = g.total + d.delta_amount
WHEN NOT MATCHED THEN INSERT (customer_id, total) VALUES (d.customer_id, d.delta_amount);

-- 4. NOT self-maintainable: MAX cannot be repaired from a delta.
--    Recompute only the groups the change feed actually touched.
CREATE OR REPLACE TABLE gold.customer_peak AS
SELECT customer_id, MAX(amount) AS max_amount
FROM silver.orders
WHERE customer_id IN (
  SELECT DISTINCT customer_id FROM table_changes('silver.orders', 812, 947)
)
GROUP BY customer_id;

На практиці

Телеком-оператор перераховував вітрину споживання на 2,1 млрд рядків повністю щоночі — 3 год 40 хв обчислень. Перехід на фід змін плюс MERGE звів оновлення до 6 хвилин кожні 15 хвилин при добовому churn 0,8%. Повний звірювальний перерахунок лишили щонеділі: за 14 місяців він двічі виявив розбіжність, спричинену запізнілими подіями роумінгу.

Антипатерн

Опитувати джерело за updated_at > :watermark і вважати це CDC. Видалені рядки не мають нового updated_at, тому вони назавжди лишаються у вітрині: звіт показує клієнтів, яких уже не існує, і розбіжність спливає лише під час зовнішнього аудиту.

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

  1. DeltaLakeOptimizationTechniquesforScalableLakehouseArchitectures
  2. Characterizing and Fixing Silent Data Loss in Spark on AWS
  3. LakeVilla A Modular and Non Invasive Toolbox for Lakehouse
  4. Proof Gated Publication Verify Before Commit Content Integ
  5. Smart Compaction Predicting Compaction Utility from Lakeho
  6. Delta Tensor Efficient Vector and Tensor Storage in Delta
  7. An Architecture for Real Time Warehousin
  8. Big Data Techniques Systems Applications
  9. Big Data Analytic Approaches Classificat

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

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

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

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