Перейти до вмісту

Шардинг і розподілена консистентність

«Крамниця» виросла, і один PostgreSQL не встигає писати. Команда ділить дані між вузлами за product_id: товар завжди на одному вузлі, запит про нього торкається одного сервера. У small (Датасет) найгарячіший товар дає 5.2 % позицій замовлень. На чотирьох шардах найбільший несе в 1.17 раза більше за середній, на 64 шардах уже в 4.81 раза. Вузлів побільшало, а навантаження не розподілилось.

Інша проблема помітна не одразу. Покупка зменшує залишок на вузлі A і створює замовлення на вузлі B. Обидва готові закомітити свою частину, але координатор, що мав сказати «комітьте», зник, і рядок залишку лишається заблокованим.

Нарешті, обіцянки виробників. У 2020 році Jepsen перевірив MongoDB, яка рекламувала «повні ACID-транзакції», і знайшов аномалії навіть на найсуворіших налаштуваннях. Щоб розібрати, що було обіцяно й виміряно, потрібні слова «лінеаризованість», «серіалізованість» і «CAP», які часто плутають. Кожне спершу показано на прикладі, а потім дано визначення.

Передумови. Ізоляція, MVCC, блокування: модуль 10. Плани й відсікання: модуль 9. Копії одного набору даних, лаг і перемикання: модуль 12.

Партиціювання і шардинг

Section titled “Партиціювання і шардинг”

Партиціювання ділить таблицю на частини всередині однієї бази: той самий сервер, той самий WAL, ті самі з’єднання. Шардинг кладе частини на окремі сервери, кожен зі своїм процесором, диском і журналом. Партиціювання полегшує обслуговування, але запис не масштабує. Шардинг масштабує, а платить транзакціями й запитами через вузли.

Таблиця view_events росте за часом і старіє, тож у PostgreSQL її ріжуть декларативно за діапазоном occurred_at. Первинний ключ мусить містити ключ партиціювання, бо унікальність перевіряється всередині партиції:

CREATE TABLE view_events_p (
event_id bigint NOT NULL, occurred_at timestamptz NOT NULL, session_id uuid NOT NULL,
customer_id bigint, product_id integer NOT NULL, event_type text NOT NULL,
device text NOT NULL, referrer text,
PRIMARY KEY (event_id, occurred_at)
) PARTITION BY RANGE (occurred_at);
DO $$
DECLARE m date := DATE '2023-12-01';
BEGIN
WHILE m < DATE '2026-01-01' LOOP
EXECUTE format(
'CREATE TABLE view_events_p_%s PARTITION OF view_events_p FOR VALUES FROM (%L) TO (%L)',
to_char(m, 'YYYY_MM'), m::text || ' 00:00+00', (m + interval '1 month')::date::text || ' 00:00+00');
m := (m + interval '1 month')::date;
END LOOP;
END $$;
INSERT INTO view_events_p SELECT * FROM view_events;

Межі записано з явним +00, бо без поясу вони залежали б від TimeZone сесії. Запит за березень 2025 відсікає решту партицій:

СпробуйPostgreSQLCtrl+Enter — виконати

У плані лишається одна партиція, view_events_p_2025_03. Відсікання працює, лише коли в умові є ключ партиціювання: запит за product_id відкриє всі 25.

Друга вигода видна, коли дані старіють. Місяць із цілої таблиці видаляє DELETE з подальшим VACUUM, а партицію можна відʼєднати: вона стає окремою таблицею, яку заархівують чи скинуть DROP:

ALTER TABLE view_events_p DETACH PARTITION view_events_p_2025_03;
SELECT count(*) FROM view_events_p;

Усе це на одному сервері, тож коли запис його перевищив, партиціювання не допоможе.

Шардинг починається з правила «рядок → вузол». Є два способи. Діапазонне: шард відповідає діапазону ключа, і BETWEEN читає сусідні шарди. Воно ламається на ключах, що зростають з часом: шардування orders за placed_at спрямовує весь новий запис на останній шард. Хешоване: шард дорівнює hash(ключ) % N. Навантаження розлітається, але діапазонний запит торкається всіх шардів.

