Triển khai MongoDB Sharded Cluster bằng Docker Compose: Phân mảnh dữ liệu cho ứng dụng lớn

Docker tutorial - IT technology blog
Docker tutorial - IT technology blog

Khi nào MongoDB Sharding thực sự cần thiết?

Khoảng 8 tháng trước, mình đang maintain một hệ thống e-commerce với MongoDB single instance. Mọi thứ chạy ổn cho đến khi bộ dữ liệu vượt 50GB và concurrent users tăng lên 3000+. Query thời gian thực chậm dần — từ 50ms lên 2–3 giây. Vertical scaling (thêm RAM, CPU) chỉ cầm cự được vài tuần trước khi tắc nghẽn trở lại.

Đó là lúc mình bắt đầu nghiêm túc với MongoDB Sharded Cluster. Với môi trường test, Docker Compose là cách nhanh nhất: spin up toàn bộ 10 container trong khoảng 2 phút, reset sạch khi cấu hình sai, không cần thuê thêm máy chủ.

Kiến trúc gồm 3 thành phần

  • Config Servers: Lưu metadata về cluster — shard nào chứa data range nào. Chạy dạng replica set 3 node.
  • Shard Servers: Nơi thực sự lưu data. Mỗi shard là một replica set 3 node để đảm bảo HA.
  • mongos (Query Router): Gateway duy nhất — client connect vào đây, mongos tự route query đến đúng shard.

Cài đặt môi trường

Yêu cầu tối thiểu

  • Docker Engine 24+ và Docker Compose v2
  • RAM tối thiểu 4GB (8GB khuyến nghị cho test thực tế)
  • MongoDB 7.0

Tạo thư mục project và keyfile — bắt buộc để các MongoDB instance xác thực lẫn nhau trong replica set:

mkdir mongo-sharded && cd mongo-sharded
mkdir -p config/keyfile scripts

openssl rand -base64 756 > config/keyfile/mongo-keyfile
chmod 400 config/keyfile/mongo-keyfile

File docker-compose.yml

Toàn bộ cluster gồm: 3 config servers, 6 shard servers (2 shards × 3 node), 1 mongos router.

version: '3.8'

networks:
  mongo-cluster:
    driver: bridge

x-mongo-common: &mongo-common
  image: mongo:7.0
  restart: unless-stopped
  networks:
    - mongo-cluster

services:
  # --- Config Servers ---
  configsvr1:
    <<: *mongo-common
    container_name: configsvr1
    command: mongod --configsvr --replSet configReplSet --port 27017 --keyFile /etc/mongo/keyfile
    volumes:
      - configsvr1_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro
    ports:
      - "27119:27017"

  configsvr2:
    <<: *mongo-common
    container_name: configsvr2
    command: mongod --configsvr --replSet configReplSet --port 27017 --keyFile /etc/mongo/keyfile
    volumes:
      - configsvr2_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro

  configsvr3:
    <<: *mongo-common
    container_name: configsvr3
    command: mongod --configsvr --replSet configReplSet --port 27017 --keyFile /etc/mongo/keyfile
    volumes:
      - configsvr3_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro

  # --- Shard 1 ---
  shard1rs1:
    <<: *mongo-common
    container_name: shard1rs1
    command: mongod --shardsvr --replSet shard1ReplSet --port 27017 --wiredTigerCacheSizeGB 0.5 --keyFile /etc/mongo/keyfile
    volumes:
      - shard1rs1_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro
    ports:
      - "27121:27017"

  shard1rs2:
    <<: *mongo-common
    container_name: shard1rs2
    command: mongod --shardsvr --replSet shard1ReplSet --port 27017 --wiredTigerCacheSizeGB 0.5 --keyFile /etc/mongo/keyfile
    volumes:
      - shard1rs2_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro

  shard1rs3:
    <<: *mongo-common
    container_name: shard1rs3
    command: mongod --shardsvr --replSet shard1ReplSet --port 27017 --wiredTigerCacheSizeGB 0.5 --keyFile /etc/mongo/keyfile
    volumes:
      - shard1rs3_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro

  # --- Shard 2 ---
  shard2rs1:
    <<: *mongo-common
    container_name: shard2rs1
    command: mongod --shardsvr --replSet shard2ReplSet --port 27017 --wiredTigerCacheSizeGB 0.5 --keyFile /etc/mongo/keyfile
    volumes:
      - shard2rs1_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro
    ports:
      - "27122:27017"

  shard2rs2:
    <<: *mongo-common
    container_name: shard2rs2
    command: mongod --shardsvr --replSet shard2ReplSet --port 27017 --wiredTigerCacheSizeGB 0.5 --keyFile /etc/mongo/keyfile
    volumes:
      - shard2rs2_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro

  shard2rs3:
    <<: *mongo-common
    container_name: shard2rs3
    command: mongod --shardsvr --replSet shard2ReplSet --port 27017 --wiredTigerCacheSizeGB 0.5 --keyFile /etc/mongo/keyfile
    volumes:
      - shard2rs3_data:/data/db
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro

  # --- Query Router ---
  mongos:
    <<: *mongo-common
    container_name: mongos
    command: mongos --configdb configReplSet/configsvr1:27017,configsvr2:27017,configsvr3:27017 --port 27017 --keyFile /etc/mongo/keyfile
    volumes:
      - ./config/keyfile/mongo-keyfile:/etc/mongo/keyfile:ro
    ports:
      - "27017:27017"
    depends_on:
      - configsvr1
      - configsvr2
      - configsvr3
    deploy:
      resources:
        limits:
          memory: 1G
        reservations:
          memory: 512M

