MongoDB Sharding Thực Chiến: Hướng Dẫn Tối Ưu Quy Mô Lớn (Kèm Docker)
Khi nào Vertical Scaling (Nâng cấp RAM/SSD) là chưa đủ?
Trong hành trình vận hành hệ thống dữ liệu, đến một thời điểm nào đó, bạn sẽ nhận ra rằng việc nâng cấp phần cứng trên một server duy nhất không còn là giải pháp bền vững. MongoDB gọi đây là giới hạn của Vertical Scaling (mở rộng theo chiều dọc).
Replication vs Sharding: Sự khác biệt cốt lõi
Nhiều kỹ sư thường nhầm lẫn giữa Replication và Sharding. Cần phải phân biệt rõ:
-
Replica Set (Replication) – Giải pháp cho High Availability (sẵn sàng cao). Dữ liệu được nhân bản giống hệt nhau trên nhiều node. Nếu node Primary chết, một node Secondary sẽ được bầu làm Primary mới. Replication không giải quyết được bài toán dung lượng vượt ngưỡng, bởi vì mỗi node trong Replica Set đều lưu trữ toàn bộ dữ liệu.
-
Sharding – Giải pháp cho Horizontal Scaling (mở rộng theo chiều ngang). Dữ liệu được phân chia thành các phần nhỏ (chunk) và phân phối trên nhiều server khác nhau. Mỗi server (Shard Node) chỉ chứa một subset (tập con) của toàn bộ dữ liệu. Sharding giúp vượt qua giới hạn lưu trữ vật lý của một server đơn lẻ.
💡 Nguyên tắc vàng: Replication để chống chịu lỗi, Sharding để mở rộng dung lượng. Trong thực tế, mỗi Shard Node thường được triển khai dưới dạng một Replica Set để đảm bảo cả hai yếu tố.
Dấu hiệu cảnh báo hệ thống của bạn bắt buộc phải lên Sharding ngay lập tức
Theo tài liệu chính thức của MongoDB, bạn nên xem xét sharding khi gặp các dấu hiệu sau:
- Working Set vượt quá dung lượng RAM – Khi tập dữ liệu thường xuyên được truy cập (working set) lớn hơn bộ nhớ RAM, MongoDB buộc phải đọc từ đĩa, làm tăng độ trễ truy vấn một cách đáng kể.
- Dung lượng Collection ≥ 3TB – Đây là ngưỡng khuyến nghị của MongoDB để bắt đầu sharding.
- CPU bị quá tải bởi query rates cao – Khi tần suất truy vấn vượt quá khả năng xử lý của CPU trên một server.
- Không thể nâng cấp thêm phần cứng – Các nhà cung cấp cloud đều có giới hạn cứng về cấu hình máy chủ.
So sánh chi phí định tính:
| Phương án | Chi phí vận hành | Giới hạn | Độ phức tạp |
|---|---|---|---|
| Vertical Scaling | Tăng theo cấp số nhân khi lên cấu hình cao nhất | Giới hạn vật lý của server (tối đa vài TB RAM, vài chục core) | Thấp |
| Sharding | Tăng tuyến tính theo số server | Về lý thuyết không giới hạn | Cao hơn (cần thiết kế Shard Key, quản lý cluster) |
Giải mã bộ ba quyền lực trong cụm Sharded Cluster