Другий вибір: колонка. Промоделюємо hash % N на позиціях замовлень і порівняємо найбільший шард із середнім.

СпробуйPostgreSQLCtrl+Enter — виконати
shards | key | busiest | average | ratio
--------+-------------+---------+---------+-------
4 | customer_id | 10333 | 9876 | 1.05
4 | order_id | 10133 | 9876 | 1.03
4 | product_id | 11534 | 9876 | 1.17
64 | customer_id | 1067 | 617 | 1.73
64 | order_id | 723 | 617 | 1.17
64 | product_id | 2970 | 617 | 4.81

Вивід із small (39 504 позиції), рядки скорочено. На малому N гарячий ключ тоне в середньому, але з ростом N середній шард легшає, а товар 2410 (2047 позицій) нікуди не дівається. Десять найгарячіших товарів дають 22.8 % позицій: це закон Ципфа в датасеті. customer_id на 64 шардах теж перекошений (постійні клієнти), а order_id унікальний і розподіляє найрівніше.

Хеш розкидає різні ключі, але гарячий лишається одним, і весь його запис іде на один шард. Допомагає розщепити ключ: лічильник переглядів зберігають у k рядках (product_id, 0..k-1), пишуть у випадковий, а читання їх додає. Змінити ключ шардування пізніше означає переселити всю таблицю.

Перебалансування і кільце хешів

Section titled “Перебалансування і кільце хешів”

При переході з hash % 4 на hash % 5 рядок лишається на місці, лише якщо обидва залишки збіглися. Беремо 3500 товарів small і кільце з 64 віртуальними позиціями на вузол:

СпробуйPostgreSQLCtrl+Enter — виконати
products | moved_modulo | moved_ring | moved_elsewhere
----------+--------------+------------+-----------------
3500 | 2818 | 695 | 0

За модулем переїжджає 80.5 % товарів, за кільцем 19.9 %, і всі вони йдуть на новий вузол. Це консистентне хешування: вузли й ключі лежать в одному просторі хешів, замкненому в кільце, а ключ належить першому вузлу за годинниковою стрілкою. Новий вузол забирає лише дугу між собою й попереднім вузлом.

Кільце хешів: чотири вузли, ключі біля них; пʼятий вузол E забирає лише дугу між B і Eпростір хешів, замкнений у кільцеk1k2k3k4ABCDEКлюч іде за годинниковоюстрілкою до першого вузла.Вузол E зʼявився між B і C.Переїжджають лише ключіна виділеній дузі: k2 від C до E.k1, k3, k4 лишились на місці.hash % N: з N = 4 на N = 5міняють шард майже всі ключі.
Вузол E з'явився між B і C. Власника міняють лише ключі на виділеній дузі (тут k2).

Без віртуальних позицій чотири точки лягають на кільце нерівно; зі 64 товари розклалися по вузлах як 846, 1059, 788 і 807. Метод описали Карґер та співавтори (1997). Інший підхід: шардів набагато більше, ніж вузлів, і переселяють шарди цілком. Citus типово робить 32 шарди на таблицю (citus.shard_count), і змінюється лише «шард → вузол».

Запит з ключем шардування іде на один вузол. Без нього його розсилають на всі й збирають відповіді (scatter-gather). Затримку визначає найповільніша відповідь: якщо шард повільний в 1 % запитів, запит до 64 шардів повільний у 1 − 0.99^64 ≈ 47 % випадків (tail at scale, Dean і Barroso).

JOIN таблиць з різних шардів ще дорожчий: рядки однієї таблиці треба переслати мережею до рядків іншої. Тому дані, що з’єднуються, кладуть разом, і це колокація. Товари одного продавця в small:

SELECT round(avg(by_product), 2) AS avg_shards_by_product,
round(avg(by_seller), 2) AS avg_shards_by_seller
FROM (
SELECT count(DISTINCT abs(hashtextextended(product_id::text, 0)) % 8) AS by_product,
count(DISTINCT abs(hashtextextended(seller_id::text, 0)) % 8) AS by_seller
FROM products GROUP BY seller_id
) AS s;
avg_shards_by_product | avg_shards_by_seller
-----------------------+----------------------
6.30 | 1.00

