In an in-memory data grid used for Space-Based Architecture, what's the practical trade-off between a fully replicated cache topology and a partitioned (sharded) cache topology?
answer
- replicated = full copy everywhere, local reads
- partitioned = sharded by key, local writes
- replicated caps at 1 node's capacity
- partitioned needs routing + rebalancing
- mix both per data set
basics
~20 sReplicated means every node has a full copy of all the data — fast reads everywhere but limited by how much fits on one machine and slower writes. Partitioned means each node holds only a slice — you can hold way more data total, but reads/writes for data you don't hold need a network hop.
solid answer
~50 sReplicated topology keeps a full copy of the entire data set on every node: any node can serve any read purely locally with zero network hop, which is great for read-heavy, relatively small, latency-critical data (e.g., product catalogs, pricing rules). The cost is that every write has to propagate to every replica before it's safely durable, so write throughput and total capacity are both bounded by a single node's memory and by fan-out replication cost. Partitioned (sharded) topology splits the data set across nodes by key, so total capacity scales with the number of nodes and writes only need to reach the partition's own replica set, not the whole cluster — much better for write-heavy or very large data sets. The cost is that a request for data outside the local partition requires a routed hop to the owning node, and you need partition-aware routing plus rebalancing logic. Real systems often mix both: small reference data replicated everywhere, large transactional data partitioned.
go deeper
Should be able to describe the basic difference — full copy everywhere vs. split across nodes — in plain language.
Should explain the capacity and write-throughput implications of each and give an example of when each is appropriate.
Should discuss backup replicas for partition resilience, consistency windows in replicated propagation, and the scatter-gather cost for cross-partition queries.
Should reason about mixed-topology designs within one deployment and how a mismatched topology choice becomes a hard scaling or latency ceiling later.
## The decision every data set forces An in-memory data grid has to decide, for each logical data set, how copies of that data are distributed across the cluster's nodes — and the two canonical topologies, **replicated** and **partitioned**, sit at opposite ends of a capacity-versus-locality trade-off that shapes almost everything downstream about latency and write cost. | | Replicated | Partitioned (sharded) | |---|---|---| | What a node holds | a complete copy of the entire data set | only a slice — a partition — of the whole | | A read | always local, with zero network round trip | needs to reach that key's owning node | | A write | has to be propagated to every other node | needs to reach the node(s) that own that key's partition | | Natural fit | small, relatively static, read-heavy reference data | large, write-heavy, transactional working sets | ## Replicated: every node owns every key In a fully replicated topology, every node in the cluster holds a complete copy of the entire data set. Mechanically, a write on any node has to be propagated to every other node before the grid considers it safely committed (or at least committed with whatever consistency guarantee is configured), typically via some form of broadcast or gossip replication. The immediate benefit is that reads are always local: any processing unit can answer any query against this data instantly, from its own memory, with zero network round trip and zero dependency on which node 'owns' a given key, because every node owns every key. This is why replicated topology is the natural fit for small, relatively static, read-heavy reference data — where the data set is small enough to fit comfortably on every single node and reads vastly outnumber writes: - product catalogs - currency exchange rates - feature flags - pricing/discount rule tables ## Partitioned: each node owns a slice In a partitioned (sharded) topology, the data set is split by key (usually via consistent hashing or a hash-mod-N scheme) so that each node owns only a slice — a partition — of the whole. A write for a given key only needs to reach the node(s) that own that key's partition (the primary and, if configured, its backup replicas), not the entire cluster; a read for a key similarly only needs to reach that key's owning node. The benefit is that total data capacity scales roughly linearly with the number of nodes — add a node, get more aggregate memory for data — and write throughput scales the same way, since writes are no longer fanned out cluster-wide. This is the natural fit for large, write-heavy, transactional working sets: active shopping carts, in-flight orders, live session state, anything where the total data volume would never fit on one machine and where writes are frequent. ## Where the trade-off surfaces The trade-off surfaces in two different places depending on topology. - **For replicated caches, the constraint** is that both total capacity and write throughput are capped by a single node's resources and by full-fan-out replication cost — doubling the cluster size doesn't get you more capacity (every node still needs the whole data set) and doesn't meaningfully speed up writes (they still have to reach every node), so this topology simply doesn't scale for large or write-heavy data no matter how many nodes you add. - **For partitioned caches, the constraint is locality:** a request for a key that isn't on the local node has to be routed across the network to the node that owns it, reintroducing exactly the network hop the in-memory grid was supposed to eliminate for that access. Multi-key operations that span partitions (e.g., 'sum all order totals for this customer' if orders are partitioned by order-id rather than customer-id) become the hardest case, potentially requiring scatter-gather across many nodes, which is slow and defeats much of the locality win. ## Failure modes Failure modes differ correspondingly. - **Replicated topology's failure mode is a consistency window during propagation:** a write acknowledged on node A but not yet propagated to node B means a read on node B returns stale data — a real risk if the replication is asynchronous, and a real cost (added write latency) if it's made synchronous to avoid that staleness. - **Partitioned topology's failure mode is availability of a specific partition:** if the node(s) owning a partition (and its backups) all go down, that specific slice of data becomes unavailable while the rest of the cluster keeps functioning normally — a more contained but still real outage, versus a replicated cluster where any single surviving node keeps serving everything. - **Partitioned topology also has rebalancing risk:** when a node joins or leaves, partitions have to be reassigned and data physically moved, and reads/writes for keys mid-migration can be delayed or, if done wrong, briefly lost. ## Choosing per data set, not globally In practice, production in-memory grids — **Hazelcast**, **Apache Geode/GemFire**, **GigaSpaces XAP** — support both topologies side by side within one deployment, and the design decision is per data set, not global: reference/lookup data goes replicated for zero-hop reads everywhere, while large transactional/session data goes partitioned (often with one or two backup copies per partition for resilience) so capacity and write throughput scale with cluster size. Choosing the wrong topology for a given data set — replicating a large, write-heavy data set, or partitioning a small, read-only lookup table unnecessarily — is a common design mistake that shows up as either a hard capacity ceiling (wrong topology: replicated for big data) or unnecessary network hops on every request (wrong topology: partitioned for tiny data).
- Would you ever partition a small, rarely-changing lookup table like a country-code list?Generally no — partitioning it would only add unnecessary network hops for reads of tiny, cheap-to-replicate data, with no real capacity benefit since the whole table already fits trivially on every node. Replicated topology is the better fit precisely because the data is small and read-heavy.
- How does a partitioned topology usually protect against losing a whole partition when its owning node fails?By configuring one or more backup copies of each partition on different nodes, so a synchronous or near-synchronous replica exists elsewhere in the cluster; on primary failure, the backup is promoted and the partition stays available with minimal data loss, at the cost of extra write latency and memory for the backup copies.
- What's the risk of running a multi-key aggregation query, like summing values for one customer, against a partitioned cache?If the data is partitioned by a key that doesn't align with the query (e.g., partitioned by order ID but querying by customer ID), the values needed can be scattered across many nodes, forcing a scatter-gather operation that hits multiple nodes over the network and is much slower than a single local lookup — this is a case where partition key choice really matters.
Replicated is like every branch library keeping a full copy of every book — instant local browsing, but you can only stock as many titles as your smallest branch can shelve. Partitioned is like a library system where each branch holds a different section of the collection — way more total titles across the system, but you sometimes have to request a book from another branch.
saying these in an interview costs you the question
- Thinks replicated topology scales capacity by adding nodes
- Doesn't mention that partitioned reads/writes for non-local keys require a network hop
- Assumes one topology is strictly 'better' rather than a per-data-set trade-off
- Can't explain why replicated topology is bounded by a single node's memory
- Ignores the need for backup replicas in partitioned topology to survive node failure