Sharding вживую

Hashed и ranged shard key, chunk splitting, balancer, роутинг через mongos, config servers, zone sharding, jumbo chunks и resharding — на живом шардированном кластере, с хотспотами и тем, как их ловить

Шардинг MongoDB — тот раздел, который нельзя закрыть чтением документации: последствия выбора shard key видны только на живом кластере, под реальной записью, когда балансировщик начинает двигать чанки, а не в теории. Здесь — полный шардированный кластер: mongos, config replica set и несколько шардов, каждый из которых сам по себе replica set. Разобраны hashed и ranged shard key и последствия выбора, chunk splitting и balancer, роутинг запросов, zone sharding, jumbo chunks и resharding — с хотспотами, пойманными вживую, а не описанными абстрактно.

Шардированный кластер: маскот-лист mongos-диспетчер в наушниках + config servers с картой чанков; шарды shard1/shard2 (каждый — PRIMARY и два SECONDARY листа); «1) HASHED ROUTING — равномерно 0.500/0.500», «2) RANGED ROUTING — риск хотспота»; JUMBO CHUNK, роборука BALANCER двигает чанк (диапазон 10–20), «RESHARDING IN PROGRESS 23%»; итог «HASHED = balanced, RANGED = locality (risk)»

В статье

Предыдущая статья серии — 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.
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

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
Шардированная топология стенда: mongos → config RS → shard1/shard2, и контраст распределения по shard key — hashed равномерно (0.500 / 0.500) против ranged write-хотспот. Числа — с живого стенда mongodb/sharding, mongo:8.2.11, seed=42

Последовательность инициации на стенде: поднимаются 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.

Обсуждение в Telegram

Присоединиться →

Комментарии