← назад к разделу

Когда данных становится много, один сервер перестаёт справляться — медленно читает, медленно пишет, и если упадёт, приложение встанет. MongoDB решает это двумя механизмами: репликация хранит копии данных на нескольких серверах (надёжность), шардинг разрезает коллекцию между серверами (горизонтальное масштабирование).

записи принимает primary, secondary догоняют его по oplog primary все записи oplog: операции дописываются в конец order-1NEW order-1PAID order-2NEW order-1SHIPPED позиция чтения узла A B primary недоступен secondary A переигрывает oplog secondary B отстал на одну запись primaryпринимает записи выборы: нужен голос большинства из трёх узлову A четыре записи, у B три — primary становится A

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Когда нужны самые актуальные данные
primaryPreferredPrimary, при его недоступности — secondaryЛучше отдать чуть устаревшее, чем ошибку на время выборов
secondaryТолько secondaryАналитика, отчёты — не нагружать primary
secondaryPreferredSecondary, при недоступности — 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 — роутер Shard A набор реплик Shard B набор реплик Shard C набор реплик конфигурационные серверы набор из трёх узлов

Клиент не знает про шарды: он говорит с роутером, а тот смотрит в конфигурационные серверы и решает, куда идти. Поэтому потеря конфигурационных серверов страшнее потери одного шарда.

  • 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 % объёма — это данные, которые никто не читает; убрав их, вопрос о шардировании откладывается на годы.

Вынос коллекции. Одна коллекция растёт быстрее остальных (события, журналы, вложения) — её переносят в отдельный набор или вообще в другую систему, рассчитанную на такой поток. Остальная база остаётся маленькой и управляемой.

Шардирование оправдано, когда после всех трёх шагов данные или запись не влезают на одну машину — и когда есть ключ, по которому они разложатся равномерно. Без второго условия кластер даст только сложность.

по _id ObjectId растёт один диапазон шард A горит по _id hashed хеш перемешал три диапазона нагрузка ровно categoryId + _id категория первой внутри хеш точечный запрос

Три варианта shard key и то, чем каждый заканчивается: монотонный _id сваливает свежие товары в один шард, хеш разносит их поровну, а составной ключ с категорией первой частью и разносит, и оставляет запрос по категории точечным.

Как выбрать хороший shard key

Пять критериев, которые нужно проверить:

  1. Высокая кардинальность — много разных значений. Поле status с тремя значениями не позволит распределить данные по десяти шардам.

  2. Равномерное распределение — иначе горячий шард. userId обычно хорошо, country для одной страны — плохо.

    Перекос видно до всякого шардинга — обычной группировкой по кандидату в ключи. Посчитайте на учебных заказах, как распределились бы они по продавцам:

    живой пример

    db.orders.aggregate([
      { $group: { _id: "$sellerId", заказов: { $sum: 1 } } },
      { $sort: { заказов: -1 } }
    ])
    
    Запустить

    Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

    Первая строка вдвое с лишним больше остальных — на таком ключе один шард получил бы половину нагрузки.

  3. Присутствует в большинстве запросов — иначе каждый запрос опрашивает все шарды (scatter-gather), что медленно.

  4. Редко меняется — shard key практически неизменный. Изменение с 5.0 возможно через reshardCollection, но это долгая операция, и у неё две неприятные особенности: база строит вторую копию коллекции рядом со старой (значит, нужен запас свободного места), а в самом конце, пока она переключает кластер на новую раскладку, запись в коллекцию ненадолго останавливается.

  5. Запись идёт на один шард — каждая операция записи попадает на конкретный шард без распределённых транзакций.

Практический совет: если коллекция 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 автоматического разделения кусков нет — диапазоны режет балансировщик, а неравномерность почти всегда означает плохой ключ.

Что почитать дальше