NivaarExam PrepOfficial exam papers ↗

25-Comp-B10 Distributed Systems · May 2014

Question 4 of 6: Operating Systems for Distributed Architectures

Nivaar worked solution (AI-drafted; not reviewed by a licensed engineer)

Notes on this paper

98-Comp-B10 Distributed Systems — National Examinations, May 2014. 3 hours, closed book, non-programmable calculator only. Candidates were instructed to answer any five of the six questions, all carrying equal weight and mostly requiring essay-format answers; all six are answered below as a complete study resource.

Reference texts: Coulouris, Dollimore, Kindberg & Blair, Distributed Systems: Concepts and Design (5th ed.) — system models and client-server architecture (ch. 2), interprocess communication and the request-reply protocol (ch. 4–5), operating system support for distributed systems (ch. 7), distributed file systems (ch. 12–12.4, AFS/NFS), and time, coordination, replication and fault tolerance (ch. 14–15, 18).

Question 4: Operating Systems for Distributed Architectures

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.

Given. Cache hit rate 70% (miss rate 30%); on a hit, a request costs 7 ms of CPU time; on a miss, an additional 10 ms of disk I/O is required on top of the same 7 ms CPU cost; three server configurations to evaluate: single-threaded, two threads on one CPU, and two threads on two CPUs.

Given data
QuantityValue
Cache hit probability0.70
Cache miss probability0.30
CPU time per request (hit or miss)7 ms
Extra disk I/O time on a miss10 ms

Find. The average throughput (requests/second) the server sustains under each of the three threading/processor configurations.

Approach. Model the server as demanding two separable resources per request — CPU time (always 7 ms) and disk time (10 ms, but only on the 30% of requests that miss) — then find each configuration's throughput as the reciprocal of whichever resource-and-concurrency combination limits it: total serial time per request when nothing can overlap, or the busier of the two resources once concurrency lets disk waits overlap with other requests' CPU work.

  1. 1. Single-threaded. With only one thread, the server cannot do anything else while a request is blocked on disk I/O — every request, hit or miss, is serviced strictly serially, so the throughput is limited by the mean total time per request. $$\bar{t}=0.70(7)+0.30(7+10)=4.9+5.1=10.0\ \text{ms}$$ $$\boxed{X_1=\dfrac{1000}{10.0}=100\ \text{requests/s}}$$ Assumption: the single thread blocks for the full disk I/O time with no other useful work possible in that window.
  2. 2. Two threads, one CPU. A miss only blocks the servicing thread on disk (10 ms); with a second thread present, the CPU is not idle during that wait — it can run the other thread's 7 ms of CPU work instead. Every request needs 7 ms of CPU regardless of hit/miss, so the CPU's total demand per request is a constant 7 ms, while the disk's average demand per request is only $0.30(10)=3$ ms. Since $7>3$, the single CPU (not the disk) is the busier, bottleneck resource, and with enough concurrent requests to keep it continuously fed, throughput saturates at the CPU's own service rate. $$\boxed{X_2=\dfrac{1000}{7}\approx 142.86\ \text{requests/s}}$$ Assumption: two threads provide enough concurrency to keep the CPU busy whenever a request is ready for it (the 30% miss rate means disk waits are infrequent enough that this holds).
  3. 3. Two threads, two processors. Doubling the number of processors doubles the CPU-limited capacity, but the (single, physical) disk may become the new bottleneck, so both limits must be checked. CPU-limited capacity with two processors, each independently servicing 7 ms of CPU work per request: $$X_{cpu}=2\times\dfrac{1000}{7}\approx 285.71\ \text{requests/s}.$$ Disk-limited capacity, assuming one disk device serves one I/O operation at a time and every miss needs exactly one 10 ms disk operation: the disk can complete at most $1000/10=100$ operations/s, and since only 30% of requests are misses, $$X_{disk}=\dfrac{100}{0.30}\approx 333.33\ \text{requests/s}.$$ Here the two limits do not coincide: the CPU remains the binding constraint even with a second processor. $$\boxed{X_3=\min(X_{cpu},X_{disk})\approx 285.71\ \text{requests/s}}$$ Assumption: a single, serially-accessed disk device (no RAID striping/multiple disk channels assumed); with it, the disk has more than enough spare capacity (333.33 req/s) that a third processor would still gain throughput (up to the disk's own 333.33 req/s ceiling) before the disk itself finally became the limiting resource.
Final Results — Q4(a)
ConfigurationBottleneckThroughput
1. Single-threadedtotal serial time (no overlap)100.00 req/s
2. Two threads, one CPUCPU (7 ms/request > 3 ms avg. disk)142.86 req/s
3. Two threads, two CPUsCPU (285.71 req/s < disk's 333.33 req/s)285.71 req/s

(b) Thread-per-request vs. worker-pool architecture. In a thread-per-request server, a fresh thread is created for every incoming request and destroyed when that request completes. This gives the simplest programming model (each thread's code reads like a straight-line sequential handler with no shared scheduling state to reason about) and naturally scales concurrency to the offered load, but thread creation/destruction is not free — under a heavy or bursty request rate, the server pays that overhead on every single request and can exhaust kernel resources (thread-table entries, stack memory) if too many requests arrive at once, since there is no built-in cap on how many threads can exist simultaneously. In a worker-pool architecture, a fixed (or bounded, elastically-sized) pool of threads is created once at startup; each incoming request is placed on a shared queue, and idle worker threads pull the next request off that queue and process it, returning to the pool afterward instead of terminating. This amortizes thread-creation cost across the server's whole lifetime and bounds resource consumption predictably (the pool size caps how much concurrent work the server will ever attempt), at the cost of a request occasionally waiting in the queue if all workers are currently busy, and of slightly more complex code (a shared, thread-safe work queue) than the thread-per-request model needs. In practice, high-throughput production servers favour the worker-pool model precisely because of its predictable resource bound under load, reserving thread-per-request for simpler or lower-volume services where its programming simplicity outweighs the creation overhead.

(c) Kernel support for user-level threads. A pure user-level threading implementation (a run-time library, such as an early JVM's "green threads" on UNIX) multiplexes many application-level threads onto one kernel-visible process/task, doing its own scheduling, context-switching and stack management entirely in user space, without the kernel ever being aware that more than one logical thread exists. For this to work at all, the kernel must still provide: (1) some mechanism the library can use to yield control back to itself periodically without kernel help — typically the library relies on cooperative yielding at library call boundaries, or on a periodic timer/alarm signal the kernel delivers to the process, which the library's signal handler intercepts to force a preemptive context switch between user threads; (2) non-blocking (or at least library-interceptable) I/O primitives — the library must be able to substitute its own wrapper around a blocking system call (e.g. read) so that when one user thread would block, the library can switch to another ready user thread instead of the whole kernel-level task blocking; and (3) enough raw process/task abstraction (address space, a stack region the library can carve up and manage itself) for the library to build its own thread control blocks on top of.

Yes, page faults are a genuine problem for user-level threads. Because the kernel schedules the whole process (not the individual user-level threads inside it) as a single unit, a page fault taken by one user thread blocks in the kernel exactly as an ordinary blocking system call would — and since the kernel has no visibility into the library's other, perfectly runnable user threads, the entire process, and therefore every user-level thread within it, is suspended until the fault is serviced, even though most of those threads had nothing to do with the faulting memory access. This defeats one of the main reasons for using threads (letting other work proceed while one thread waits) and is a core reason production systems moved toward kernel-level threads (or a hybrid many-to-many model) rather than pure user-level threading, since a kernel-level scheduler can simply run a different kernel thread while one blocks on a page fault.