Когда данных становится много, один сервер перестаёт справляться — медленно читает, медленно пишет, и если упадёт, приложение встанет. MongoDB решает это двумя механизмами: репликация хранит копии данных на нескольких серверах (надёжность), шардинг разрезает коллекцию между серверами (горизонтальное масштабирование).
Oplog — журнал операций primary: каждый secondary читает его со своей позиции и переигрывает у себя. Отставший узел держит состояние постарше, поэтому на выборах primary становится тот, кто прочитал журнал дальше. Пока идут выборы, записей не принимает никто.
Почему нельзя просто держать один сервер
Интернет-магазин с миллионами товаров на одном сервере:
- Если сервер упадёт — сайт недоступен, продажи стоят.
- Если данных больше, чем помещается в оперативную память — каждый запрос идёт на диск, скорость падает в десятки раз.
- Нагрузка на запись ограничена скоростью одного диска.
Репликация решает первую проблему, шардинг — вторую и третью.
Replica set — резервные копии в реальном времени
Что делать, когда единственный сервер упал, у MongoDB отвечает replica set: группа серверов, которые хранят одни и те же данные. Минимальная конфигурация — три узла: один primary принимает все записи, два secondary автоматически копируют изменения.
Почему три, а не два? Потому что при сбое узлы голосуют, кто станет новым primary. Нужно большинство — при двух серверах одного голоса недостаточно. Три узла дают кворум.
Встречается и третий вид узла — арбитр: он данных не хранит, только голосует. Его ставят, когда на третий полноценный сервер жалко денег, но у такого набора есть неприятное свойство: узлов с данными остаётся два, и большинство из них при падении одного уже не набирается — запись с w: "majority" в этот момент просто перестаёт проходить. Поэтому арбитра стараются не заводить вовсе.
Как работает oplog
Механизм репликации построен на oplog — специальной коллекции в базе local, куда primary записывает каждую успешную операцию. Secondary постоянно читают этот журнал и воспроизводят операции у себя.
Если secondary отстал и снова подключился — он догоняет по oplog с той позиции, на которой остановился.
Oplog имеет фиксированный размер (по умолчанию около 5% свободного диска). Если secondary отстал сильно и нужных записей в oplog уже нет — придётся копировать весь массив данных заново (initial sync). Поэтому при большой нагрузке oplog стоит увеличивать.
Что происходит при падении primary
Secondary замечают, что primary молчит, когда истекает electionTimeoutMillis — по умолчанию это 10 секунд. Дальше начинается голосование, и побеждает тот, у кого самый свежий oplog. При настройках по умолчанию медианное время до появления нового primary не превышает 12 секунд — всё это время кластер не принимает записи.
Вторая сторона переключения: записи, которые старый primary успел принять, но не успел передать большинству узлов, при его возвращении в кластер откатываются. Именно от этого страхует w: "majority": с ним клиент получает подтверждение только после того, как запись дошла до большинства.
Важное правило: если из трёх узлов два недоступны — оставшийся один переходит в режим только для чтения. Это защита от split brain: если бы оба «выживших» сегмента продолжали принимать записи независимо, данные разошлись бы.
Куда идут запросы на чтение
По умолчанию все запросы идут на primary. Но можно направить чтение на secondary — чтобы разгрузить primary или читать из ближайшего географически узла. Это называется read preference, и режимов у него пять:
| Режим | Откуда читаем | Когда использовать |
|---|---|---|
primary (по умолчанию) | Только primary | Когда нужны самые актуальные данные |
primaryPreferred | Primary, при его недоступности — secondary | Лучше отдать чуть устаревшее, чем ошибку на время выборов |
secondary | Только secondary | Аналитика, отчёты — не нагружать primary |
secondaryPreferred | Secondary, при недоступности — primary | Нагрузка на чтение, данные могут немного отставать |
nearest | Узел с минимальной задержкой | Распределённые кластеры в разных регионах |
Подводный камень: secondary может немного отставать от primary. Если вы только что записали данные и сразу читаете — используйте primary, иначе можете прочитать устаревшее.
db.product.find({ categoryId: 1 })
.readPref("secondaryPreferred");
Когда одной replica set уже мало
Replica set решает проблему надёжности, но не масштабирования: каждый узел хранит полную копию, поэтому в память одного сервера данные по-прежнему должны помещаться.
Практический ориентир: если активный набор данных (те документы, к которым обращаются чаще всего) занимает больше 70% памяти — пора думать о следующем шаге.
Следующий шаг в MongoDB — sharded cluster: данные физически разрезаются между несколькими replica set.
Sharded cluster — как устроено изнутри
Данные перестали помещаться на один replica set, и их режут на части по разным машинам. Чтобы приложение этого не заметило, в шардированном кластере работают четыре типа компонентов, у каждого своя работа:
Клиент не знает про шарды: он говорит с роутером, а тот смотрит в конфигурационные серверы и решает, куда идти. Поэтому потеря конфигурационных серверов страшнее потери одного шарда.
- mongos — точка входа. Клиент подключается к mongos, не зная о шардах. Mongos смотрит в метаданные и направляет запрос на нужный шард; данных он не хранит, поэтому его запускают в нескольких экземплярах.
- Config servers — три узла, которые хранят карту: какие данные на каком шарде. Без них кластер не работает.
- Shards — обычные replica set, каждый хранит свою часть данных.
- Balancer — фоновый процесс, который следит за равномерностью. Если один шард перегружен — перемещает часть данных на другой.
Данные внутри коллекции делятся на chunks — куски по 128 МБ по умолчанию. Balancer перемещает их между шардами так, чтобы выровнялся объём данных, а не число кусков.
Важно понимать цену: минимальный боевой кластер — это 3 узла config servers + 2 шарда по 3 реплики + 2 mongos = 11 серверов. Это заметно дороже одной replica set, поэтому шардинг включают, только когда без него не обойтись.
Shard key — самое важное решение при шардинге
Shard key — поле (или несколько полей), по которому MongoDB решает, на какой шард положить документ. Это решение принимается один раз и менять его сложно — поэтому выбор shard key требует внимания.
Равномерное распределение против горячего шарда
Представьте: магазин шардирует коллекцию товаров по полю _id типа ObjectId. ObjectId начинается с отметки времени в секундах, и по ней он растёт (внутри одной секунды порядок задают случайные байты и счётчик). Для шардинга по диапазону этого хватает с лихвой: все товары, заведённые сегодня, попадают на один и тот же шард. Он перегружен, остальные простаивают. Это называется hot shard — главная ошибка при шардинге.
Ranged sharding — документы распределяются по диапазонам значений. Подходит когда значения равномерно распределены или когда запросы часто фильтруют по диапазону (например, по дате):
sh.shardCollection("shop.product", { categoryId: 1 });
Hashed sharding — MongoDB сама хеширует значение ключа, и соседние значения расходятся по разным шардам. Хорошо для монотонных ключей вроде ObjectId или timestamp. Минус: запросы по диапазону вынуждены опрашивать все шарды:
sh.shardCollection("shop.product", { _id: "hashed" });
Разница видна на счёте: разложим 900 подряд идущих идентификаторов по трём шардам — сначала по диапазону, потом по хешу:
живой пример
import java.util.Map;
import java.util.TreeMap;
public class ShardRouting {
static final String[] SHARDS = {"shardA", "shardB", "shardC"};
static String ranged(long objectId) {
if (objectId < 1_000_000_000L) return SHARDS[0];
if (objectId < 2_000_000_000L) return SHARDS[1];
return SHARDS[2];
}
static String hashed(long objectId) {
long mixed = objectId * 0x9E3779B97F4A7C15L;
return SHARDS[Math.floorMod(Long.hashCode(mixed), SHARDS.length)];
}
public static void main(String[] args) {
Map<String, Integer> byRange = new TreeMap<>();
Map<String, Integer> byHash = new TreeMap<>();
long nextId = 2_000_000_100L;
for (int i = 0; i < 900; i++) {
long objectId = nextId + i;
byRange.merge(ranged(objectId), 1, Integer::sum);
byHash.merge(hashed(objectId), 1, Integer::sum);
}
System.out.println("ranged: " + byRange);
System.out.println("hashed: " + byHash);
}
}
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
живой пример
package main
import (
"fmt"
"maps"
"slices"
"strings"
)
var shards = []string{"shardA", "shardB", "shardC"}
func ranged(objectID int64) string {
if objectID < 1_000_000_000 {
return shards[0]
}
if objectID < 2_000_000_000 {
return shards[1]
}
return shards[2]
}
func hashed(objectID int64) string {
mixed := uint64(objectID) * 0x9E3779B97F4A7C15
h := int(int32(uint32(mixed ^ (mixed >> 32))))
return shards[((h%len(shards))+len(shards))%len(shards)]
}
func show(counts map[string]int) string {
var parts []string
for _, shard := range slices.Sorted(maps.Keys(counts)) {
parts = append(parts, fmt.Sprintf("%s=%d", shard, counts[shard]))
}
return "{" + strings.Join(parts, ", ") + "}"
}
func main() {
byRange, byHash := map[string]int{}, map[string]int{}
nextID := int64(2_000_000_100)
for i := int64(0); i < 900; i++ {
objectID := nextID + i
byRange[ranged(objectID)]++
byHash[hashed(objectID)]++
}
fmt.Println("ranged:", show(byRange))
fmt.Println("hashed:", show(byHash))
}
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
живой пример
const SHARDS = ['shardA', 'shardB', 'shardC'];
function ranged(objectId) {
if (objectId < 1_000_000_000n) return SHARDS[0];
if (objectId < 2_000_000_000n) return SHARDS[1];
return SHARDS[2];
}
function hashed(objectId) {
const mixed = BigInt.asUintN(64, objectId * 0x9E3779B97F4A7C15n);
const h = Number(BigInt.asIntN(32, mixed ^ (mixed >> 32n)));
return SHARDS[((h % SHARDS.length) + SHARDS.length) % SHARDS.length];
}
const show = (m) => '{' + [...m.keys()].sort().map((k) => `${k}=${m.get(k)}`).join(', ') + '}';
const byRange = new Map(), byHash = new Map();
const nextId = 2_000_000_100n;
for (let i = 0n; i < 900n; i++) {
const objectId = nextId + i;
byRange.set(ranged(objectId), (byRange.get(ranged(objectId)) ?? 0) + 1);
byHash.set(hashed(objectId), (byHash.get(hashed(objectId)) ?? 0) + 1);
}
console.log('ranged: ' + show(byRange));
console.log('hashed: ' + show(byHash));
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
живой пример
from collections import Counter
SHARDS = ["shardA", "shardB", "shardC"]
MASK64 = (1 << 64) - 1
def ranged(object_id: int) -> str:
if object_id < 1_000_000_000:
return SHARDS[0]
if object_id < 2_000_000_000:
return SHARDS[1]
return SHARDS[2]
def hashed(object_id: int) -> str:
mixed = (object_id * 0x9E3779B97F4A7C15) & MASK64
h = (mixed ^ (mixed >> 32)) & 0xFFFFFFFF
if h >= 1 << 31:
h -= 1 << 32
return SHARDS[h % len(SHARDS)]
def show(counts: Counter) -> str:
return "{" + ", ".join(f"{k}={counts[k]}" for k in sorted(counts)) + "}"
by_range, by_hash = Counter(), Counter()
next_id = 2_000_000_100
for i in range(900):
object_id = next_id + i
by_range[ranged(object_id)] += 1
by_hash[hashed(object_id)] += 1
print("ranged: " + show(by_range))
print("hashed: " + show(by_hash))
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Все 900 записей по диапазону ушли в один шард, по хешу разошлись примерно поровну. MongoDB делает то же самое, только хеширует по-своему.
Составной shard key
Часто лучший вариант — составной ключ: первая часть даёт разнообразие (равномерность), вторая — локальность для типичных запросов.
Если запросы обычно идут по categoryId, а в каждой категории товаров разное количество:
sh.shardCollection("shop.product", { categoryId: 1, _id: "hashed" });
Товары одной категории хешируются по _id и распределяются по шардам равномерно. При этом запросы по categoryId идут на меньшее число шардов, чем при чистом хешировании по _id.
Zoned sharding — данные в нужном регионе
Если регуляторика требует хранить данные пользователей в определённой стране — используется zoned sharding: диапазонам значений shard key назначаются зоны, шардам присваиваются те же зоны:
sh.addShardToZone("shardEU", "EU");
sh.addShardToZone("shardUS", "US");
sh.updateZoneKeyRange(
"shop.product",
{ region: "EU", productId: MinKey },
{ region: "EU", productId: MaxKey },
"EU"
);
Данные с region: "EU" физически останутся на европейских серверах. В старых руководствах те же действия записаны как sh.addShardTag() и sh.addTagRange() — это прежние имена, объявленные устаревшими ещё в 3.4, когда зоны пришли на смену тегам; в современном mongosh на них лучше не рассчитывать.
Роли узлов: скрытые, отложенные, без права стать primary
В производственном наборе узлы редко одинаковы. Четыре настройки, которые используют постоянно.
priority: 0 — узел не может стать primary. Так помечают узел в удалённом регионе: он держит копию и участвует в голосовании, но становиться ведущим ему нельзя, потому что все клиенты тогда получат задержку.
hidden: true — узел не виден драйверам, то есть на него не уйдёт ни одно чтение, даже если приложение просит secondaryPreferred. Это узел «для служебных задач»: снятие резервных копий, тяжёлая аналитика, выгрузки. Скрытый узел обязан быть и priority: 0.
slaveDelay / secondaryDelaySecs — узел применяет журнал с задержкой (например, на час). Это защита от человеческой ошибки: удалили коллекцию по скрипту — у вас есть час, чтобы достать данные из отложенного узла до того, как удаление до него доедет. Такой узел тоже делают скрытым.
Теги. Узлам приписывают метки ({ region: "eu", use: "analytics" }), а в запросе указывают, с каких узлов читать. Это единственный способ управляемо направить аналитику на выделенный узел, а пользовательские чтения — на ближайший.
Типовая рабочая конфигурация на пять узлов: два обычных в основной зоне, один в соседней с priority: 0, один скрытый для копий и один отложенный на час. Голосующих при этом должно быть нечётное число.
Как сделать чтение с реплики управляемым
Совет «secondary отстаёт, читайте с primary» верен по умолчанию, но есть две настройки, которые делают чтение с реплик пригодным к использованию.
maxStalenessSeconds — драйвер не будет читать с узла, который отстал больше указанного (минимум 90 секунд). Это защита от худшего случая: узел, отставший на час, просто исключается из выбора. Гарантии «не старше 90 секунд» это не даёт (оценка приблизительная), но отсекает мёртвые узлы.
Теги в предпочтении чтения. readPreference: secondary с тегом направляет чтение на конкретную группу узлов — например, только на аналитические. Тогда тяжёлые отчёты не мешают ни primary, ни чтениям пользователей.
И то, что надо мерить, чтобы этим пользоваться: отставание. rs.printSecondaryReplicationInfo() показывает по каждому узлу, насколько он позади, rs.status() даёт метки времени последней применённой операции, а в метриках мониторинга это отдельный показатель. Оповещение ставят на отставание больше допустимого для чтения — иначе «читаем с реплик» однажды начнёт отдавать вчерашние данные молча.
Что изменилось с chunks
Читатель, пришедший с версий 4.x, ищет привычные операции и не находит: начиная с версии 6.1 автоматическое разделение кусков (autosplit) убрано, и балансировщик сам режет диапазоны, когда переносит данные. Ручные команды разделения остались, но нужны редко.
Практическое следствие: размер куска (по умолчанию 128 МБ) теперь настройка балансировщика, а не то, за чем следят вручную; а состояние распределения смотрят по sh.status() и по метрикам балансировщика. Если данные разложены неравномерно — вопрос почти всегда в ключе, а не в кусках.
Что попробовать до шардирования
Порог «одиннадцать серверов» — это уже решение, а перед ним есть три дешёвых шага.
Вертикальный рост. Набор на машине с большим объёмом памяти обслуживает десятки терабайт: MongoDB упирается прежде всего в то, влезают ли индексы и горячие данные в память. Удвоить память дешевле, чем поднять шардированный кластер с серверами конфигурации и маршрутизаторами.
Архивирование. Старые данные уезжают в отдельную коллекцию, в объектное хранилище или удаляются по сроку жизни (индекс с истечением). Обычно 80 % объёма — это данные, которые никто не читает; убрав их, вопрос о шардировании откладывается на годы.
Вынос коллекции. Одна коллекция растёт быстрее остальных (события, журналы, вложения) — её переносят в отдельный набор или вообще в другую систему, рассчитанную на такой поток. Остальная база остаётся маленькой и управляемой.
Шардирование оправдано, когда после всех трёх шагов данные или запись не влезают на одну машину — и когда есть ключ, по которому они разложатся равномерно. Без второго условия кластер даст только сложность.
Три варианта shard key и то, чем каждый заканчивается: монотонный _id сваливает свежие товары в один шард, хеш разносит их поровну, а составной ключ с категорией первой частью и разносит, и оставляет запрос по категории точечным.
Как выбрать хороший shard key
Пять критериев, которые нужно проверить:
-
Высокая кардинальность — много разных значений. Поле
statusс тремя значениями не позволит распределить данные по десяти шардам. -
Равномерное распределение — иначе горячий шард.
userIdобычно хорошо,countryдля одной страны — плохо.Перекос видно до всякого шардинга — обычной группировкой по кандидату в ключи. Посчитайте на учебных заказах, как распределились бы они по продавцам:
живой пример
db.orders.aggregate([ { $group: { _id: "$sellerId", заказов: { $sum: 1 } } }, { $sort: { заказов: -1 } } ])Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Первая строка вдвое с лишним больше остальных — на таком ключе один шард получил бы половину нагрузки.
-
Присутствует в большинстве запросов — иначе каждый запрос опрашивает все шарды (scatter-gather), что медленно.
-
Редко меняется — shard key практически неизменный. Изменение с 5.0 возможно через
reshardCollection, но это долгая операция, и у неё две неприятные особенности: база строит вторую копию коллекции рядом со старой (значит, нужен запас свободного места), а в самом конце, пока она переключает кластер на новую раскладку, запись в коллекцию ненадолго останавливается. -
Запись идёт на один шард — каждая операция записи попадает на конкретный шард без распределённых транзакций.
Практический совет: если коллекция product шардирована, а category нет — запросы с $lookup между ними вынуждены опрашивать все шарды. Лучше шардировать обе коллекции по categoryId: тогда связанные документы окажутся физически рядом и соединение останется локальным.
Глубже: эксплуатация: explain, профайлер, резервные копии и памятьрасширенное
Раздел про Elasticsearch получил статью про эксплуатацию, MongoDB нет, хотя вопросы те же: почему запрос медленный, как не потерять данные и почему сервер ест всю память.
Медленный запрос. db.orders.find({ customerId: 42 }).explain("executionStats") показывает план: COLLSCAN означает полный обход коллекции, IXSCAN работу по индексу. Главная пара чисел это totalDocsExamined против nReturned: осмотрели сто тысяч, вернули десять, значит индекса нет или он не тот. Если SORT в плане стоит отдельной ступенью, сортировка идёт в памяти, и индекс не покрывает порядок. Какие запросы вообще медленные, показывает профайлер: db.setProfilingLevel(1, { slowms: 100 }) пишет всё дольше ста миллисекунд в коллекцию system.profile той же базы, а по умолчанию медленные запросы и без профайлера попадают в лог сервера с тем же порогом. Что происходит прямо сейчас, видно в db.currentOp(), а зависшую операцию снимают db.killOp(id).
Память. Движок WiredTiger держит кэш данных и индексов и по умолчанию берёт под него половину оперативной памяти минус гигабайт. Остальное нужно самому процессу, файловому кэшу операционной системы и соседям, поэтому на сервере с MongoDB не ставят второй прожорливый процесс, а в контейнере задают storage.wiredTiger.engineConfig.cacheSizeGB явно: без этого движок считает память хоста, а не лимит контейнера, и его убивает OOM. Признак нехватки кэша в db.serverStatus().wiredTiger.cache: страницы, которые читаются с диска, растут, bytes read into cache не отстаёт от записи.
Резервные копии. mongodump и mongorestore это логическая копия документами: годится для баз в гигабайты и для переноса между версиями, но на терабайте идёт часами и нагружает рабочий сервер. Для больших баз снимают файловую систему узла реплики (снимок диска в облаке или LVM) при включённом журнале; копию снимают с secondary, чтобы не трогать primary. Шардированный кластер это отдельная беда: у каждого шарда своя копия, и они должны быть сделаны близко по времени при остановленном балансировщике (sh.stopBalancer()), плюс копия конфигурационных серверов; иначе после восстановления часть чанков будет в двух местах. Восстановление на любой момент времени требует ещё и oplog, и здесь управляемые сервисы вроде Atlas честно выигрывают у своей установки.
Обновление версий. Только на одну мажорную версию за раз, сначала secondary по одному, потом rs.stepDown() и старый primary. После обновления бинарников новые возможности выключены, пока не поднята версия совместимости: db.adminCommand({ setFeatureCompatibilityVersion: "8.0", confirm: true }); пока она старая, можно откатиться. Размер oplog проверяют до долгих операций: rs.printReplicationInfo() показывает, на сколько часов его хватает, и если secondary отстанет дальше, он уйдёт в полную пересинхронизацию.
Коротко
- Replica set — три узла и больше с одними данными: primary принимает записи, secondary переигрывают их по oplog.
- Выборы нового primary укладываются примерно в 12 секунд, и всё это время записей нет. Непереданные большинству записи откатываются — от этого страхует
w: "majority". - Oplog фиксированного размера: secondary, отставший дальше его границы, копирует данные заново (initial sync).
- Чтение можно увести на secondary через
readPreference, но он отстаёт: своё только что записанное читают только с primary, а чтение с реплик делают управляемымmaxStalenessSeconds, тегами и оповещением на отставание (rs.printSecondaryReplicationInfo()). - Sharded cluster берут, когда активные данные не влезают в память одного сервера: это минимум 11 узлов, и до него пробуют вертикальный рост, архивирование старых данных и вынос самой большой коллекции.
- Shard key — решение почти без права на ошибку: высокая кардинальность, равномерность, присутствие в запросах. Монотонный ключ без хеша даёт горячий шард.
- Эксплуатация:
explain("executionStats")и параtotalDocsExamined/nReturned, профайлер поslowms; кэш WiredTiger половина памяти, в контейнере задать явно; копии с secondary, в шардированном кластере с остановленным балансировщиком; обновление по одной мажорной версии с FCV. - В производственном наборе узлы не одинаковы:
priority: 0для удалённых,hiddenдля копий и аналитики, отложенный узел как защита от ошибочного удаления, теги для маршрутизации. - С версии 6.1 автоматического разделения кусков нет — диапазоны режет балансировщик, а неравномерность почти всегда означает плохой ключ.
Что почитать дальше
- ACID и согласованность в MongoDB — какие гарантии работают в sharded cluster и что на самом деле обещает
w: "majority". - Моделирование документов — правильная схема снижает нагрузку и откладывает потребность в шардинге.
- Секционирование — те же горячие точки и ребалансировка, но без привязки к MongoDB.
- Модели репликации — почему реплика отстаёт и как считают кворумы.