За product_id каталог продавця розкиданий у середньому по 6.3 шардах із 8, за seller_id лежить на одному. Плата: найбільший продавець має 10.8 % товарів, тож ключ гарячіший. А замовлення з кількома позиціями за product_id лежать на кількох шардах у 92.5 % випадків (8882 із 9606).

Розподілені транзакції: двофазний коміт

Section titled “Розподілені транзакції: двофазний коміт”

Покупка зачіпає вузли A і B, які не ділять ні пам’яті, ні журналу, а завершитися має на обох або ні на одному. Протокол двофазного коміту (2PC, two-phase commit) веде координатор у два кроки. На першому кожен учасник виконує роботу, записує стан на диск, тримає блокування й відповідає «готовий»: відкотитися сам він уже не може. На другому координатор розсилає коміт, якщо всі готові, або відкат.

Двофазний коміт: обидва шарди відповіли «готовий», координатор зник, і рядки лишаються заблокованимикоординаторшард A: залишкишард B: замовленняфаза 1: підготовкаPREPARE TRANSACTIONPREPARE TRANSACTIONготовийготовийфаза 2: рішеннязникCOMMIT PREPARED не надійшоврядок заблоковано,вирішити сам шард не може
Обидва шарди відповіли «готовий», координатор зник до другої фази. Виділено очікування: шард не може ні закомітити, ні відкотити.

Тому 2PC не люблять. Між фазами учасник тримає блокування й чекає, а розблокувати його може лише людина чи скрипт. Коміт потребує живих усіх учасників і координатора, тож із кожним шардом імовірність відмови росте. До того ж він коштує двох кіл обміну та двох fsync на учасника.

Тому шардують так, щоб транзакція лишалась на одному вузлі (колокація), або замінюють її локальними транзакціями з компенсаціями (saga). Spanner зберігає 2PC, але робить координатора й учасників групами консенсусу, які не зникають разом.

Що означає «консистентність» у розподіленій системі

Section titled “Що означає «консистентність» у розподіленій системі”

У модулі 10 консистентність була про правила бази (CHECK, FOREIGN KEY). У розподіленій системі це ще й те, що бачать читачі різних копій. Спершу три сцени, назви потім.

Сцена 1. Залишок товару лежить на вузлах A і B. Покупець купив останню одиницю: A показує 0, і йому відповіли «успіх». Через секунду запит іншого покупця потрапляє на B, а той ще показує 1. Чи може читання після завершеного запису повернути старе? Якщо ні, система поводиться так, ніби копія одна. Це лінеаризованість (linearizability): кожна операція діє миттєво десь між початком і відповіддю, а порядок збігається з реальним часом.

Сцена 2. Покупка змінює залишок і створює замовлення, а звіт одночасно читає обидві таблиці. Чи може він побачити залишок зменшеним, а замовлення ще ні? Транзакції мають давати такий результат, ніби йшли по черзі. Це серіалізованість (serializability): вона про кілька об’єктів і не вимагає збігу з реальним часом. Порушує її, наприклад, write skew з модуля 10.

Сцена 3. Між A і B пропала мережа, а покупець зайшов на B. Вузол або відповість старим, або відмовить до відновлення зв’язку, і разом обидва варіанти неможливі: це теорема CAP, до неї дійдемо нижче.

Це різні властивості, і система може мати одну без іншої. Разом вони дають строгу серіалізованість, як у Spanner (нижче).

Моделі консистентності

Section titled “Моделі консистентності”

Шкала моделей від сильної до слабкої, на яку спирається частина IV:

