Шардинг MongoDB — тот раздел, который нельзя закрыть чтением документации: последствия выбора shard key видны только на живом кластере, под реальной записью, когда балансировщик начинает двигать чанки, а не в теории. Здесь — полный шардированный кластер: mongos, config replica set и несколько шардов, каждый из которых сам по себе replica set. Разобраны hashed и ranged shard key и последствия выбора, chunk splitting и balancer, роутинг запросов, zone sharding, jumbo chunks и resharding — с хотспотами, пойманными вживую, а не описанными абстрактно.
В статье
- Топология стенда:
mongos, config servers, шарды - Shard key: hashed против ranged — выбор и его последствия
- Chunk splitting и работа балансировщика
- Роутинг запросов через
mongos - Zone sharding — привязка данных к географии или классу узлов
- Jumbo chunks — как возникают и что с ними делать
- Resharding вживую
- Хотспоты — как их поймать на реальном кластере
- Граница и что дальше
Предыдущая статья серии — Replica sets, oplog и concerns — держала в кадре один replica set: кворумные выборы, oplog, цену w:majority. Здесь этих replica set становится несколько, и каждый из них — уже не самостоятельный кластер, а шард. Всё измеренное ниже снято на живом шардированном кластере mongo:8.2.11: роутер mongos, config-server replica set и два шарда, на сквозном датасете серии (seed=42: 200 000 заказов), стенд mongodb/sharding. Бинарник sharding/main.go подключается к роутеру mongos, а не к отдельному шарду, и меряет то, что видно только на живом кластере: как shard key раскладывает чанки, в сколько шардов уходит запрос, что делает балансировщик и во что обходится resharding.
Топология стенда: mongos, config servers, шарды
Шардированный кластер MongoDB — это три роли, а не «просто несколько узлов». На стенде каждая роль поднята вживую:
mongos— роутер запросов. Не хранит данных сам; знает, какой чанк на каком шарде, и направляет туда операции. Клиент подключается кmongos, а не к шардам напрямую.- config server replica set (
csrs) — метаданные кластера: карта чанков (config.chunks), список шардов, настройки балансировщика (config.settings), список шардированных коллекций (config.collections). На стенде — 1 узел с флагом--configsvr(в проде это отдельный 3-узловой replica set). - шарды
shard1/shard2— собственно данные. Каждый шард — сам по себе replica set (те самые из статьи #5), на стенде по 1 узлу с флагом--shardsvr.
роутер запросов"] MG --> CS["config RS (csrs)
карта чанков, метаданные"] MG --> S1["shard1 (RS)"] MG --> S2["shard2 (RS)"] MG --> SK["shardCollection:
выбор shard key"] SK -->|"_id hashed"| EVEN["shard1 50% · shard2 50%
оба шарда, равномерно"] SK -->|"_id ranged (монотонный)"| HOT["все новые вставки → верхний maxKey-чанк
один шард: write-хотспот"] EVEN -.-> S1 EVEN -.-> S2 HOT -.-> S1 style EVEN fill:#c9e4c5,stroke:#5b8a5e style HOT fill:#f4d9c6,stroke:#c67a4a
flowchart TD
APP["клиент"] --> MG["mongos
роутер запросов"]
MG --> CS["config RS (csrs)
карта чанков, метаданные"]
MG --> S1["shard1 (RS)"]
MG --> S2["shard2 (RS)"]
MG --> SK["shardCollection:
выбор shard key"]
SK -->|"_id hashed"| EVEN["shard1 50% · shard2 50%
оба шарда, равномерно"]
SK -->|"_id ranged (монотонный)"| HOT["все новые вставки → верхний maxKey-чанк
один шард: write-хотспот"]
EVEN -.-> S1
EVEN -.-> S2
HOT -.-> S1
style EVEN fill:#c9e4c5,stroke:#5b8a5e
style HOT fill:#f4d9c6,stroke:#c67a4a
Последовательность инициации на стенде: поднимаются 4 контейнера (csrs1, shard1a, shard2a, mongos1); каждый replica set инициируется в своей роли (rs.initiate() с configsvr:true для config-сервера); роутер mongos1 стартует с --configdb csrs/csrs1:27017 и регистрирует оба шарда через sh.addShard(...). Дальше включается шардирование базы (sh.enableSharding("cookbook")), размер чанка снижается до 1 МБ (config.settings, _id:"chunksize") — чтобы на 200 000 документов получилось осмысленное число чанков, а не один на всё, — и шардируются две коллекции. Балансировщик перед импортом останавливается (sh.stopBalancer()), чтобы pre-balance срез отражал чистую раскладку роутера на вставке.
// на mongos, после готовности config-сервера и обоих шардов
sh.addShard("shard1/shard1a:27017");
sh.addShard("shard2/shard2a:27017");
sh.enableSharding("cookbook");
// уменьшаем размер чанка до 1 МБ — иначе весь датасет уместится в один чанк
db.getSiblingDB("config").settings.updateOne(
{ _id: "chunksize" },
{ $set: { value: 1 } },
{ upsert: true }
);
sh.stopBalancer(); // pre-balance срез = чистая раскладка роутера
Честная оговорка про топологию (важно для интерпретации всех чисел ниже): каждый шард на стенде — это 1-узловой replica set, config-server тоже 1 узел. В продакшене каждый шард был бы 3-узловым replica set (иначе теряется отказоустойчивость самого шарда — тема статьи #5). Для демонстрации именно шардирования — роутинга, чанков, балансировки, resharding — минимальной топологии достаточно, но она напрямую влияет на стоимость тяжёлых операций (см. resharding ниже, где эта упрощённость проявляется остро).
Shard key: hashed против ranged — выбор и его последствия
Это главный контраст стенда и главное решение при проектировании шардированной коллекции. Обе коллекции — один и тот же датасет (200 000 заказов, import_orders_hashed=200000, import_orders_ranged=200000), отличаются они только shard key. Срез снят при выключенном балансировщике (balancer_enabled_at_start=false) — то есть распределение ровно такое, каким его сделал роутер по shard key на вставке, без последующей балансировки:
| Коллекция | shard key | чанков всего | shard1 | shard2 | max-доля (по чанкам) |
|---|---|---|---|---|---|
orders_hashed |
{_id: "hashed"} |
2 | 1 | 1 | 0.500 (равномерно) |
orders_ranged |
{_id: 1} (ranged) |
1 | 1 | 0 | 1.000 (весь на одном шарде) |
Ключевая деталь датасета: _id заказов — монотонный синтетический ObjectId-совместимый идентификатор, сгенерированный hex-счётчиком (первый …200000000001). Это худший случай для ranged-ключа, и стенд показывает почему:
- hashed-ключ (
{_id: "hashed"}) наshardCollectionпустой коллекции сразу пресплитит начальные чанки по шардам — по одному на шард, — и роутер раскладывает документы по хэшу_idравномерно. Результат: оба шарда покрыты,max_share=0.500— ровно половина на каждый. - ranged-ключ (
{_id: 1}) стартует с одного чанка[minKey, maxKey)на primary-шарде. Из-за монотонного_idвсе вставки попадают в верхний чанк (тот, что граничит сmaxKey):max_share=1.000— 100% записи в один шард. Это классический write-хотспот.
Ассерт стенда (прошёл): hashed охватывает оба шарда, ranged концентрированнее hashed (ranged.maxShare=1.000 > hashed.maxShare=0.500). Практический вывод прямой: под монотонный первичный ключ ranged shard key превращает горизонтально масштабируемый кластер обратно в один узел на запись. Hashed убивает эту монотонность ценой того, что диапазонные запросы по ключу перестают быть локальными (об этом — в роутинге).
Chunk splitting и работа балансировщика
Данные шарда физически живут в чанках — непрерывных диапазонах значений shard key. Когда чанк перерастает лимит (chunksize, на стенде 1 МБ), он разрезается (split) на два. Разрезание — операция над метаданными, дешёвая; она не двигает данные. Данные двигает балансировщик: фоновый процесс, который выравнивает распределение между шардами, мигрируя чанки с перегруженного шарда на недогруженный.
Стенд включает балансировщик (balancerStart) и опрашивает его статус вместе с config.chunks, пока раунд не сойдётся (три замера подряд при inBalancerRound=false). Сошлось за ~1m20s (balance_settle_time_approx=1m20s):
| Коллекция | pre (shard1/shard2) | post (shard1/shard2) | покрытие шардов | max-доля (по чанкам) |
|---|---|---|---|---|
orders_hashed |
1 / 1 | 1 / 1 | 2→2 (без изменений) | 0.500→0.500 |
orders_ranged |
1 / 0 | 1 / 35 | 1→2 | 1.000→0.972 |
hashed уже был равномерен — балансировщику делать нечего, распределение не изменилось. ranged: единственный чанк на shard1 разрезан на 36 чанков, и бо́льшая часть мигрирована на shard2 — покрытие выросло с одного шарда до двух (ranged_movement pre{covered=1} post{covered=2}).
Здесь — главная находка стенда, и её нельзя прочитать наивно. Post-срез ranged по числу чанков выглядит перекошенным: 1 на shard1 против 35 на shard2. Соблазн сказать «данные не сбалансированы» — и это была бы ошибка. Начиная с MongoDB 6.0 балансировщик оперирует объёмом данных, а не числом чанков. Само число чанков — и max_share, которую стенд по нему считает, — уже не метрика баланса объёма. Что стенд действительно доказал: исходный единственный чанк на shard1 разрезан, часть чанков мигрирована, и покрытие выросло с одного шарда до двух (ranged_movement pre{covered=1} post{covered=2}). А чего он не измерял — фактического распределения байтов и документов между шардами: max_share=0.972 — это доля чанков на busiest-шарде (35 из 36), а не доля данных, и балансировку именно объёма этими числами подтвердить нельзя. «1 против 35 чанков» — не про дисбаланс, а артефакт того, как разрезался исходный монотонный диапазон; о реальной балансировке объёма судят по размеру данных на шардах, чего минимальный стенд не снимал.
И ещё: балансировка не отменяет урок ranged-ключа. Балансировщик разрезал и мигрировал уже накопившиеся чанки постфактум — но write-хотспот был и остаётся про то, куда попадают новые вставки. У монотонного _id новая запись всегда идёт в верхний maxKey-чанк, а это всегда один шард. Балансировщик потом растащит накопившееся, но пиковая запись как била в один шард, так и бьёт. Выбор shard key лечит хотспот на входе; балансировщик — только уборщик после.
Роутинг запросов через mongos
Роутер знает карту чанков и по ней решает, в сколько шардов направить запрос. Это видно на explain с самого mongos (verbosity:"queryPlanner") — по верхней стадии winningPlan и числу целевых шардов:
Запрос к orders_hashed |
winningPlan.stage | целевых шардов |
|---|---|---|
по shard key (_id == <конкретный>) |
SINGLE_SHARD |
1 |
без shard key (status == "paid") |
SHARD_MERGE |
2 |
// по shard key — роутер знает шард по хэшу ключа
db.orders_hashed.find({ _id: someId }).explain("queryPlanner");
// winningPlan.stage = "SINGLE_SHARD", целевых шардов: 1
// без shard key — веером по всем шардам + merge
db.orders_hashed.find({ status: "paid" }).explain("queryPlanner");
// winningPlan.stage = "SHARD_MERGE", целевых шардов: 2
Запрос по shard key роутится в ровно один шард (SINGLE_SHARD, shards_targeted=1): роутер по хэшу ключа точно знает, где лежит документ. Запрос без shard key — это scatter-gather (SHARD_MERGE, shards_targeted=2): mongos рассылает запрос на все шарды, каждый выполняет его локально, роутер сливает ответы.
Это ключевой практический довод при проектировании: shard key нужно выбирать под самые частые запросы приложения. На двух шардах разница между targeting’ом в 1 шард и в 2 кажется мелкой — но кластер на два шарда никто не строит. На 20 шардах запрос без shard key превращается в 20 параллельных запросов с последующим merge, и латентность определяется самым медленным из двадцати. Здесь же виден и подвох hashed-ключа из предыдущего раздела: он даёт идеальное распределение записи, но диапазонный запрос по _id (_id > X) по хэшу локализовать нельзя — он неизбежно scatter-gather. Ranged-ключ дал бы такому запросу SINGLE_SHARD, но ценой write-хотспота. Выбор shard key — это всегда компромисс между равномерностью записи и локальностью чтения, и решается он не в теории, а по профилю запросов конкретного приложения.
Zone sharding — привязка данных к географии или классу узлов
Балансировщик по умолчанию стремится к равномерному распределению объёма. Zone sharding (ранее tag-aware sharding) позволяет это переопределить: привязать диапазон значений shard key к зоне, а зону — к конкретным шардам. Балансировщик тогда держит данные из диапазона только на шардах нужной зоны.
// пометить шарды зонами
sh.addShardToZone("shard1", "EU");
sh.addShardToZone("shard2", "US");
// диапазон shard key → зона
sh.updateZoneKeyRange(
"cookbook.orders_by_region",
{ region: "EU", _id: MinKey },
{ region: "EU", _id: MaxKey },
"EU"
);Типичные сценарии: data residency (данные европейских пользователей физически на шардах в EU-датацентре — требование регуляторов), tiered storage (горячие свежие данные на шардах с NVMe, холодные архивные — на шардах с ёмкими HDD). На минимальной топологии стенда (два 1-узловых шарда) осмысленной зональной картины не построить, поэтому zone sharding здесь — на уровне механики, не живых чисел; общий приём «данные ближе к их потребителю» разобран в статье про шардинг вне MongoDB (см. границу ниже).
Jumbo chunks — как возникают и что с ними делать
Балансировщик двигает чанки, но у него есть предел: он не может мигрировать чанк, который невозможно разрезать. Если множество документов делит одно и то же значение shard key и суммарно перерастает лимит чанка, разрезать этот диапазон не по чему — граница split’а прошла бы посередине одного значения ключа, а так нельзя. Такой чанк помечается jumbo и застревает на своём шарде: балансировщик его пропускает.
Классический источник jumbo-чанков — shard key с низкой кардинальностью: например, {country: 1}, где на одну популярную страну приходятся миллионы документов. Все они делят значение country: "RU", диапазон ["RU", "RU"] неделим, и чанк растёт неограниченно, оставаясь на одном шарде — снова хотспот, теперь по хранению. Отсюда правило выбора shard key: высокая кардинальность (много различных значений) плюс низкая частотность (ни одно значение не доминирует). Монотонный _id из нашего стенда кардинальность имеет отличную (каждое значение уникально) — его беда была не в jumbo, а в монотонности (write-хотспот). Это два разных провала одного выбора: низкая кардинальность даёт jumbo-чанки, монотонность даёт write-хотспот.
Практический разбор застрявшего jumbo-чанка — либо составной shard key, добавляющий кардинальности ({country: 1, _id: 1}), либо, если коллекция уже живёт в проде с плохим ключом, — resharding, к которому и переходим.
Resharding вживую
Раньше сменить shard key на живой коллекции было нельзя — только выгрузить, пересоздать, залить заново. С MongoDB 5.0 появился reshardCollection, а к 8.x он дозрел до операции, которую можно запустить на работающей коллекции без остановки. Стенд меняет shard key orders_ranged с монотонного {_id: 1} на {user_id: "hashed"} — то есть буквально лечит write-хотспот из раздела про shard key:
| Параметр | Значение (этот прогон) |
|---|---|
| key_before | {_id: 1} |
| new_key | {user_id: "hashed"} |
| результат | OK (config.collections.key подтверждает смену) |
| длительность | ~5m29s |
db.adminCommand({
reshardCollection: "cookbook.orders_ranged",
key: { user_id: "hashed" }
});
// результат: OK, elapsed ~5m29s
// после успеха config.collections.key = { user_id: "hashed" }
Внутри resharding — тяжёлая операция. Координатор клонирует всю коллекцию под новый ключ (фаза cloning доминирует по времени), затем applying (докатывает изменения, накопившиеся за время клонирования) → blocking-writes (короткая заморозка записи) → committing (атомарная подмена). На 200 000 документов на этой топологии успешный прогон занял ~5m29s.
Честная оговорка про устойчивость — и это не единичный сбой, а свойство минимальной топологии. Стоимость resharding здесь сидит прямо на границе используемого окна ожидания, и исход нестабилен от прогона к прогону. Три живые попытки одного и того же reshardCollection на одном и том же датасете и топологии дали три разных исхода:
- успех за 5m29s (число в таблице выше);
MaxTimeMSExpiredна ~4m37s — при клиентском потолке 5 минут операция не успела и вернула таймаут (из-за чего стенд поднял потолок до 8 минут);- обрыв соединения с
mongosна 4m10s (incomplete read of message header: connection closed unexpectedly ... EOF) — стенд честно зафиксировалreshard_result=NOT_REPRODUCED, при этом все остальные сценарии того же прогона (распределение чанков, targeting, балансировщик) прошли идентично числам выше.
Важный нюанс: клиентский таймаут не отменяет операцию на сервере — resharding resumable и продолжается независимо от того, дождался ли его клиент. Практический вывод для продакшена прямой: resharding крупных коллекций нужно либо планировать с большим запасом по времени, либо отслеживать прогресс асинхронно ($currentOp, поле reshardingFields.state в config.collections) вместо блокирующего ожидания ответа команды. Число «5m29s» — это реальный успешный прогон, а не гарантия; на минимальной топологии (1-узловые shard RS, один хост) стоимость операции настолько близка к границе окна ожидания, что рассчитывать на конкретное время нельзя. На проде с 3-узловыми шардами и параллельным вводом-выводом картина будет иной — но урок «resharding дорогой и требует асинхронного контроля» от этого только крепнет.
Хотспоты — как их поймать на реальном кластере
Хотспот — это когда нагрузка, которая должна была растечься по шардам, концентрируется на одном. На стенде он пойман вживую в двух формах:
- Write-хотспот —
orders_rangedс монотонным_id:max_share=1.000на вставке, 100% записи в один шард. Ловится поconfig.chunks, сгруппированному по shard: если один шард держит долю, близкую к 1.0, а покрытие = 1 шард — это он. Признак на уровне приложения: shard key монотонно растёт со временем (_id, timestamp, автоинкремент). - Query-хотспот (scatter-gather) — запрос без shard key:
SHARD_MERGEвместоSINGLE_SHARD, targeting во все шарды. Ловится черезexplainнаmongos: если частый запрос даётSHARD_MERGE, а неSINGLE_SHARD, значит его фильтр не содержит shard key, и он веером идёт по всем шардам.
Инструменты диагностики, использованные на стенде: config.chunks (реальная раскладка чанков по шардам и max-доля), explain с mongos (стадия SINGLE_SHARD / SHARD_MERGE и число целевых шардов), balancerStatus (идёт ли раунд балансировки). Лечение зависит от формы: write-хотспот от монотонного ключа лечится hashed-ключом (ценой локальности диапазонных чтений) или составным ключом с высокой энтропией в первом поле; query-хотспот лечится выбором shard key под профиль частых запросов; jumbo-хотспот по хранению — повышением кардинальности ключа. Во всех трёх случаях, если коллекция уже в проде с плохим ключом, крайнее средство — resharding (с оговорками про его стоимость выше).
Граница и что дальше
Всё выше — про то, как шардирование устроено именно в MongoDB: mongos, config server replica set, чанки, reshardCollection, hashed/ranged shard key, targeting через explain. Общая теория шардинга — вне MongoDB — это отдельный разговор: consistent hashing, range- vs hash- vs directory-based партиционирование, rebalancing-стратегии, вертикальное vs горизонтальное партиционирование, шардинг в реляционных БД и вручную на уровне приложения. Он ведётся в отдельной статье Sharding в проде — она про принципы, применимые к любому хранилищу, тогда как эта — про конкретную реализацию MongoDB. Кто пришёл за общими паттернами распределения данных — туда; кто за механикой MongoDB — уже здесь.
Что осталось в руках после этой статьи: выбор shard key — самое дорогое и самое трудно-отменяемое решение шардированной коллекции. Hashed даёт равномерную запись, но теряет локальность диапазонных чтений; ranged на монотонном ключе даёт локальность, но убивает масштабирование записи в один хотспот. Балансировщик выравнивает объём (не число чанков — с 6.0 это уже не метрика), но лечит только уже лежащие данные, а не хотспот на входе. Targeting запроса определяется тем, есть ли в фильтре shard key: один шард против веера по всем. А resharding, хоть и работает вживую, — операция настолько тяжёлая, что на минимальной топологии её исход балансирует на грани таймаута; в проде её планируют с запасом и контролируют асинхронно.
Дальше в серии — Эксплуатация и карта выбора: пулы соединений, change streams как CDC, retryable writes, бэкапы и итоговая карта, когда MongoDB — верный выбор, а когда нет. Живой стенд к этой статье — mongodb/sharding.
Комментарии