volumes:
  configsvr1_data:
  configsvr2_data:
  configsvr3_data:
  shard1rs1_data:
  shard1rs2_data:
  shard1rs3_data:
  shard2rs1_data:
  shard2rs2_data:
  shard2rs3_data:

Flag --wiredTigerCacheSizeGB 0.5 trên mỗi shard là bài học từ thực tế. Lần deploy đầu tiên mình bỏ qua — không có giới hạn nào cả. Sau khoảng 3 tiếng chạy load test, mongos ngốn gần hết RAM của host. Mất 2 ngày debug mới tìm ra thủ phạm. Với Docker, memory limit là thứ khai báo từ đầu, không phải sau khi bị sự cố.

Cấu hình chi tiết sau khi khởi động

Khởi động toàn bộ cluster:

docker compose up -d
# Đợi ~30 giây cho các instance ready

Khởi tạo replica set cho Config Servers:

docker exec -it configsvr1 mongosh --port 27017 --eval '
rs.initiate({
  _id: "configReplSet",
  configsvr: true,
  members: [
    { _id: 0, host: "configsvr1:27017" },
    { _id: 1, host: "configsvr2:27017" },
    { _id: 2, host: "configsvr3:27017" }
  ]
})'

Khởi tạo replica set cho Shard 1 và Shard 2:

docker exec -it shard1rs1 mongosh --port 27017 --eval '
rs.initiate({
  _id: "shard1ReplSet",
  members: [
    { _id: 0, host: "shard1rs1:27017" },
    { _id: 1, host: "shard1rs2:27017" },
    { _id: 2, host: "shard1rs3:27017" }
  ]
})'

docker exec -it shard2rs1 mongosh --port 27017 --eval '
rs.initiate({
  _id: "shard2ReplSet",
  members: [
    { _id: 0, host: "shard2rs1:27017" },
    { _id: 1, host: "shard2rs2:27017" },
    { _id: 2, host: "shard2rs3:27017" }
  ]
})'

Đăng ký các shard vào cluster qua mongos:

docker exec -it mongos mongosh --port 27017 --eval '
sh.addShard("shard1ReplSet/shard1rs1:27017,shard1rs2:27017,shard1rs3:27017");
sh.addShard("shard2ReplSet/shard2rs1:27017,shard2rs2:27017,shard2rs3:27017");'

Enable Sharding cho Database và Collection