Một cụm Sharded Cluster trong MongoDB được cấu thành từ ba thành phần cốt lõi:
Shard Nodes
Mỗi Shard là một Replica Set chứa một subset dữ liệu của toàn bộ cluster. Trong môi trường Production, mỗi Shard nên được triển khai với ít nhất 3 thành viên (một Primary, hai Secondary) để đảm bảo High Availability.
Shard Nodes là nơi dữ liệu thực sự được lưu trữ và xử lý truy vấn. Khi bạn thực hiện một câu lệnh find() hoặc insert(), công việc cuối cùng đều diễn ra tại các Shard này.
Config Servers
Config Servers lưu trữ toàn bộ metadata và cấu hình của cluster. Đây là nơi ghi nhận:
- Thông tin về các Shard trong cluster
- Ánh xạ giữa các chunk dữ liệu và Shard chứa chúng
- Cấu trúc Shard Key của từng collection
Config Servers bắt buộc phải được triển khai dưới dạng Replica Set (CSRS) với ít nhất 3 thành viên. Không được phép chạy Config Server dưới dạng standalone trong Production.
Mongos Router
mongos đóng vai trò là query router – cổng giao tiếp duy nhất giữa ứng dụng client và cụm sharded cluster.
Luồng xử lý một câu lệnh query:
Client → mongos → Config Servers (lấy metadata) → Shard Node (thực thi) → mongos → Client
Khi nhận được request, mongos sẽ:
- Tra cứu metadata từ Config Servers (hoặc từ cache của chính nó)
- Xác định chính xác Shard nào chứa dữ liệu cần truy vấn (nếu query có chứa Shard Key)
- Định tuyến request đến đúng Shard
- Tổng hợp kết quả và trả về cho client
⚠ Lưu ý quan trọng: Trong MongoDB 8.0, bạn không thể kết nối trực tiếp đến một Shard để thực hiện các câu lệnh trên collection đã được sharding. Hệ thống sẽ trả về lỗi: “You are connecting to a sharded cluster improperly by connecting directly to a shard. Please connect to the cluster via a router (mongos)”.
Kỹ thuật chọn Shard Key: Quyết định sống còn tránh thảm họa ‘Hotspot Shard’
Shard Key là yếu tố quan trọng nhất quyết định sự thành bại của một cụm Sharded Cluster. Một Shard Key được chọn sai có thể dẫn đến Hotspot – tình trạng toàn bộ dữ liệu và truy vấn dồn vào một Shard duy nhất, khiến các Shard khác bị bỏ phí.
Ranged Sharding vs Hashed Sharding: Ưu và nhược điểm

| Tiêu chí | Ranged Sharding | Hashed Sharding |
|---|---|---|
| Cơ chế | Phân phối dữ liệu theo khoảng giá trị của Shard Key | Băm (hash) giá trị Shard Key để phân phối ngẫu nhiên |
| Ưu điểm | Tối ưu cho truy vấn range (ví dụ: find({age: {$gt: 18, $lt: 30}})) |
Phân phối dữ liệu đều, tránh hotspot |
| Nhược điểm | Nguy cơ hotspot nếu Shard Key là monotonic (tăng đều) | Truy vấn range kém hiệu quả vì dữ liệu nằm rải rác |
| Phù hợp | Truy vấn range chiếm ưu thế | Ghi dữ liệu với tần suất cao (log, IoT) |
Sai lầm chết người: ID tự tăng hoặc Timestamp
Đây là một trong những sai lầm phổ biến nhất của người mới bắt đầu với Sharding.
Giả sử bạn chọn created_at (timestamp) làm Shard Key với Ranged Sharding. Các document mới luôn có created_at lớn hơn các document cũ. MongoDB sẽ liên tục ghi dữ liệu mới vào chunk có giá trị cao nhất (gần MaxKey), và chunk này nằm trên cùng một Shard.
Kết quả: một Shard phải gánh toàn bộ workload ghi, trong khi các Shard khác hầu như không hoạt động. Đây chính là Hotspot Shard – thảm họa hiệu năng.
Giải pháp: Sử dụng Hashed Sharding cho các trường monotonic như _id (ObjectId) hoặc timestamp. Hashed Sharding sẽ băm giá trị và phân phối đều trên toàn bộ cluster.
// ✅ ĐÚNG - Hashed Sharding cho trường _id (monotonic)
sh.shardCollection("myDB.orders", { "_id": "hashed" })
// ✅ ĐÚNG - Hashed Sharding cho user_id để phân phối đều
sh.shardCollection("myDB.user_activities", { "user_id": "hashed" })
// ❌ SAI - Ranged Sharding trên timestamp gây hotspot
sh.shardCollection("myDB.logs", { "timestamp": 1 }) // NGUY HIỂM!
Thiết lập Sharded Cluster với Docker Compose

