UNIT 3: ADVANCED PARALLEL COMPUTING SYSTEMS AND PROGRAMMING
I. FUNDAMENTALS OF PARALLELISM AND CLASSIFICATION
A. Core Conceptual Distinctions
[!TIP] Exam Focus: These distinctions are very frequently asked as 3-4 mark "differentiate" questions.
| Feature | Parallelism | Pipelining |
|---|---|---|
| Core Idea | Spatial: Multiple tasks/instructions execute simultaneously on different hardware units. | Temporal: A single task is broken into stages; different instructions occupy these stages concurrently over time. |
| Hardware | Requires multiple processing units (cores, processors). | Uses a single processing unit with stage registers. |
| Goal | Increase throughput (tasks per unit time). | Reduce latency (time per task) by overlapping stages. |
| Analogy | Multiple workers building different houses at the same time. | Assembly line: one house moves through framing, wiring, painting stages. |
| Feature | Serial Processing | Parallel Processing |
| :--- | :--- | :--- |
| Execution | Instructions execute one after another on a single processor. | Instructions execute concurrently on two or more processors. |
| Speedup | Limited by clock speed and instruction-level parallelism. | Potential for linear or super-linear speedup with more processors (ignoring overhead). |
| Complexity | Simpler control flow, easier to program/debug. | Complex synchronization, communication, and load balancing. |
| Feature | Control Flow Computer | Data Flow Computer |
| :--- | :--- | :--- |
| Driving Paradigm | Program Counter (PC) driven. Execution follows a predefined instruction sequence. | Data availability driven. An instruction fires when its input data tokens are available. |
| Control Structure | Implicit in the program (if, for, while). | Explicit in the dataflow graph. No PC; inherent parallelism. |
| Synchronization | Managed by the sequential program flow. | Automatic via data token matching. |
| Example | Von Neumann architecture (all modern CPUs). | Experimental architectures; not mainstream. |
B. System Taxonomy
Flynn's Taxonomy classifies computers by instruction and data streams:
-
SISD (Single Instruction, Single Data): Traditional uniprocessor (e.g., single-core CPU).
-
SIMD (Single Instruction, Multiple Data): Same instruction on multiple data elements (e.g., GPU cores, vector processors).
-
MISD (Multiple Instruction, Single Data): Rare; used in some fault-tolerant systems.
-
MIMD (Multiple Instruction, Multiple Data): Most common parallel architecture. Each processor executes its own instruction stream on its own data. Includes multiprocessors and multi-computers.
Key Point for MIMD: A multiprocessor is a single computer with multiple CPUs sharing memory and clock. A multiple-computer system/network is a collection of independent computers connected by a network, with distributed memory and no global clock.
II. PARALLEL ARCHITECTURES AND HARDWARE
A. MIMD Multiprocessors
-
Defining Characteristics (vs. Networks/Clusters):
-
Shared Physical Memory: All processors can directly access a common main memory.
-
Single Operating System: Typically runs a single OS instance managing all processors.
-
Tight Coupling: Processors connected via high-speed bus or switch fabric; low communication latency.
-
Global Clock: All processors are synchronized to a single clock source.
-
Common I/O: Processors share I/O devices.
-
-
Architectural Models:
-
UMA (Uniform Memory Access): Access time to any memory location from any processor is the same. Achieved via a shared bus or crossbar switch. Simpler but not scalable.
-
NUMA (Non-Uniform Memory Access): Access time depends on the location of the memory relative to the processor. Local memory access is fast; remote memory (attached to another processor) is slower. Scalable but requires NUMA-aware programming.
-
B. Graphics Processing Units (GPUs)
-
Architectural Design (SIMD Focus):
-
Many-Core: Hundreds to thousands of smaller, efficient cores (e.g., NVIDIA CUDA cores, AMD stream processors).
-
SIMD Execution: Groups of cores (e.g., Warps in NVIDIA, Wavefronts in AMD) execute the same instruction on different data (SIMT - Single Instruction, Multiple Threads).
-
Deep Memory Hierarchy: Registers per thread โ Shared Memory (per block) โ Global Memory (device) โ Host Memory.
-
Massive Thread Parallelism: Designed to manage thousands of lightweight threads to hide memory latency.
-
-
Application Domains (GP-GPU): Beyond graphics: scientific simulation, machine learning, image processing, financial modelingโany compute-intensive, data-parallel workload.
C. Architectural Characteristics of Future Systems
-
Heterogeneity: Integration of different core types (e.g., high-performance CPU cores + low-power GPU/accelerator cores) on a single chip (e.g., AMD APU, Apple M-series).
-
Many-Core: Shift from few powerful cores to hundreds/thousands of simpler cores.
-
Memory-Centric: Moving computation closer to memory (e.g., Processing-in-Memory, Near-Data Processing) to overcome the "memory wall."
-
Chiplet-Based Design: Using smaller, specialized dies (chiplets) connected via high-speed interconnects (e.g., AMD Infinity Fabric, Intel Foveros) for better yield and flexibility.
III. MEMORY SYSTEMS AND CACHE COHERENCE
A. Cache Hierarchy Design
-
Virtual vs. Physical Addressing:
-
Virtual-Address Cache: Cache is indexed by virtual address. Disadvantages:
-
Homonym Problem: Same virtual address in different processes maps to different physical pages โ cache inconsistency.
-
Synonym Problem: Different virtual addresses map to the same physical page โ cache incoherence (multiple copies).
-
Requires tag comparison on virtual addresses before TLB lookup, complicating design.
-
-
Physical-Address Cache: Indexed by physical address. Solves homonym/synonym issues but requires TLB lookup before cache access, potentially increasing latency. Virtually-indexed, physically-tagged (VIPT) caches are a common compromise.
-
-
Centralized vs. Distributed Shared Caches:
| Centralized Shared Cache (L3) | Distributed Shared Cache | | :--- | :--- | | Single, large cache shared by all cores. | Cache is physically distributed; each core or group has a local portion. | | Pros: Simple coherence, single point for data. | Pros: Lower access latency/contention, better scalability. | | Cons: Contention, bandwidth bottleneck, not scalable to many cores. | Cons: Complex coherence (needs directory or snooping across nodes), non-uniform access. |
B. The Multicache Coherence Problem
-
Challenge: In a shared-memory multiprocessor with private L1/L2 caches, multiple cached copies of the same memory block can exist. If one processor writes to its copy, all other cached copies must be invalidated or updated to maintain a single, coherent view of memory.
-
Coherence Requirements (Fundamental):
-
Write Propagation: Writes to a block must eventually be visible to all caches.
-
Write Serialization: All processors must agree on the order of writes to the same block.
-
-
Two Main Strategies:
-
Write-Invalidate: A write to a shared block invalidates all other cached copies. Subsequent reads fetch the new data from the writer's cache or memory.
-
Write-Update (Write-Broadcast): A write updates all other cached copies with the new value. More traffic but can reduce subsequent read misses.
-
C. Cache Coherence Protocols
-
Snooping-Based (Bus-Based):
-
All caches monitor (snoop) a shared broadcast medium (bus).
-
Each cache has a coherence state per block (e.g., Modified, Exclusive, Shared, Invalid - MESI protocol).
-
On a read/write, caches snoop the bus transaction and act according to their state and the transaction type.
-
Write-Invalidate (Common): Processor issues a "Bus Write" transaction; all other caches invalidate their copy.
-
Pros: Simple, fast for small systems (UMA). Cons: Does not scale (bus broadcast traffic grows with processors).
-
-
Directory-Based:
-
A centralized or distributed directory tracks the state and location (which caches have a copy) of each shared block.
-
On a read/write, the processor queries the directory. The directory coordinates invalidations/updates only to caches that have a copy.
-
Pros: Scalable to large NUMA/many-core systems; reduces unnecessary broadcast traffic. Cons: Higher latency for directory lookup, directory itself can be a bottleneck/point of failure.
-
Comparison: Snooping is implicit (all caches listen to all traffic). Directory is explicit (directory knows who has copies).
IV. PARALLEL PROGRAMMING MODELS AND PARADIGMS
A. Distributed-Memory Programming (MPI)
-
MPI Fundamentals: Standard library (API) for message-passing in distributed memory systems (clusters). Processes have private address spaces; communication via
send()andreceive(). -
Communication Types:
-
Point-to-Point:
MPI_Send,MPI_Recvbetween two specific processes. -
Collective: All processes in a communicator participate.
-
MPI_Bcast: Broadcast from root to all. -
MPI_Reduce: Combine data from all to root (sum, max, etc.). -
MPI_Scatter: Distribute chunks from root to all. -
MPI_Gather: Collect chunks from all to root. -
MPI_Allreduce: Reduce + broadcast result to all. -
MPI_Barrier: Synchronize all processes.
-
-
B. Shared-Memory Programming (Threads)
-
Thread Management:
-
Creation:
pthread_create(POSIX),std::thread(C++11). Creates a new flow of control within the same process (shared memory space). -
Synchronization: Primitives to coordinate access to shared data.
-
Mutex/Lock:
pthread_mutex_lock. Ensures mutual exclusion for a critical section. -
Condition Variable:
pthread_cond_wait. Blocks a thread until a condition is true. -
Semaphore: Counting semaphore for resource pool management.
-
Barrier:
pthread_barrier_wait. All threads wait until all arrive.
-
-
Scheduling: OS thread scheduler maps user threads to kernel threads (1:1, M:1, M:N models). Modern systems typically use 1:1.
-
-
Key Issues:
-
Race Condition: Incorrect outcome due to unpredictable timing of accesses to shared data. Solved by synchronization.
-
Deadlock: Circular wait condition (Thread A holds Lock1 & waits for Lock2, Thread B holds Lock2 & waits for Lock1). Prevention: Lock ordering, timeout.
-
Starvation: A thread is perpetually denied access to a resource.
-
C. Transactional Memory (TM)
-
Transaction vs. Transactional Memory:
-
Transaction (DB): A sequence of operations on a database that satisfies ACID properties.
-
Transactional Memory (TM): A concurrency control mechanism for shared-memory programming. A transaction is a block of code that executes atomically and in isolation with respect to other transactions.
-
-
Properties:
-
Atomicity: Transaction either commits all its writes or aborts and retries (no partial effects).
-
Isolation: Intermediate states of a transaction are invisible to other transactions.
-
-
Conflict Resolution: If two transactions access the same location and at least one is a write, a conflict occurs. System aborts one (usually the younger/smaller) and allows the other to commit. Retry logic is handled by the runtime system.
-
Goal: Replace error-prone lock-based synchronization with simpler, composable atomic blocks.
V. PERFORMANCE EVALUATION AND METRICS
A. Key Performance Metrics
- Speedup ($$\displaystyle S_p $$):
$$S_p = \frac{T_1}{T_p}$$
* $$\displaystyle T_1 $$: Execution time on **1** processor.
* $$\displaystyle T_p $$: Execution time on **p** processors.
* **Ideal Linear Speedup:** $$\displaystyle S_p = p $$.
- Efficiency ($$\displaystyle E_p $$):
$$E_p = \frac{S_p}{p} = \frac{T_1}{p \cdot T_p}$$
* Fraction of time processors are usefully employed. Ranges from 0 to 1.
-
Scalability: How speedup/efficiency change as p increases. A system is scalable if efficiency remains constant or degrades slowly with increasing p.
-
Amdahl's Law (Fixed Workload):
$$S_p \leq \frac{1}{(1 - P) + \frac{P}{p}}$$
* $P$: Proportion of program that is **parallelizable**.
* $1-P$: **Serial fraction**.
* **Limitation:** Speedup is **bounded by the serial part**, even with infinite processors ($$\displaystyle S_\infty = \frac{1}{1-P} $$).
> [!TIP] **Common Pitfall:** Amdahl's Law assumes a **fixed problem size**. It pessimistically limits speedup for large **p**.
- Gustafson's Law (Scaled Workload):
$$S_p = p - \alpha (p - 1) \approx p \text{ for large } p$$
* Assumes problem size **scales with p**. Users solve a larger problem in the same time.
* $\alpha$: Serial fraction of the **scaled** problem. Often much smaller than Amdahl's $(1-P)$.
B. Pipeline Performance Analysis
-
Theoretical Speedup of k-stage Linear Pipeline:
-
Non-pipelined (Serial) Time for n tasks: $$\displaystyle T_{serial} = n \cdot k \cdot \tau $$ (k cycles per task).
-
Pipelined Time for n tasks: $$\displaystyle T_{pipe} = (k + n - 1) \cdot \tau $$ (k cycles for first task, 1 cycle for each subsequent).
-
Speedup ($$\displaystyle S_n $$):
-
$$S_n = \frac{T_{serial}}{T_{pipe}} = \frac{n \cdot k}{k + n - 1}$$
* **As $n \to \infty$:**
$$S_\infty = \lim_{n \to \infty} \frac{n \cdot k}{k + n - 1} = k$$
> \boxed{S_\infty = k}
* **Interpretation:** Maximum possible speedup is **k**, the number of pipeline stages. Achieved when pipeline is **fully utilized** (many tasks).
-
Factors Limiting Real Speedup:
-
Pipeline Hazards: Structural (resource conflict), Data (dependencies), Control (branches).
-
Imbalanced Stages: Slowest stage determines clock cycle time.
-
Stalls/Bubbles: Due to hazards or synchronization.
-
Pipeline Fill/Drain Overhead: First/last tasks don't achieve full throughput.
-
C. I/O Handling in Parallel Systems
-
Challenges:
-
Mismatched Granularity: Parallel programs operate on fine-grained data; I/O is often coarse-grained.
-
Contention: Many processes accessing shared file system simultaneously.
-
Non-Contiguous Access: Data needed by a process may be scattered in the file.
-
Metadata Overhead: Opening/closing files, directory lookups.
-
-
Approaches:
-
Collective I/O: Processes cooperate to perform a single, large, contiguous I/O operation (e.g.,
MPI_File_read_all). Reduces metadata ops and improves disk throughput. -
Two-Phase I/O (Aggregation):
-
Phase 1 (Aggregation): Processes with non-contiguous requests first exchange data (via MPI) to form larger, contiguous chunks.
-
Phase 2 (I/O): A subset of processes (I/O leaders) perform the actual file read/write for their aggregated group.
-
-
Data Sieving: A process reads a larger contiguous chunk from file containing its needed non-contiguous elements, then extracts needed data in memory.
-
VI. PARALLEL ALGORITHM DESIGN AND DATA MANAGEMENT
A. Categories of Parallelism
-
Data Parallelism: Same operation applied concurrently to different elements of a data set. (e.g., vector addition, matrix multiply). Natural for SIMD/SIMT.
-
Task Parallelism: Different tasks (functions/threads) execute concurrently on different data. (e.g., web server handling requests, pipeline stages).
-
Pipeline Parallelism: A task is decomposed into a sequence of subtasks (stages); different data items flow through these stages concurrently. (e.g., GPU graphics pipeline, CPU instruction pipeline).
B. Steps to Design a Parallel Program
-
Decomposition: Break the problem into smaller, concurrent tasks (using data, task, or pipeline parallelism).
-
Mapping: Assign tasks to processes/threads/cores. Consider load balancing and data locality.
-
Orchestration: Manage interaction: communication (message-passing, shared memory), synchronization (barriers, locks), and ordering.
-
Analysis: Model performance (speedup, efficiency), identify bottlenecks (communication, load imbalance), and refine design.
C. Parallel Data Management
-
Data Distribution (for distributed memory):
-
Block Distribution: Contiguous chunk of array assigned to each processor. Good for local computations.
-
Cyclic Distribution: Elements distributed round-robin (p0 gets 0, p, 2p...). Good for load balancing if access pattern is irregular.
-
Block-Cyclic: Blocks of size b distributed cyclically. Combines benefits of both.
-
-
Load Balancing: Ensure all processors have roughly equal work. Static (at design time) vs. Dynamic (at runtime, e.g., work queues).
-
Communication Optimization: Minimize volume and frequency of messages. Overlap communication with computation.
D. Parallel Primitives and Operations
-
Map Operation: Applies a function f independently to each element of a collection.
-
Definition: $$\displaystyle \text{map}(f, [x_1, x_2, ..., x_n]) = [f(x_1), f(x_2), ..., f(x_n)] $$
-
Parallelism: Trivially parallel; no dependencies between elements.
-
Example: Square all elements in an array.
-
-
Scan (Prefix Sum) Operation: Computes all intermediate sums of a binary associative operator (like +).
-
Definition: $$\displaystyle \text{scan}(+, [a,b,c,d]) = [a, a+b, a+b+c, a+b+c+d] $$
-
Parallelism: Has dependencies (each output depends on all previous inputs). Requires efficient parallel algorithms (e.g., Hillis-Steele, Blelloch).
-
-
Fusion of Map and Scan:
-
Concept: Combine a
mapfollowed by ascaninto a single, more efficient parallel operation to avoid an extra full pass over the data. -
Example:
map(f, data)thenscan(+)is equivalent toscan_general(f, +, data)where the operator is modified to incorporate f. -
Benefit: Reduces memory traffic and synchronization overhead, improving performance. A key optimization in parallel libraries (e.g., Thrust, C++ Parallel STL).
-