Sharding không tự động kích hoạt — cần chỉ định rõ database và collection cần shard. Bước quan trọng hơn cả là chọn shard key. Quyết định này gần như không thể đổi lại: một khi collection đã có data, muốn thay shard key là phải dump toàn bộ ra ngoài, xóa collection, rồi import lại từ đầu.

docker exec -it mongos mongosh --port 27017 --eval '
// Enable sharding cho database
sh.enableSharding("myapp");

// Hashed sharding: phân tán đều, phù hợp với query theo ID cụ thể
sh.shardCollection("myapp.orders", { user_id: "hashed" });

// Range-based: hiệu quả hơn cho range query, dễ hotspot hơn
// sh.shardCollection("myapp.logs", { timestamp: 1 });'

Cách phân biệt: dùng hashed khi thường query theo ID cụ thể (user_id: 42), dùng range-based khi query theo khoảng (timestamp từ ngày A đến ngày B). Với orders của e-commerce, hashed theo user_id cho phân tán đều hơn — không có nhóm user nào chiếm quá nhiều traffic. Còn collection logs thì range-based theo timestamp hợp lý hơn, vì phần lớn query kiểu lấy log 7 ngày gần nhất.

Kiểm tra và Monitoring

Xem trạng thái cluster

# Tổng quan cluster và phân phối chunks
docker exec -it mongos mongosh --port 27017 --eval 'sh.status()'

# Chi tiết phân phối data trên từng shard
docker exec -it mongos mongosh --port 27017 --eval '
use myapp;
db.orders.getShardDistribution()'

Test phân tán thực tế

docker exec -it mongos mongosh --port 27017 --eval '
use myapp;
for (let i = 0; i < 10000; i++) {
  db.orders.insertOne({
    user_id: i,
    product: "item_" + Math.floor(Math.random() * 100),
    amount: Math.random() * 1000,
    created_at: new Date()
  });
}
print("Total:", db.orders.countDocuments());
db.orders.getShardDistribution();'

Sau khi insert 10,000 documents, getShardDistribution() sẽ cho thấy data phân tán gần đều giữa 2 shards. Ban đầu cluster chỉ có 1 chunk — MongoDB tự split và migrate khi data tăng.

Phân tích query với explain()

docker exec -it mongos mongosh --port 27017 --eval '
use myapp;
db.orders.find({ user_id: 42 }).explain("executionStats")'

Xem trường queryPlanner.winningPlan.shards trong output: chỉ 1 shard xuất hiện — query đang targeted, mongos biết chính xác cần hỏi shard nào. Cả 2 shard cùng xuất hiện — scatter-gather, mongos hỏi tất cả rồi merge kết quả. Với 2 shard thì overhead chưa đáng kể. Scale lên 8–10 shard mà query vẫn scatter-gather, latency tăng theo vì phụ thuộc vào shard chậm nhất. Lúc đó cần xem lại shard key hoặc bổ sung index phù hợp.

Script health check tổng hợp

#!/bin/bash
# scripts/health-check.sh
echo "=== Config Servers ==="
docker exec configsvr1 mongosh --quiet --eval \
  'rs.status().members.forEach(m => print(m.name, m.stateStr))'

echo "=== Shard 1 ==="
docker exec shard1rs1 mongosh --quiet --eval \
  'rs.status().members.forEach(m => print(m.name, m.stateStr))'

echo "=== Shard 2 ==="
docker exec shard2rs1 mongosh --quiet --eval \
  'rs.status().members.forEach(m => print(m.name, m.stateStr))'

echo "=== Cluster Shards ==="
docker exec mongos mongosh --quiet --eval \
  'sh.status()' | grep -E "(shards|currently|chunks)"

Cluster này đang chạy production được 6 tháng — 10 triệu documents, query response time trung bình dưới 30ms. Hai điều quyết định thành bại: chọn shard key kỹ từ đầu, và khai báo memory limit cho tất cả container. Cả hai đều không phải thứ có thể vá sau — mình đã học theo cách khó nhất.

Share: