Шардинг — горизонтальне масштабування БД: дані розбиваються на частини (шарди) і розміщуються на окремих серверах. Один сервер не справляється з мільярдами рядків — шардинг дозволяє розподілити навантаження.
Стратегії розбиття
Range-based — шард 1: id 1–1M, шард 2: id 1M–2M. Просто, але нерівномірне навантаження
Hash-based — shard = hash(id) % N. Рівномірний розподіл, але складний resharding
Directory-based — окремий сервіс-каталог знає, де який запис. Гнучко, але додатковий hop
Складнощі
JOIN між шардами неможливий або дуже дорогий
Транзакції через кілька шардів потребують distributed transaction протоколів (2PC)
Resharding при додаванні нового сервера — болюча операція