Шардинг — горизонтальное масштабирование БД: данные разбиваются на части (шарды) и размещаются на отдельных серверах. Один сервер не справляется с миллиардами строк — шардинг позволяет распределить нагрузку.
Стратегии разбиения
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 при добавлении нового сервера — болезненная операция