Модель Гарантія «Крамниця»
лінеаризованість система поводиться так, ніби копія одна; порядок збігається з реальним часом після «успіху» покупки жодне читання не покаже старий залишок
послідовна усі бачать один порядок операцій, що поважає порядок кожного клієнта, але може відставати від реального часу двоє бачать стрічку подій однаково, хоч і з запізненням
causal consistency якщо одна операція могла вплинути на іншу, усі бачать їх у цьому порядку відповідь продавця не з’являється раніше відгуку
монотонні читання клієнт не бачить минуле після теперішнього оновили сторінку, а лічильник переглядів зменшився: порушення
read-your-writes клієнт бачить власні записи (модуль 12) автор одразу бачить свій відгук
eventual consistency якщо записи припинились, усі копії врешті збігаються лічильник переглядів

Eventual consistency обіцяє мало: колись копії зійдуться, а що бачить читач до того, не сказано. Лічильнику переглядів цього досить, залишку товару мало: двоє покупців одночасно прочитають «1». Конфліктним записам потрібне правило злиття, а найпростіше, «виграє останній», мовчки губить один із них.

Сцену 3 описує теорема CAP (гіпотеза Брюера 2000 року, доведена Гілбертом і Лінчем 2002-го): у мережі, що може губити повідомлення, не можна одночасно гарантувати (C) лінеаризованість і (A) відповідь на кожен запит від кожного живого вузла. C тут означає лінеаризованість, а не консистентність з ACID. A означає відповідь кожного живого вузла, а не «здебільшого доступно», під що багато «AP-систем» не підпадають.

«Обери два з трьох» вводить в оману, бо розриву мережі не обирають: він буває. На 6000 км лише поширення сигналу дає 60 мс туди й назад (модуль 2 курсу мереж). Вибір виникає лише під час розділення: відмовити частині клієнтів і зберегти лінеаризованість або відповісти тим, що є. Клеппман зауважує, що ярлики «CP» і «AP» грубі: одна система має різні гарантії для різних операцій.

CAP доповнює PACELC (Абаді, 2012). Без розділення система все одно вибирає між затримкою (latency) і консистентністю. Лінеаризований запис чекає реплік, а швидше за відстань вони не відповідять: між регіонами з 60 мс туди й назад кожен синхронний коміт коштує не менше. Цю ціну видно в тарифах: Cosmos DB продає кілька рівнів консистентності, а сильно консистентне читання DynamoDB вдвічі дорожче (хмарний курс, модуль 7).

Вибір лідера, журнал і склад кластера вимагають, щоб вузли домовились про одне значення. Це задача консенсусу. Вона розв’язує split brain з модуля 12 правилом більшості: дві більшості з n вузлів завжди перетинаються, тож два лідери одного терму неможливі.

Raft (Онгаро й Остерхаут, 2014) має три частини. Терми: пронумеровані періоди, кожен починається з виборів і має не більше одного лідера; вузол, що бачить вищий терм, підкоряється. Вибори: репліка, що не чула лідера довше за випадковий тайм-аут (у статті 150–300 мс), збільшує терм і просить голосів, перший голос за себе. Вузол голосує в терму раз, перемагає більшість. Журнал: запис закомічено, коли його має більшість. Кандидат із відсталим журналом голосів не отримає, тож новий лідер має всі закомічені записи.

Raft: лідер терму 3 і закомічений запис: він у журналах більшості вузлівS3 · лідерзапис x = 5S1 · репліканедоступнийS2 · реплікамає записS4 · реплікамає записS5 · репліказапису ще немаєЗ пʼяти вузлів запис має троє: лідер і дві репліки.Це більшість, запис закомічено. Вузол без звʼязку наздожене згодом.терми: кожен починається з виборівтерм 1: лідер S1терм 2: голоси розділились,лідера немаєтерм 3: лідер S3S1 зник
Лідер S3 терму 3. Запис закомічено: він у трьох журналах із п'яти. Унизу терми; у терму 2 голоси розділились.

Кластер із 2f + 1 вузлів переживає відмову f: три вузли зносять один, п’ять двох. Кожен коміт чекає більшості. Старий лідер, відрізаний від більшості, нічого не закомітить, а побачивши вищий терм, підкориться. На консенсусних сховищах (etcd, Consul, ZooKeeper) тримається автоматичне перемикання з модуля 12.

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

