Skip to content
IT-803 (D) ยท Parallel Computing/Quick Revision Short Notes

Parallel Computing (IT-803 (D)) - Unit 3 Short Notes

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:

  1. SISD (Single Instruction, Single Data): Traditional uniprocessor (e.g., single-core CPU).

  2. SIMD (Single Instruction, Multiple Data): Same instruction on multiple data elements (e.g., GPU cores, vector processors).

  3. MISD (Multiple Instruction, Single Data): Rare; used in some fault-tolerant systems.

  4. 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):

    1. Shared Physical Memory: All processors can directly access a common main memory.

    2. Single Operating System: Typically runs a single OS instance managing all processors.

    3. Tight Coupling: Processors connected via high-speed bus or switch fabric; low communication latency.

    4. Global Clock: All processors are synchronized to a single clock source.

    5. 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

  1. 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).

  2. Many-Core: Shift from few powerful cores to hundreds/thousands of simpler cores.

  3. Memory-Centric: Moving computation closer to memory (e.g., Processing-in-Memory, Near-Data Processing) to overcome the "memory wall."

  4. 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:

      1. Homonym Problem: Same virtual address in different processes maps to different physical pages โ†’ cache inconsistency.

      2. Synonym Problem: Different virtual addresses map to the same physical page โ†’ cache incoherence (multiple copies).

      3. 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):

    1. Write Propagation: Writes to a block must eventually be visible to all caches.

    2. 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() and receive().

  • Communication Types:

    • Point-to-Point: MPI_Send, MPI_Recv between 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:

    1. Pipeline Hazards: Structural (resource conflict), Data (dependencies), Control (branches).

    2. Imbalanced Stages: Slowest stage determines clock cycle time.

    3. Stalls/Bubbles: Due to hazards or synchronization.

    4. 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):

      1. Phase 1 (Aggregation): Processes with non-contiguous requests first exchange data (via MPI) to form larger, contiguous chunks.

      2. 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

  1. Data Parallelism: Same operation applied concurrently to different elements of a data set. (e.g., vector addition, matrix multiply). Natural for SIMD/SIMT.

  2. Task Parallelism: Different tasks (functions/threads) execute concurrently on different data. (e.g., web server handling requests, pipeline stages).

  3. 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

  1. Decomposition: Break the problem into smaller, concurrent tasks (using data, task, or pipeline parallelism).

  2. Mapping: Assign tasks to processes/threads/cores. Consider load balancing and data locality.

  3. Orchestration: Manage interaction: communication (message-passing, shared memory), synchronization (barriers, locks), and ordering.

  4. 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 map followed by a scan into a single, more efficient parallel operation to avoid an extra full pass over the data.

    • Example: map(f, data) then scan(+) is equivalent to scan_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).


Go to where you left off?

Quick Add to Notes

Save questions, your own notes and screenshots into notes filed by unit. It takes a free account.

Create free account

Have an account? Log in