The need: sustain an incoming stream
A metrics platform needs to absorb writes, retain useful data and serve reads without concentrating every workload on one instance. Sizing starts with ingestion rate, document size, peaks, retention and actual query patterns.
Architecture: distribute and replicate
Sharding distributes data across multiple shards. Each shard is a replica set: its primary receives writes and secondaries replicate data, enabling an election after a failure. Config servers form their own replica set and hold cluster metadata. Member placement accounts for failure domains and quorum.
Redundant mongos routers
Applications access the cluster through multiple mongos processes. These routers direct operations to the relevant shards; they do not store metrics and do not themselves form a replica set. Client configuration supports multiple routers and reconnection handling. Shard failover and router loss are tested separately.
Choose the shard key from actual workloads
The key must distribute writes and support targeted reads. A strictly increasing timestamp used alone can concentrate new writes on a single range. Cardinality, source distribution, potential hashing and query filters are assessed using representative data. Good ingestion distribution does not guarantee efficient reads.
Sustain capacity over time
Indexes, batched writes, retention and disk throughput influence capacity as much as node count. Monitoring tracks latency, queues, replication lag, oplog window and shard balance. Write concern reflects durability requirements; backup and recovery cover the cluster consistently. Load tests validate performance rather than assuming sharding alone will deliver it.
Key considerations
- Sharding and key selection
- Replica sets and quorum
- Redundant mongos routers
- Ingestion, retention and observability