Spanner (Corbett та ін., 2012) для лінеаризованості між центрами потребує спільного годинника. TrueTime повертає не момент, а інтервал, що гарантовано містить справжній час: джерела GPS і атомні годинники, невизначеність ε за статтею від 1 до 7 мс, у середньому 4 мс. Транзакція отримує мітку часу коміту й чекає, поки вона гарантовано мине (commit wait), і лише тоді показує результат. Тож транзакція, що почалась після чужого коміту, має більшу мітку: це зовнішня консистентність, за статтею те саме, що лінеаризованість. Ціна за мікротестами статті: близько 5 мс очікування й близько 9 мс на Paxos при одній репліці, а 2PC між учасниками коштував у середньому 17.0 мс з одним учасником, 24.5 з двома і 31.5 мс з п’ятьма.

CockroachDB і YugabyteDB обходяться без GPS. За документацією, Cockroach виконує транзакції на SERIALIZABLE, користується гібридними логічними годинниками й вимикає вузол, що розійшовся з годинниками половини кластера на 80 % допустимого зсуву. Запис у діапазон комітить більшість реплік через Raft. Yugabyte ділить таблиці на tablet-и, що реплікуються через Raft. Його SQL-шар YSQL побудовано на коді планувальника й виконавця PostgreSQL, а зберігання замінено власним.

Платять усі однаково: коміт чекає більшості реплік, тож затримка запису не менша за відстань до найближчої більшості. Ці бази виграють, коли потрібні транзакції й масштаб разом; на простому навантаженні дешевший PostgreSQL.

Citus лишається розширенням PostgreSQL: create_distributed_table('orders', 'customer_id') розподіляє таблицю за колонкою, а запис через кілька воркерів комітиться двофазно, тим самим PREPARE TRANSACTION.

medium піднімають так: DATASET=medium docker compose --profile postgres up -d --wait. Числа партиціювання з нього (610 тисяч подій, 25 партицій, PostgreSQL 18.6); березень 2025 має 27 148 подій:

-- view_events (без партицій): Buffers: shared hit=6784 Execution Time: 12.0 ms
-- view_events_p (одна партиція): Buffers: shared hit=300 Execution Time: 2.4 ms

Видалення того ж березня з копії таблиці проти DETACH і DROP партиції (medium, один прогін, порядок величин). DELETE 27 148 рядків: 27.5 мс, 1489 кБ WAL, файл лишився 66 МБ. DETACH: 1.0 мс і 8960 байтів WAL, DROP: 1.7 мс і 6000 байтів. Мертві версії після DELETE чекають VACUUM.

Тепер 2PC. Типово max_prepared_transactions = 0, і PREPARE TRANSACTION відповідає prepared transactions are disabled. Після ALTER SYSTEM SET max_prepared_transactions = 10 і перезапуску змоделюємо два шарди двома базами одного сервера: shop_small із залишками й shop_orders із order_log. gid унікальний на весь сервер, тож у кожного учасника свій.

-- шард A (shop_small)
BEGIN;
UPDATE inventory SET quantity = quantity - 1 WHERE product_id = 39 AND quantity > 0;
PREPARE TRANSACTION 'buy_1001_stock';
-- шард B (shop_orders)
BEGIN;
INSERT INTO order_log VALUES (1001, 39, 1);
PREPARE TRANSACTION 'buy_1001_orders';
-- COMMIT PREPARED не надходить

У пісочниці цього не відтворити: там немає ні підготовлених транзакцій, ні двох баз. З іншої сесії на шарді A:

SELECT gid, database FROM pg_prepared_xacts ORDER BY gid;
SELECT quantity FROM inventory WHERE product_id = 39;
SET lock_timeout = '2s';
UPDATE inventory SET quantity = quantity - 1 WHERE product_id = 39 AND quantity > 0;
gid | database
-----------------+-------------
buy_1001_orders | shop_orders
buy_1001_stock | shop_small
quantity
----------
2
ERROR: canceling statement due to lock timeout
CONTEXT: while updating tuple (0,39) in relation "inventory"

