Part 6 · 1 chapters · ~8 min

MongoDB Distribution

Replica sets with primaries, secondaries, elections and the oplog, rollback, read preference and read concern, write concern and journaling, causal consistency with sessions, sharding with chunks, the balancer and jumbo chunks, hashed and ranged shard keys and resharding, mongos and config servers, change streams, and cross-shard transactions.

10

Replication and sharding

code
db.transfers.insertOne(doc, { writeConcern: { w: "majority", j: true, wtimeout: 5000 } })
db.transfers.find(q).readConcern("majority").readPref("primaryPreferred")
sh.shardCollection("bank.events", { account_id: "hashed" })            // even writes, scatter-gather ranges
sh.shardCollection("bank.statements", { account_id: 1, month: 1 })     // ranged: targeted queries per account

// change streams: CDC from the oplog, resumable with a token
const cs = db.transfers.watch([{ $match: { operationType: "insert" } }], { fullDocument: "updateLookup" });
for await (const change of cs) { publish(change); saveResumeToken(change._id); }
REPLICA SETS AND SHARDING
a primary per shard, an oplog to secondaries, and mongos routing by shard key
applicationdrivermongos routersconfig serverschunk mapshard 1 replica setP + 2 Sshard 2 replica setP + 2 Sbalancermoves chunks
swipe the figure sideways, or tap expand for full screen
1/5
replica sets
Each replica set has one primary taking writes and secondaries replicating the primary's oplog. If the primary fails, members elect a new one using a Raft-like protocol, usually within seconds.
one primary, secondaries tail the oplogelections in seconds