Phần này sẽ hướng dẫn bạn thiết lập một cụm Sharded Cluster hoàn chỉnh bằng Docker Compose để thực hành.
⚠ Yêu cầu: Docker ≥ 26.0.0 và Docker Compose ≥ v2.27.0. Sử dụng MongoDB 7.0 (LTS) để đảm bảo tương thích.
Cấu hình Docker Compose
Dưới đây là file docker-compose.yml cho cụm Sharded Cluster với:
- 3 Config Servers (Replica Set)
- 2 Shards, mỗi Shard là một Replica Set với 3 thành viên
- 1 Mongos Router
version: '3.8'
services:
# ============ CONFIG SERVERS (Replica Set: cfg-rs) ============
config-1:
image: mongo:7.0
command: mongod --configsvr --replSet cfg-rs --port 27017 --bind_ip_all
container_name: config-1
volumes:
- config-1-data:/data/db
networks:
- shard-network
config-2:
image: mongo:7.0
command: mongod --configsvr --replSet cfg-rs --port 27017 --bind_ip_all
container_name: config-2
volumes:
- config-2-data:/data/db
networks:
- shard-network
config-3:
image: mongo:7.0
command: mongod --configsvr --replSet cfg-rs --port 27017 --bind_ip_all
container_name: config-3
volumes:
- config-3-data:/data/db
networks:
- shard-network
# ============ SHARD 1 (Replica Set: shard1-rs) ============
shard1-1:
image: mongo:7.0
command: mongod --shardsvr --replSet shard1-rs --port 27017 --bind_ip_all
container_name: shard1-1
volumes:
- shard1-1-data:/data/db
networks:
- shard-network
shard1-2:
image: mongo:7.0
command: mongod --shardsvr --replSet shard1-rs --port 27017 --bind_ip_all
container_name: shard1-2
volumes:
- shard1-2-data:/data/db
networks:
- shard-network
shard1-3:
image: mongo:7.0
command: mongod --shardsvr --replSet shard1-rs --port 27017 --bind_ip_all
container_name: shard1-3
volumes:
- shard1-3-data:/data/db
networks:
- shard-network
# ============ SHARD 2 (Replica Set: shard2-rs) ============
shard2-1:
image: mongo:7.0
command: mongod --shardsvr --replSet shard2-rs --port 27017 --bind_ip_all
container_name: shard2-1
volumes:
- shard2-1-data:/data/db
networks:
- shard-network
shard2-2:
image: mongo:7.0
command: mongod --shardsvr --replSet shard2-rs --port 27017 --bind_ip_all
container_name: shard2-2
volumes:
- shard2-2-data:/data/db
networks:
- shard-network
shard2-3:
image: mongo:7.0
command: mongod --shardsvr --replSet shard2-rs --port 27017 --bind_ip_all
container_name: shard2-3
volumes:
- shard2-3-data:/data/db
networks:
- shard-network
# ============ MONGOS ROUTER ============
mongos:
image: mongo:7.0
command: mongos --configdb cfg-rs/config-1:27017,config-2:27017,config-3:27017 --bind_ip_all
container_name: mongos
ports:
- "27018:27017"
depends_on:
- config-1
- config-2
- config-3
- shard1-1
- shard1-2
- shard1-3
- shard2-1
- shard2-2
- shard2-3
networks:
- shard-network
volumes:
config-1-data:
config-2-data:
config-3-data:
shard1-1-data:
shard1-2-data:
shard1-3-data:
shard2-1-data:
shard2-2-data:
shard2-3-data:
networks:
shard-network:
driver: bridge
Giải thích:
--configsvr và --shardsvr là hai flag bắt buộc để phân biệt Config Server và Shard Node
--replSet khai báo tên của Replica Set mà node này thuộc về
- Mongos sử dụng tham số
--configdb để chỉ định địa chỉ của các Config Servers
Khởi tạo Config Server Replica Set
Sau khi chạy docker compose up -d, bạn cần khởi tạo Replica Set cho Config Servers:
# Kết nối vào config-1
docker exec -it config-1 mongosh --port 27017
Trong Mongo Shell, chạy:
rs.initiate({
_id: "cfg-rs",
configsvr: true,
members: [
{ _id: 0, host: "config-1:27017" },
{ _id: 1, host: "config-2:27017" },
{ _id: 2, host: "config-3:27017" }
]
})
// Kiểm tra trạng thái replica set
rs.status()
// Đợi đến khi PRIMARY xuất hiện (khoảng 5-10 giây)
Khởi tạo các Shard Replica Set
Tương tự, khởi tạo Replica Set cho Shard 1:
docker exec -it shard1-1 mongosh --port 27017
rs.initiate({
_id: "shard1-rs",
members: [
{ _id: 0, host: "shard1-1:27017" },
{ _id: 1, host: "shard1-2:27017" },
{ _id: 2, host: "shard1-3:27017" }
]
})
rs.status() // Kiểm tra
Và cho Shard 2:
docker exec -it shard2-1 mongosh --port 27017
rs.initiate({
_id: "shard2-rs",
members: [
{ _id: 0, host: "shard2-1:27017" },
{ _id: 1, host: "shard2-2:27017" },
{ _id: 2, host: "shard2-3:27017" }
]
})
rs.status()
Kết nối Mongos và thêm Shard
Kết nối vào Mongos:
docker exec -it mongos mongosh --port 27017
Thêm các Shard vào cluster:
sh.addShard("shard1-rs/shard1-1:27017,shard1-2:27017,shard1-3:27017")
sh.addShard("shard2-rs/shard2-1:27017,shard2-2:27017,shard2-3:27017")
Kiểm tra trạng thái:
sh.status()
Bạn sẽ thấy danh sách các Shard đã được thêm thành công.
Kịch bản Sharding Collection
Sau khi cluster đã sẵn sàng, kích hoạt sharding cho database và collection:
// Bước 1: Enable sharding cho database
sh.enableSharding("myDB")
// Bước 2: Chọn Shard Key và sharding method
// Ví dụ: Hashed Sharding trên trường user_id
// Với collection rỗng, chỉ định numInitialChunks để phân phối đều ban đầu
sh.shardCollection("myDB.orders", { "user_id": "hashed" }, { numInitialChunks: 4 })
// Bước 3: Kiểm tra phân phối chunk
sh.status()
// Bước 4: Xem thông tin chi tiết về các chunk
use myDB
db.orders.getShardDistribution()
Lưu ý quan trọng: Nếu collection đã có dữ liệu, bạn phải tạo index hỗ trợ Shard Key trước khi chạy sh.shardCollection():
// Tạo index trước khi sharding collection có dữ liệu
db.orders.createIndex({ "user_id": "hashed" })
// Sau đó mới shard (không cần numInitialChunks nếu đã có dữ liệu)
sh.shardCollection("myDB.orders", { "user_id": "hashed" })
⚠️ Lỗi thường gặp: Scatter‑Gather Query

Vấn đề
Lỗi “Scatter‑Gather Query” xảy ra khi câu lệnh truy vấn không chứa Shard Key trong điều kiện .find(). Khi đó, mongos không biết dữ liệu nằm ở Shard nào và buộc phải gửi query đến TẤT CẢ các Shard, sau đó tổng hợp kết quả.
Hậu quả
- Tốn thời gian chờ Shard chậm nhất phản hồi
- Tiêu tốn tài nguyên CPU/IO trên toàn bộ cluster
- Khi số lượng Shard tăng lên, hiệu năng của Scatter‑Gather Query càng tệ hơn
Giải pháp
Rà soát toàn bộ mã nguồn Backend, đảm bảo mọi câu lệnh đọc/ghi quan trọng đều phải chứa Shard Key để Mongos định tuyến trúng đích ngay lập tức.
// ❌ SAI - Scatter-Gather Query (không có Shard Key)
db.orders.find({ "status": "pending" })
// Mongos phải query TẤT CẢ các Shard
// ✅ ĐÚNG - Query có chứa Shard Key (user_id)
db.orders.find({ "user_id": "u12345", "status": "pending" })
// Mongos biết chính xác dữ liệu nằm ở Shard nào
💡 Mẹo: Shard Key của bạn nên được chọn sao cho khớp với pattern truy vấn phổ biến nhất của ứng dụng.
✅ Best Practices
1. Luôn giữ kích thước chunk mặc định (64MB) ổn định
Kích thước chunk mặc định là 64MB. MongoDB tự động split chunk khi vượt quá ngưỡng này và cân bằng (balance) giữa các Shard. Việc thay đổi kích thước chunk có thể gây ra hành vi không mong muốn cho Auto-Balancer.
2. Mỗi Shard Node bắt buộc phải là một Replica Set
Trong Production, mỗi Shard phải là một Replica Set với ít nhất 3 thành viên. Điều này đảm bảo tính High Availability: nếu một node trong Shard bị lỗi, các node khác vẫn tiếp tục phục vụ.
3. Sử dụng Auto-Balancer và AutoMerger (MongoDB 7.0+)
Từ MongoDB 7.0, khi bật Balancer, AutoMerger cũng được kích hoạt tự động. AutoMerger sẽ tự động gộp các chunk nhỏ liền kề, giảm số lượng chunk không cần thiết và cải thiện hiệu năng.
// Bật Balancer (tự động bật AutoMerger từ MongoDB 7.0)
sh.startBalancer()
4. Thiết kế Schema tốt là nền tảng
Trước khi nghĩ đến Sharding, hãy đảm bảo Schema của bạn đã được tối ưu. Việc lựa chọn giữa Embedding và Referencing ảnh hưởng trực tiếp đến hiệu năng sau khi Sharding.
5. Sử dụng Zone Sharding khi cần kiểm soát vị trí dữ liệu
Zone Sharding cho phép bạn gán các khoảng giá trị Shard Key cụ thể vào một subset các Shard. Hữu ích cho các yêu cầu:
- Dữ liệu của khách hàng tại khu vực nào thì lưu trên server gần khu vực đó
- Dữ liệu “nóng” lưu trên hardware mạnh hơn
- Tuân thủ quy định về vị trí dữ liệu (data residency)
6. Cân nhắc Write Concern và Read Preference
- Trong sharded cluster, write concern
"majority"đảm bảo dữ liệu được ghi vào đa số các bản sao trước khi xác nhận, nhưng có thể làm chậm hiệu năng. - Read preference nên được đặt phù hợp: nếu bạn chấp nhận đọc dữ liệu không mới nhất, hãy dùng
secondaryPreferredđể giảm tải cho Primary.
// Ví dụ đặt read preference cho kết nối
db.getMongo().setReadPref('secondaryPreferred')
❓ FAQ
1. Tôi có thể thay đổi Shard Key sau khi collection đã chứa hàng trăm triệu bản ghi không?
Câu trả lời ngắn: Có, nhưng với các ràng buộc nhất định.
- MongoDB 4.2 trở về trước: KHÔNG thể thay đổi Shard Key sau khi đã thiết lập.
- MongoDB 4.4 trở lên: Có thể sử dụng lệnh
refineCollectionShardKeyđể thêm hậu tố (suffix) vào Shard Key hiện có. - MongoDB 5.0 trở lên: Hỗ trợ
reshardCollection– reshard toàn bộ collection với Shard Key mới. - MongoDB 8.0: Resharding nhanh hơn đáng kể, sử dụng ít bộ nhớ hơn và có tùy chọn
forceRedistributionđể phân phối lại dữ liệu trên các Shard mới.
⚠ Ràng buộc quan trọng với refineCollectionShardKey:
- Feature Compatibility Version (FCV) phải là 4.4 trở lên
- Phải có index hỗ trợ Shard Key mới trước khi chạy lệnh
- Nếu index hiện tại có ràng buộc unique, Shard Key mới không thể chỉ định
Ví dụ refine Shard Key:
// Shard Key hiện tại: { "user_id": 1 }
// Refine thành: { "user_id": 1, "region": 1 }
// Bước 1: Tạo index hỗ trợ
db.orders.createIndex({ "user_id": 1, "region": 1 })
// Bước 2: Refine Shard Key
db.adminCommand({
refineCollectionShardKey: "myDB.orders",
key: { "user_id": 1, "region": 1 }
})
2. Chi phí phần cứng tối thiểu để duy trì một cụm Sharding chuẩn Production là bao nhiêu server?
Tối thiểu cho Production: 9 servers (chưa tính mongos)
| Thành phần | Số lượng | Lý do |
|---|---|---|
| Config Servers | 3 | Bắt buộc là Replica Set 3 thành viên |
| Shard 1 | 3 | Replica Set 3 thành viên |
| Shard 2 | 3 | Replica Set 3 thành viên (cần ít nhất 2 Shard để có ý nghĩa) |
Tối thiểu cho Development/Testing: 7 servers
- 3 Config Servers
- 2 Shard, mỗi Shard chỉ 2 thành viên (không khuyến nghị cho Production)
Ngoài ra, bạn cần ít nhất 1 mongos (thường chạy chung với ứng dụng hoặc trên hardware riêng).
💡 Khuyến nghị Production: Triển khai Config Server và Shard Replica Sets trên ít nhất 3 data centers khác nhau để đảm bảo High Availability trong trường hợp một data center gặp sự cố.
3. Làm thế nào để chọn số lượng Shard ban đầu?
Số lượng Shard nên được tính dựa trên:
- Dung lượng dữ liệu dự kiến – mỗi Shard nên giữ khoảng 1-2TB dữ liệu để đảm bảo hiệu năng.
- Tốc độ ghi – nếu ghi 10.000 ops/s, một Shard đơn có thể xử lý được, nhưng nếu cao hơn, cần thêm Shard để phân tán.
- Khả năng mở rộng – nên bắt đầu với 2-3 Shard và tăng dần khi cần.
🏁 Kết luận
MongoDB Sharding là một giải pháp mạnh mẽ để mở rộng hệ thống vượt qua giới hạn phần cứng của một server đơn lẻ. Tuy nhiên, thành công của một cụm Sharded Cluster phụ thuộc phần lớn vào ba yếu tố:
- Hiểu đúng kiến trúc – Phân biệt rõ vai trò của Shard Nodes, Config Servers và Mongos Router
- Chọn Shard Key thông minh – Tránh hotspot, ưu tiên Hashed Sharding cho trường monotonic
- Viết query đúng cách – Luôn kèm Shard Key trong điều kiện truy vấn để tránh Scatter-Gather
🔗 Tham khảo thêm tài liệu:
- MongoDB Official Docs – Sharded Cluster Components – Tài liệu chính thức về kiến trúc Sharded Cluster.
- MongoDB Official Docs – Shard Keys – Hướng dẫn chi tiết về Shard Key và các lựa chọn.
- MongoDB Official Docs – Sharded Cluster Balancer – Quản trị Auto-Balancer.
- MongoDB 8.0 New Features – Percona – Phân tích tính năng mới của MongoDB 8.0 liên quan đến Sharding.