Читання проходить і бачить стару версію: MVCC читачів не блокує. UPDATE без lock_timeout висів би з wait_event = transactionid. Підготовлена транзакція пережила розрив сесії й docker restart: таймера в PostgreSQL немає. Закрити її можна з будь-якої сесії. Вона ще й тримає горизонт xmin: після UPDATE 3499 рядків VACUUM (VERBOSE) написав 3499 are dead but not yet removable, а після COMMIT PREPARED для обох 3500 removed, 0 are dead but not yet removable.

Jepsen запускає систему в кластері, випадкові клієнти виконують операції, тестувальник рве мережу й зупиняє вузли, а записану історію перевіряє програма проти формальної моделі (для транзакцій це Elle, що шукає цикли в графі залежностей). Метод знаходить помилки, але не доводить їх відсутності.

Звіт від 15 травня 2020 року порівнює обіцянку з виміром: MongoDB заявляла «одні з найсильніших гарантій консистентності, коректності й безпеки серед усіх сучасних баз» і «повні ACID-транзакції». Тестували кластер із дев’яти вузлів (два шарди по три й три вузли метаданих) з розривами мережі, націленими на мастерів.

  • Типово запис підтверджує один вузол (w: 1), тож при перемиканні мастера підтверджений запис може зникнути, а читання local бачить незакомічені дані. За звітом, близько 80 % користувачів хмарної версії лишають типовий write concern і 99.6 % типовий read concern.
  • Усередині транзакції рівні, задані на базі чи колекції, ігноруються: діє рівень транзакції, типово знову local і w: 1. Рівень snapshot не гарантує знімка без write concern: majority, навіть для транзакцій, що лише читають.
  • На найсуворіших налаштуваннях бачили асиметрію читання (read skew), циклічний потік інформації, подвоєні записи й транзакцію, що читала власні записи з майбутнього. В одній історії близько 100 транзакцій на секунду 1461 із 13 914 мали циклічні залежності. Snapshot isolation такі цикли (write skew, модуль 10) дозволяє, але «повною ізоляцією» це не назвеш.
  • Через 11 днів MongoDB знайшла в повторах транзакцій помилку, яку вважає причиною аномалій; виправлення планувалось у 4.2.8.

Підсумок: лінеаризованість окремого документа була, але лише з majority і читанням linearizable, а транзакції робили менше, ніж обіцяла реклама. «ACID» у рекламі, snapshot isolation у документації й типова поведінка — три різні речі. Пастка не лише нереляційна: за звітом про PostgreSQL 12.3 (12 червня 2020) SERIALIZABLE на одному вузлі пропускав аномалію G2-item через помилку у виявленні конфліктів SSI; її виправили в 12.4 того ж року.

Типові помилки розуміння

Section titled “Типові помилки розуміння”

«CAP: обери два з трьох». Розрив мережі не обирають. Вибір є лише під час розділення: відмовити чи відповісти застарілим. Без розділення лишається вибір між затримкою й консистентністю (PACELC).

«Лінеаризованість і серіалізованість — те саме». Перша про свіжість копій одного об’єкта, друга про еквівалентний порядок транзакцій. Обидві разом дає Spanner.

«Хеш розв’язує гарячі ключі». Він розкидає різні ключі, а один гарячий лишається на одному шарді: на small 4.81 від середнього при 64 шардах.

«2PC гарантує завершення». Він гарантує атомарність. Якщо координатор зникне між фазами, учасники чекають із блокуваннями.

Перевір себе

1. Кластер із `hash % 4` розширили до `hash % 5`. Яка частка ключів переїде, і що змінює кільце хешів?
2. Після `PREPARE TRANSACTION` сесію розірвано, а координатор зник. Що з рядками, які змінила транзакція?
3. Після «успішної» покупки залишок на одній репліці 0, але читач на іншій репліці бачить 1. Яку модель порушено?
4. Кластер Raft із п’яти вузлів, два недоступні. Що із записами?
5. За звітом Jepsen, що з читанням `snapshot` у транзакції MongoDB 4.2.6, яка комітиться без `write concern: majority`?

Лабораторної до модуля немає. Суміжні теми відпрацьовують L7 (реплікація, лаг, перемикання мастера) і L9 (одна предметна область у трьох моделях).