Пошук уроків, статей та іншого контенту
Розберете горизонтальне розподілення даних між вузлами, вибір shard key, маршрутизацію запитів і проблеми ребалансування.
Шардінг — це горизонтальне розподілення даних між кількома вузлами — шардами.
Кожен шард зберігає лише частину повного набору даних, але для клієнта система може виглядати як одна база даних.
Наприклад, замість одного вузла з усіма замовленнями:
orders
├── shard-1: user_id 1–1 000 000
├── shard-2: user_id 1 000 001–2 000 000
└── shard-3: user_id 2 000 001–3 000 000Шардінг допомагає:
збільшити обсяг даних, який може зберігати система;
розподілити навантаження на запис і читання;
ізолювати частину навантаження або даних на окремих вузлах.
Шардінг відрізняється від реплікації:
реплікація зберігає копії одних і тих самих даних на різних вузлах;
шардінг розподіляє різні частини даних між вузлами.
Ці підходи часто використовують разом: кожен шард може мати власні репліки для відмовостійкості та масштабування читання.
Типова система з шардінгом складається з таких компонентів:
Клієнт або сервіс, який надсилає запит.
Роутер, який визначає потрібний шард.
Шарди, що зберігають частини даних.
Конфігурація топології, яка описує доступні шарди та правила маршрутизації.
Client
|
v
Router
|
+--> Shard A
+--> Shard B
+--> Shard CРоутер може бути:
окремим сервісом;
частиною application-сервісу;
компонентом драйвера або проксі перед базою даних.
Важливо, щоб усі компоненти однаково обчислювали місце розташування запису. Якщо запис потрапляє на один шард, а читання шукає його на іншому, система втрачає коректність.
Shard key — це поле або набір полів, за якими система визначає шард для запису.
Наприклад:
user_id = 4815
shard = hash(user_id) mod кількість_шардівВід вибору shard key залежать:
рівномірність розподілення даних;
рівномірність навантаження;
кількість запитів до одного шарда;
можливість виконувати запити без звернення до всіх вузлів;
складність перенесення даних під час масштабування.
Хороший ключ має:
Високу кардинальність
Він повинен мати багато різних значень. Наприклад, user_id зазвичай кращий за country, оскільки країн мало.
Рівномірний розподіл
Значення не повинні концентруватися на одному шарді.
Стабільність
Значення ключа не повинно часто змінюватися. Якщо ключ змінюється, запис може знадобитися перенести на інший шард.
Відповідність основним запитам
Якщо більшість запитів виконується для конкретного користувача, user_id може бути хорошим ключем.
Захист від hot shard
Один шард не повинен отримувати непропорційно багато записів або читань.
Для сервісу замовлень:
user_idПереваги:
усі замовлення користувача можуть бути на одному шарді;
запити історії замовлень легко маршрутизувати;
операції в межах одного користувача простіше виконувати атомарно.
Для платформи з подіями:
tenant_idПереваги:
дані одного клієнта ізольовані;
запити конкретного клієнта не потребують fan-out.
Недолік — великий клієнт може створити hot shard, якщо всі дані одного tenant зберігаються разом.
Значення shard key діляться на діапазони:
shard-1: user_id < 1 000 000
shard-2: 1 000 000 <= user_id < 2 000 000
shard-3: user_id >= 2 000 000Переваги:
легко виконувати запити за діапазоном;
зрозуміло, де розташовані дані;
зручно переносити окремі діапазони.
Недоліки:
нерівномірний розподіл;
нові значення можуть постійно потрапляти в один діапазон;
можливі hot shard.
Наприклад, якщо ключем є час створення, усі нові записи спочатку потраплятимуть до найновішого діапазону.
Спочатку обчислюється хеш ключа, а потім визначається шард:
shard = hash(shard_key) mod NПереваги:
зазвичай рівномірніший розподіл;
зменшується ризик hot shard через послідовні значення;
простий маршрут для точкового запиту.
Недоліки:
діапазонні запити стають складними;
зміна кількості шардів може перемістити багато ключів;
запит без shard key може потребувати звернення до всіх шардів.
Окрема таблиця або сервіс зберігає відповідність між ключем і шардом:
tenant-a -> shard-1
tenant-b -> shard-3
tenant-c -> shard-2Переваги:
можна переміщувати окремі ключі без зміни алгоритму хешування;
зручно підтримувати нерівномірні розміри клієнтів;
маршрутизація може враховувати додаткові правила.
Недоліки:
директорія стає критично важливим компонентом;
її потрібно кешувати, реплікувати та оновлювати узгоджено;
помилка або застарілі дані директорії можуть спричинити неправильну маршрутизацію.
Іноді одного поля недостатньо. Тоді використовують складений ключ:
(tenant_id, user_id)Важливий порядок полів. Він впливає на:
спосіб хешування;
локальність даних;
можливість маршрутизувати запити;
рівномірність розподілення.
Наприклад, для багатоклієнтської системи можна спочатку враховувати tenant_id, а потім user_id. Але це не гарантує захисту від великого tenant, якщо всі його дані мають потрапляти на один шард.
Якщо запит містить shard key, роутер може звернутися до одного шарда:
GET /users/4815/orders
|
v
hash(4815) -> shard-2Це найкращий випадок:
мінімальна кількість мережевих операцій;
немає об’єднання результатів;
простіше контролювати затримку;
легше виконувати транзакційні операції.
Наприклад:
SELECT * FROM orders
WHERE status = 'pending';Якщо status не є shard key і немає додаткового індексу або каталогу, роутер може бути змушений виконати запит на всіх шардах.
Такий підхід називають scatter-gather:
запит розсилається на всі шарди;
кожен шард повертає локальний результат;
роутер об’єднує результати;
роутер сортує, обмежує або агрегує загальну відповідь.
+--> Shard A --+
Client -> Router -> Shard B ---+-> merged result
+--> Shard C --+Проблеми scatter-gather:
затримка залежить від найповільнішого шарда;
кількість запитів зростає разом із кількістю шардів;
складніше виконувати ORDER BY, LIMIT, COUNT і пагінацію;
відмова одного шарда може зірвати весь запит.
Тому API та схему даних бажано проєктувати так, щоб основні запити містили shard key.
Нижче наведено спрощену модель хеш-шардінгу. Вона не є реалізацією розподіленої бази даних, але показує головний принцип: однаковий ключ стабільно маршрутизується до одного шарда.
const shards = [
new Map(),
new Map(),
new Map(),
];
function hashKey(key) {
const value = String(key);
let hash = 2166136261;
for (let i = 0; i < value.length; i += 1) {
hash ^= value.charCodeAt(i);
hash = Math.imul(hash, 16777619);
}
return hash >>> 0;
}
function getShardIndex(userId) {
return hashKey(userId) % shards.length;
}
function saveOrder(order) {
const shardIndex = getShardIndex(order.userId);
shards[shardIndex].set(order.id, order);
return shardIndex;
}
function findOrder(orderId, userId) {
const shardIndex = getShardIndex(userId);
return shards[shardIndex].get(orderId) ?? null;
}
const orders = [
{ id: "order-1", userId: 101, total: 49.99 },
{ id: "order-2", userId: 202, total: 120.00 },
{ id: "order-3", userId: 101, total: 15.50 },
];
for (const order of orders) {
const shardIndex = saveOrder(order);
console.log(`${order.id} записано на shard-${shardIndex}`);
}
const order = findOrder("order-3", 101);
console.log("Знайдене замовлення:", order);
console.log(
"Кількість записів на шардах:",
shards.map((shard) => shard.size),
);У реальній системі замість Map будуть окремі бази даних або вузли. Також потрібно зберігати конфігурацію топології, обробляти недоступність вузлів і контролювати зміни кількості шардів.
Навіть якщо кількість записів розподілена рівномірно, навантаження може бути нерівномірним.
Наприклад:
один популярний користувач отримує більшість запитів;
один tenant має набагато більше даних за інших;
один ключ використовується для великої кількості одночасних записів.
Такий ключ називають hot key, а перевантажений вузол — hot shard.
Можливі рішення:
додати випадковий суфікс до ключа та розподілити записи між кількома логічними ключами;
використовувати складений ключ;
розподілити великого tenant між кількома шардами;
додати кеш для дуже популярних читань;
обмежити або ізолювати навантаження великих клієнтів.
Додавання суфікса ускладнює читання: щоб знайти всі записи, потрібно знати можливі суфікси або виконувати кілька запитів.
Ребалансування — це перенесення даних між шардами, коли змінюється їхня кількість або розмір.
Причини ребалансування:
додавання нового шарда;
видалення або заміна вузла;
нерівномірний розподіл даних;
зростання окремих tenant або діапазонів;
зміна shard key.
Ребалансування не повинно надовго блокувати запис і читання. Зазвичай його виконують поступово:
визначають частину даних для перенесення;
копіюють дані на новий шард;
доганяють зміни, які відбулися під час копіювання;
короткочасно синхронізують або блокують перемикання;
оновлюють маршрутизацію;
перевіряють новий шард;
видаляють стару копію після підтвердження.
Під час перенесення важливо не допустити:
втрати записів;
дублювання записів;
читання застарілої копії;
маршрутизації частини запитів до старого місця;
перевищення допустимого навантаження на джерело або призначення.
hash(key) mod NПроста схема:
shard = hash(key) % Nмає суттєвий недолік. Якщо кількість шардів змінюється з N на N + 1, більшість ключів отримає новий результат.
Наприклад:
hash(user_id) % 3
hash(user_id) % 4Це означає, що велика частина даних логічно «переїхала» навіть до фізичного перенесення.
Для зменшення обсягу переміщень використовують:
узгоджене хешування;
віртуальні вузли;
логічні партиції;
directory-based маршрутизацію.
Замість одного великого діапазону для кожного фізичного шарда створюють багато малих логічних партицій.
partition-1 -> shard-A
partition-2 -> shard-C
partition-3 -> shard-A
partition-4 -> shard-BПід час додавання нового вузла переміщують лише частину партицій, а не всю структуру ключів.
Переваги:
менший обсяг даних для одного кроку міграції;
поступове ребалансування;
простіше виводити вузол з експлуатації;
кращий контроль навантаження.
Кількість віртуальних партицій потрібно вибирати з урахуванням обсягу даних і вартості керування ними. Надто мало партицій обмежує точність ребалансування, а надто багато збільшує службові витрати.
Під час переміщення одного логічного діапазону можуть одночасно існувати дві копії:
старий шард: запис
новий шард: записСистема повинна мати чіткий момент зміни власника даних.
Поширений підхід:
старий шард залишається джерелом істини;
дані копіюються на новий шард;
зміни після початкового копіювання передаються окремо;
після синхронізації оновлюється маршрутизація;
новий шард стає власником;
старий шард більше не приймає записи для цього діапазону.
Подвійний запис може допомогти під час міграції, але він створює нові ризики:
запис успішний лише на одному шарді;
повторна доставка створює дублікати;
операції виконуються в різному порядку;
помилки потребують повторної синхронізації.
Тому записи під час міграції мають бути ідемпотентними, а процес повинен мати перевірку повноти та узгодженості даних.
Шардінг ускладнює транзакції, якщо одна операція зачіпає кілька шардів.
Наприклад, переказ коштів між двома користувачами може вимагати зміни на різних шардах:
user-A -> shard-1
user-B -> shard-3Локальна транзакція одного шарда вже не охоплює обидві зміни.
Транзакції між шардами:
мають більшу затримку;
складніші в реалізації;
можуть блокувати ресурси на кількох вузлах;
потребують складнішого відновлення після помилок.
Практичний підхід — проєктувати shard key так, щоб найважливіші атомарні операції виконувалися в межах одного шарда. Наприклад, усі дані одного замовлення можна зберігати за order_id, а всі зміни його стану виконувати локально.
Якщо операція все одно має охоплювати кілька шардів, застосовують:
саги з компенсувальними діями;
стан операції в окремій сутності;
ідемпотентні команди;
повторні спроби;
журнал подій для відновлення.
Пагінація на одному шарді проста: база повертає наступну сторінку локального набору.
Для запиту на всіх шардах потрібно:
виконати локальний запит на кожному шарді;
отримати достатню кількість кандидатів;
об’єднати результати;
відсортувати їх;
повернути потрібну сторінку.
LIMIT 20 не завжди означає, що достатньо отримати по 20 записів із кожного шарда. Якщо сортування глобальне, на одному шарді можуть знаходитися всі перші 20 записів.
У розподілених системах часто використовують cursor-based pagination із детермінованим сортуванням. Але курсор має містити стан, необхідний для продовження запиту на всіх задіяних шардах.
Перед вибором ключа:
Випишіть основні запити.
Для кожного запиту визначте, чи містить він можливий ключ.
Оцініть кількість значень і їхній розподіл.
Перевірте, чи змінюється ключ.
Визначте максимальний очікуваний розмір одного логічного власника: користувача, tenant або пристрою.
Оцініть операції, які повинні бути атомарними.
Перевірте сценарій додавання та видалення шардів.
Окремо перевірте найгірший випадок, а не лише середнє навантаження.
Ознака поганого вибору — коли основні запити не містять shard key і майже кожен запит перетворюється на scatter-gather.
Для системи з шардінгом недостатньо стежити лише за загальним навантаженням. Потрібні метрики для кожного шарда:
обсяг даних;
кількість записів і читань;
latency за перцентилями;
кількість помилок;
використання CPU, пам’яті та диска;
довжина черг;
кількість scatter-gather запитів;
відсоток hot key;
прогрес міграцій;
розбіжності між копіями під час перенесення.
Важливо також логувати:
обчислений shard key;
вибраний шард;
версію конфігурації маршрутизації;
ідентифікатор міграції;
причину повторної спроби.
Вибір ключа лише за кардинальністю. Велика кількість значень не гарантує рівномірного навантаження.
Використання монотонного ключа в діапазонній схемі. Нові записи можуть створити один hot shard.
Маршрутизація за полем, якого немає в запитах. У результаті більшість запитів стає scatter-gather.
Ігнорування великих tenant або користувачів. Середній розподіл може виглядати добре, але один клієнт перевантажить шард.
Припущення, що додавання шарда автоматично розподілить дані. Потрібні міграція, контроль прогресу та оновлення маршрутизації.
Зміна N у hash(key) mod N без плану міграції. Велика кількість ключів змінить своє логічне розташування.
Подвійний запис без ідемпотентності. Повторні спроби можуть створити дублікати або різні версії даних.
Відсутність плану відновлення. Потрібно знати, що робити при зупинці міграції, недоступності шарда або частковому успіху запису.
Оцінювання лише середньої latency. Scatter-gather і hot shard часто проявляються через високі p95 та p99.
Шардінг розподіляє дані горизонтально між кількома вузлами.
Shard key визначає розташування запису та впливає на баланс, маршрутизацію і транзакції.
Хеш-шардінг зазвичай дає рівномірний розподіл, а діапазонний — кращу підтримку range-запитів.
Запити із shard key можна спрямувати на один шард.
Запити без shard key часто стають scatter-gather і мають вищу latency.
Hot shard виникає через нерівномірний розподіл даних або навантаження.
Ребалансування потрібно виконувати поступово, контролюючи копіювання, зміни та перемикання маршрутизації.
Віртуальні партиції та directory-based маршрутизація спрощують міграцію порівняно з простим hash(key) mod N.
Найкращий shard key відповідає реальним шаблонам запитів і дозволяє локалізувати найважливіші операції.