25-Comp-B10 Distributed Systems · May 2017
Nivaar worked solution (AI-drafted; not reviewed by a licensed engineer)
Question text not reproduced: the examination questions are © Engineers and Geoscientists BC. Open the official past paper (linked at the top of this page) to read the question, then follow the worked solution below.
(a) Why replication is key to the effectiveness of distributed systems. Replication — maintaining multiple copies of the same data or service on independent nodes — addresses two of the central goals a distributed system exists to provide. Performance. Placing replicas close to where they are used (geographically, or simply on separate servers that can each serve requests independently) lets read load be spread across many copies instead of funnelling through one server, cutting both response latency and the load any single node must bear. Availability and fault tolerance. If a service exists as only one copy, that single node is a single point of failure — its crash, or a network partition isolating it, takes the whole service down. With data or a service replicated across independently-failing nodes, clients can be served by any surviving replica, so the system as a whole keeps functioning even while individual replicas are down, which is essential for any service expected to be continuously available despite the partial failures that are normal in a distributed setting. Replication is therefore not an optional performance optimization but a structural requirement for a distributed system to deliver on the availability and scalability promises that justify building it as distributed in the first place — the cost it introduces in return is the need to keep replicas consistent with one another as updates occur.
(b) Active vs. passive replication. Active replication sends every update request to all replica managers, each of which independently executes the same operation on its own copy of the state; because every replica performs the same computation, the replicas stay synchronized without any one of them acting as a coordinator, but this requires the operation to be executed identically and deterministically at every replica (the same inputs applied in the same order must produce the same result everywhere), which typically needs a total-order multicast to deliver every update in the same sequence to every replica. Passive (primary-backup) replication instead designates one replica manager as the primary: clients send requests only to the primary, which executes the operation and then propagates the resulting updated state (or a log of the state changes) to the backup replicas, which simply apply what they are told rather than independently re-executing the operation. If the primary fails, one of the backups is promoted to become the new primary. The key difference is where the work happens and who decides the outcome: active replication relies on every replica computing the same result independently (more resilient to any one replica failing mid-operation, since the others are already doing the same work, but requires strict determinism and ordering), while passive replication relies on a single primary computing the result once and distributing it (simpler to reason about and does not require deterministic operations, but the primary is briefly a single point of failure until a backup takes over, and failover itself takes time and must not lose or duplicate the update in progress).
(c) The gossip architecture. In a gossip-style replicated service, a collection of replica managers (RMs) each accept updates directly from nearby clients (rather than routing every update through one master) and periodically exchange ("gossip") messages with one another, comparing what updates each has seen and forwarding any the other is missing. Over time, updates propagate transitively RM-to-RM until every replica has (eventually) applied every update — this is the classic eventual-consistency model, chosen because it tolerates individual RMs or links being temporarily unreachable (an update simply propagates once connectivity returns) without requiring all RMs to be reachable synchronously the way strict, immediate consensus would.
Each RM keeps two distinct vector timestamps because they answer two different questions. The value timestamp records exactly which updates are reflected in the data value the RM would currently return to a client — it is what a client compares against its own prior timestamp to enforce session guarantees such as "read your writes" or "monotonic reads." The replica timestamp records every update the RM has received (from a client directly, or via gossip from a peer), whether or not it has actually been applied to the value yet — because gossip can deliver an update the RM is not yet ready to apply (e.g. it is causally dependent on an earlier update this RM has not received yet, and updates must be applied in a causally consistent order). Without the replica timestamp, an RM could not correctly tell a gossiping peer "here is everything I have seen so far" (which drives what the peer sends it next), and without the value timestamp it could not correctly tell a client what guarantee its returned data actually satisfies — conflating the two would either apply updates out of causal order (corrupting the value) or under-report to peers what has already been received (causing redundant re-gossip of updates already known).
(d) Read-only replica managers and gossip performance. Yes — making some replica managers read-only can improve the performance of a gossip system. A read-only RM never accepts front-end update operations directly — only reads. This removes two costs a read-write RM must otherwise bear: it never has to withhold applying a gossiped update pending proof that no conflicting update arrived directly from one of its own clients (since none ever do), so it can apply gossiped updates as soon as they are known stable, with less coordination overhead; and it never originates updates of its own that must be gossiped out to every other RM, so it adds no new update traffic to the system, only propagating what it receives. A read-only RM can also answer a local read immediately from its own cached, eventually-consistent copy without coordinating with any other RM at read time. Concentrating writes at a smaller set of read-write RMs while replicating widely via cheap, read-only RMs lets the system scale its read capacity (by adding more read-only replicas close to clients) largely independently of its write throughput and gossip traffic — a scalability trade a design with every RM accepting writes cannot offer as cleanly.