UNIT 4: Advanced Distributed and Parallel Systems
I. Blockchain Technology
A. Cryptographic Foundations
Cryptographic Hash Function
A hash function $H$ that maps arbitrary-sized input to a fixed-size output (hash/digest). Essential for blockchain integrity.
Properties:
-
Deterministic: Same input → same output.
-
Pre-image Resistance: Given hash $h$, infeasible to find input $x$ such that $$\displaystyle H(x)=h $$.
-
Second Pre-image Resistance: Given $$\displaystyle x_1 $$, infeasible to find $$\displaystyle x_2 \neq x_1 $$ with $$\displaystyle H(x_1)=H(x_2) $$.
-
Collision Resistance: Infeasible to find any two distinct inputs $$\displaystyle x_1, x_2 $$ with $$\displaystyle H(x_1)=H(x_2) $$.
-
Avalanche Effect: A single-bit change in input drastically changes output (~50% bits flip).
-
Pseudorandom: Output appears random.
Common Examples: SHA-256 (Bitcoin), SHA-3 (Keccak), RIPEMD-160.
Public Key Cryptography (Asymmetric Cryptography)
Uses a key pair: public key (shared) and private key (secret).
-
Digital Signatures: Sender signs message $M$ with private key $$\displaystyle Priv_S $$ → signature $Sig$. Anyone verifies with public key $$\displaystyle Pub_S $$. Provides authentication, integrity, non-repudiation.
-
Encryption: Sender encrypts message $M$ with receiver's public key $$\displaystyle Pub_R $$. Only receiver with private key $$\displaystyle Priv_R $$ can decrypt.
-
Key Exchange: Protocols like ECDH (Elliptic Curve Diffie-Hellman) to establish shared secret.
HashCash & Bitcoin Proof of Work (PoW)
-
Goal: Sybil attack resistance, spam prevention, leader election in permissionless networks.
-
Mechanism: Find a nonce $n$ such that $$\displaystyle H(header\_data || n) < target $$, where
targetis a dynamically adjusted difficulty value. This is a brute-force search requiring significant computational power (work). -
HashCash (1997): Original anti-spam proposal. Required sender to compute PoW for each email.
-
Bitcoin PoW: Miners compete to find a valid nonce for the current block header. The first to find it broadcasts the block, gets block reward + fees. Difficulty adjusts ~every 2016 blocks to maintain ~10 min block time.
[!TIP] Exam Focus: Be ready to contrast HashCash (email anti-spam) vs Bitcoin PoW (distributed consensus/currency issuance). Emphasize the computational puzzle and difficulty adjustment.
B. Data Structures in Blockchain
Merkle Tree (Hash Tree)
A binary tree where every non-leaf node is the hash of its children's hashes.
Root Hash (Merkle Root)
/ \
Hash(AB) Hash(CD)
/ \ / \
H(A) H(B) H(C) H(D)
/ \ / \ / \ / \
Tx1 Tx2 Tx3 Tx4 Tx5 Tx6 Tx7 Tx8
Importance & Usage:
-
Efficiency (Merkle Proofs): To verify a transaction $$\displaystyle Tx_i $$ is in a block, only need $$\displaystyle \log_2(N) $$ hashes (the Merkle path) instead of all $N$ transactions. Crucial for lightweight clients (SPV nodes).
-
Data Integrity: Any change in a single transaction changes its hash, propagating up to alter the Merkle Root. Root is stored in block header; any tampering is detectable.
-
Consistent Hashing: Enables efficient set difference comparisons (e.g., between two block versions).
Block Structure
A block is a container for transactions, cryptographically linked to its predecessor.
-
Block Header: Contains metadata.
-
Version -
Previous Block Hash(hash of parent block header) → forms the chain -
Merkle Root(root hash of all transactions in this block) -
Timestamp -
Difficulty Target -
Nonce(PoW solution)
-
-
Transaction Counter: Number of transactions.
-
Transactions: List of transaction data (e.g., inputs, outputs, amounts in Bitcoin).
[!DIAGRAM]
DiagramCANVAS: Draw a block with two sections: Header (list fields) and Body (list of transactions). Show an arrow from Previous Block Hash in header to a previous block. Show Merkle Root connecting to a small Merkle Tree diagram inside the body.
C. Blockchain Types and Design
Public vs Private Blockchain
| Feature | Public Blockchain | Private Blockchain |
|---|---|---|
| Access | Permissionless (anyone can join/read/participate) | Permissioned (invitation/authorization required) |
| Control | Decentralized, no single owner | Centralized or consortium-controlled |
| Consensus | Often PoW/PoS (resource-intensive, slow) | Often BFT-based (Raft, PBFT) (faster, efficient) |
| Transparency | Fully transparent (all transactions public) | Varies; can be private/confidential |
| Throughput | Low (e.g., Bitcoin: ~7 TPS) | High (e.g., Hyperledger: 1000s TPS) |
| Example | Bitcoin, Ethereum (mainnet) | Hyperledger Fabric, Corda, R3 Corda |
Permissioned Blockchain Design Issues
-
Identity Management: Must have a robust Membership Service Provider (MSP) to issue/revoke cryptographic identities (X.509 certificates).
-
Access Control & Policies: Fine-grained policies defining who can read/write/endorse which data/channels.
-
Scalability & Performance: Design for known, limited nodes → can use efficient consensus (no Sybil resistance needed).
-
Privacy & Confidentiality: Need channels/private data collections to hide transactions from non-participants.
-
Governance: Clear rules for network upgrades, member onboarding/offboarding.
D. Consensus Mechanisms
Bitcoin Consensus (Nakamoto Consensus)
-
Transaction Propagation: User signs transaction → broadcasts to P2P network.
-
Validation: Nodes validate: syntax, inputs unspent (UTXO), signatures, fees.
-
Mempool: Valid transactions wait in memory pool.
-
Block Creation (Mining): Miner assembles block (transactions + header). Repeatedly modifies
nonce(and extra nonce) to find $$\displaystyle H(header) < target $$. -
Block Propagation: Winning miner broadcasts new block. Other nodes independently verify every aspect (PoW, transactions, Merkle root).
-
Chain Selection: Nodes always adopt the longest valid chain (most cumulative PoW). Temporary forks resolved by next block found on one fork.
-
Transaction Finality: Probabilistic. More confirmations (blocks on top) → lower probability of reversal. 6 confirmations (~1 hour) considered secure.
Proof of Work (PoW) Variants
-
Bitcoin PoW: SHA-256 based, ASIC-friendly. Computationally intensive.
-
HashCash: Original PoW for email. Uses partial hash inversion (leading zeros).
-
Cuckoo Cycle: Memory-hard PoW, ASIC-resistant. Uses graph theory.
Proof of Elapsed Time (PoET)
-
Goal: Leader election in permissioned/trusted execution environments (TEEs like Intel SGX).
-
Mechanism:
-
Each validator requests a random wait time from a trusted execution environment (TEE).
-
TEE generates wait time from a verifiable random function (VRF).
-
Validator with shortest wait time becomes leader for next block.
-
Leader proposes block. Others verify TEE attestation and wait time fairness.
-
-
Advantages: Low energy, high scalability. Disadvantage: Requires trusted hardware (centralization risk if TEE vendor compromised).
Raft Algorithm
-
Goal: Understandable consensus for replicated state machines in non-Byzantine (crash-only) environments.
-
Roles: Leader (handles all client requests, log replication), Follower (passive), Candidate (in election).
-
Key Mechanisms:
-
Leader Election: If follower timeout without heartbeat, becomes Candidate, votes for self. Wins if gets majority votes.
-
Log Replication: Leader appends entry to its log, sends
AppendEntriesRPC to followers. Entry committed when replicated on majority. -
Safety: Only logs from current term can be committed. Leader never overwrites/removes entries.
-
-
Guarantee: Strong consistency (linearizable) if majority of nodes are non-faulty.
Byzantine Fault Tolerance (BFT) - Lamport-Shostak-Pease (Classic BFT)
-
Problem: Reaching consensus in a distributed system where some nodes may be arbitrary/malicious (Byzantine).
-
Assumptions: $n$ total nodes, up to $f$ Byzantine. Requires $$\displaystyle n > 3f $$.
-
Oral Messages Algorithm (Simplified):
-
Recursion: A commander sends value to lieutenants.
-
Recursive Forwarding: Each lieutenant, upon receiving a value, forwards it to all others (excluding sender) in the next "round".
-
Majority Rule: After $f+1$ rounds, each lieutenant decides on the majority value among all received values (including original).
-
-
Result: All loyal lieutenants agree on same value if commander is loyal. If commander is Byzantine, loyal lieutenants may disagree but will not accept a bad value if $$\displaystyle n>3f $$.
-
Practical BFT: PBFT (Practical BFT) is state machine replication variant used in permissioned blockchains (e.g., Hyperledger Fabric).
Distributed Consensus in Closed Environment
Refers to consensus in permissioned/private networks where node identities are known. No Sybil attack. Can use:
-
Crash Fault Tolerant (CFT): Raft, Paxos. Efficient, fast finality. Assumes nodes only crash.
-
Byzantine Fault Tolerant (BFT): PBFT, Tendermint. Tolerates arbitrary/malicious nodes. Slower, more message overhead.
-
Hybrid: Some systems use BFT for ordering, then CFT for execution.
Consensus Protocols for Permissioned Blockchains (Overview)
-
Raft: Used by Hyperledger Fabric (ordering service). CFT, leader-based, fast.
-
PBFT: Used by some Hyperledger Fabric setups, early versions of Stellar. BFT, message-heavy $$\displaystyle O(n^2) $$.
-
Tendermint: BFT consensus engine + application interface. Used by Cosmos SDK. $$\displaystyle O(n^2) $$ messages, instant finality.
-
IBFT (Istanbul BFT): BFT variant used by Quorum (JPMorgan). $$\displaystyle O(n^2) $$.
-
SBFT (Scalable BFT): Attempts to reduce message complexity to $O(n)$.
-
Aardvark: BFT protocol with better DoS resistance.
E. Smart Contracts
Definition & Essential Characteristics
A smart contract is a self-executing, deterministic, tamper-proof program stored on a blockchain that automatically enforces the rules/terms of an agreement when predefined conditions are met.
Essential Characteristics:
-
Self-Executing: Runs automatically without intermediary upon trigger (e.g., transaction, time).
-
Deterministic: Same input (state, transaction) → same output/state change. Must be pure (no random, no external I/O during execution).
-
Immutable: Once deployed, code cannot be changed (unless upgrade pattern built-in).
-
Tamper-Proof: Stored on decentralized ledger; cannot be altered by any single party.
-
Trustless: Parties don't need to trust each other; trust is in the code and network.
-
Stateful: Can read/write to persistent storage (blockchain state).
-
Event-Driven: Typically triggered by transactions.
Bitcoin Script
-
Language: Stack-based, non-Turing complete, no loops, limited operations (cryptographic primitives, arithmetic, conditional).
-
Purpose: Simple transaction validation logic (e.g., multisig, timelocks). Not for complex logic.
-
Limitations:
-
Cannot express complex business logic (no state, loops, complex data).
-
Not Turing complete → no general computation.
-
No native state storage beyond locking/unlocking UTXOs.
-
Minimalist by design for security and verifiability.
-
Ethereum Smart Contracts
-
Writing Process:
-
Write contract in Solidity (high-level, Turing-complete, JavaScript-like) or Vyper.
-
Compile to Ethereum Virtual Machine (EVM) bytecode.
-
Deploy via transaction (paying gas). Deployment creates contract address.
-
Interact by sending transactions to contract address (calling functions).
-
-
Key Features: Turing-complete, rich state variables, events, inheritance, libraries. Gas limits computation/storage to prevent infinite loops/DoS.
-
Execution: Every full node executes contract code to update global state (Ethereum state trie).
Hyperledger Fabric Smart Contracts (Chaincode)
-
Writing Process:
-
Write chaincode in Go, Java, or JavaScript.
-
Package into a deployment spec (chaincode definition).
-
Install on endorsing peers.
-
Instantiate/Approve chaincode on a channel (defines endorsement policy, init args).
-
Invoke via proposal to endorsing peers → simulate → collect endorsements → submit to ordering service → commit to ledger.
-
-
Key Features: Execute-Order-Validate architecture. Chaincode runs in Docker containers (isolated). Private data collections for confidentiality. No native cryptocurrency/gas.
F. Blockchain Platforms
Bitcoin
-
P2P Network: unstructured gossip protocol. ~10k-100k nodes. New nodes discover via DNS seeds.
-
Transaction Processing: UTXO model. Transactions consume unspent outputs, create new ones. Mempool for unconfirmed txs.
-
Mining: PoW (SHA-256). ASIC-dominated. Block reward halves every 210k blocks (~4 years). Max supply 21M BTC.
-
Primary Use: Decentralized digital currency/store of value.
Ethereum
-
Architecture: Account-based (not UTXO). Two types: Externally Owned Accounts (EOA) and Contract Accounts.
-
EVM: Turing-complete virtual machine. Every node executes all contract code.
-
Consensus: Originally PoW (Ethash, memory-hard). Transitioned to PoS (Casper/Beacon Chain) for scalability/energy.
-
Key Features: Smart contracts as first-class citizens. Gas for resource metering. State bloat is a challenge.
-
Use Cases: DeFi, NFTs, DAOs, dApps.
Hyperledger Fabric
-
Architecture: Modular, permissioned. Key components:
-
Peers: Endorsing peers (execute chaincode, simulate), committing peers (store ledger).
-
Ordering Service: Raft/PBFT. Orders transactions into blocks (separate from validation).
-
Channels: Private subnets for confidential transactions.
-
Membership Service Provider (MSP): Manages identities (X.509 certs).
-
Chaincode (Smart Contracts): Business logic.
-
Ledger: World state (LevelDB/CouchDB) + transaction log (blockchain).
-
-
Execute-Order-Validate: Transactions executed (simulated) by endorsers before ordering. Validation (endorsement policy check, MVCC) happens on commit.
-
Use Cases: Supply chain tracking, trade finance, healthcare records, inter-bank settlements.
Ripple (XRP Ledger)
-
Consensus: RPCA (Ripple Protocol Consensus Algorithm). Federated consensus (UNL - Unique Node List). No mining. Fast (~3-5 sec), low-cost.
-
Architecture: Permissioned-like (validators known/trusted). Native cryptocurrency XRP used as bridge currency.
-
Use Cases: Cross-border payments, interbank settlements (xCurrent, On-Demand Liquidity). Focus on financial institutions.
Corda
-
Architecture: Permissioned, peer-to-peer. Not a blockchain (no global broadcast). Data shared only on need-to-know basis.
-
Key Concepts: States (fact on ledger), Contracts (validate state transitions), Flows (business logic for agreement). Notary for uniqueness consensus (prevents double-spend).
-
Use Cases: Trade finance, syndicated loans, insurance, capital markets. Designed for regulated financial institutions.
G. Blockchain Applications and Use Cases
Mortgage Process
-
Traditional: Paper-heavy, multiple intermediaries (lender, broker, title company, appraiser, insurer), slow (30-60 days), fraud risk, lack of transparency.
-
Blockchain-Based:
-
Digitized Assets: Property title as NFT/unique token.
-
Shared Ledger: All parties (lender, borrower, title, escrow) access same immutable record.
-
Smart Contracts: Automate steps: income/asset verification, appraisal trigger, title search, escrow release upon conditions.
-
Benefits: Faster closing (days), reduced fraud, lower costs, transparent audit trail.
-
Supply Chain Finance
-
Improvements:
-
Provenance & Traceability: Track goods from origin to consumer (food safety, luxury goods).
-
Invoice Financing: Suppliers can tokenize invoices on blockchain → financiers can verify authenticity → faster funding.
-
Automated Payments: Smart contracts trigger payment upon IoT sensor confirmation of delivery (e.g., temperature, location).
-
Transparency: All parties see same data → reduces disputes, builds trust.
-
Reduced Costs: Disintermediation of manual reconciliation.
-
Cross-border Payments (Enterprise Application)
-
Problems: Slow (days), expensive (fees, forex), opaque, multiple correspondent banks.
-
Blockchain Solution (e.g., RippleNet):
-
Instant Settlement: Using digital asset (XRP) as bridge currency → settle in seconds.
-
Lower Cost: Fewer intermediaries, lower fees.
-
Transparency: Track payment end-to-end.
-
Use Case: Banks/enterprises sending international payments.
-
Trade Finance (Blockchain-Enabled Trade)
-
Traditional: Letters of Credit (LCs) are paper-based, slow, require manual verification.
-
Blockchain (e.g., we.trade, Marco Polo):
-
Digitized LCs: Smart contracts represent LCs.
-
Automated Compliance: Documents (bill of lading, invoice) hashed on-chain → all parties verify instantly.
-
Faster Settlement: Automated payment upon document matching.
-
Reduced Fraud: Immutable, shared document history.
-
Identity Management Systems (Decentralized Identity - DID)
-
Problems: Centralized providers (Google, Facebook) = single point of failure, data silos, user has no control.
-
Blockchain Solution (e.g., Sovrin, uPort):
-
Self-Sovereign Identity (SSI): User owns identity (DID document stored on-chain/off-chain).
-
Verifiable Credentials (VCs): Issuer (e.g., university) signs credential. User presents VCs to verifier without revealing unnecessary data (zero-knowledge proofs possible).
-
Benefits: Privacy, user control, reduced fraud, portable identity.
-
H. Security and Attacks
Double Spending Problem
-
Explanation: In digital cash, the same digital token/coin can be spent more than once because data can be copied.
-
Prevention in Blockchain:
-
Global Ordering: Consensus orders transactions into a single, agreed-upon sequence.
-
UTXO/Account Model: Bitcoin's UTXO model: an output can only be spent once. Once spent in a confirmed block, it's invalid in any later block.
-
Longest Chain Rule: If two conflicting txs (double spend) are in competing blocks, the one in the longest valid chain wins. The other becomes orphaned.
-
Confirmation Wait: Merchants wait for $k$ confirmations (blocks built on top) before considering payment final. Probability of attacker reversing chain drops exponentially with $k$.
-
Monopoly Problem & Attacks on Proof of Work
-
Monopoly Problem: In PoW, mining power tends to centralize due to economies of scale (ASICs, cheap electricity). Leads to pool centralization (few large pools control >50% hash power).
-
51% Attack:
-
What: An entity controls >50% of total network hash power.
-
Capabilities:
-
Double Spend: Can reorganize recent blocks to replace his own payment with one returning coins to himself (on his private fork).
-
Censor Transactions: Can exclude specific transactions from blocks.
-
Earn All Rewards: Gets all block rewards + fees.
-
-
Cannot: Steal coins from others (need their private keys), create new coins beyond protocol rules, break cryptographic signatures.
-
Economic Disincentive: Attack devalues the cryptocurrency they hold (as attacker). More profitable to be honest.
-
-
Selfish Mining: Miner with >33% hash power withholds blocks to create forks, reducing honest miners' revenue. Can increase revenue even with <50%.
[!TIP] Common Pitfall: 51% attacker cannot create money from thin air or steal coins from wallets. They can only reorganize the chain to reverse their own recent transactions (double spend).
II. Parallel Computing
A. Fundamentals of Parallelism
Parallelism vs Pipelining
| Parallelism | Pipelining | |
|---|---|---|
| Concept | Multiple tasks/instructions executed simultaneously on multiple resources. | Single task broken into sequential stages; different instructions in different stages overlap in time. |
| Hardware | Multiple processors/cores/ALUs. | Single processor with stage registers. |
| Speedup | Up to $p$ times (for $p$ processors) if perfectly parallelizable. | Up to $k$ times (for $k$ stages) if perfectly balanced. |
| Goal | Reduce execution time of a single large job. | Increase throughput (instructions per cycle). |
| Example | 4-core CPU running 4 threads. | 5-stage RISC pipeline (IF, ID, EX, MEM, WB). |
Serial vs Parallel Processing
-
Serial Processing: One instruction executes at a time. $$\displaystyle T_s $$ = sum of all execution times.
-
Parallel Processing: Multiple instructions execute concurrently on multiple processors. $$\displaystyle T_p $$ = time with $p$ processors.
-
Speedup: $$\displaystyle S = T_s / T_p $$. Ideal: $$\displaystyle S = p $$ (linear speedup).
Flynn's Taxonomy (Categories of Parallelism)
| Taxonomy | Description | Example |
|---|---|---|
| SISD | Single Instruction stream, Single Data stream. Classic serial CPU. | Single-core pre-pipeline processor. |
| SIMD | Single Instruction stream, Multiple Data streams. One instruction operates on multiple data elements simultaneously. | Vector processors, GPU cores, MMX/SSE/AVX instructions. |
| MISD | Multiple Instruction streams, Single Data stream. Rare. Different instructions process same data stream (fault tolerance). | Space shuttle flight control (historical). |
| MIMD | Multiple Instruction streams, Multiple Data streams. Most common parallel systems. Each processor has own instruction stream & data. | Multicore CPUs, clusters, NUMA systems. |
Control Flow vs Data Flow Computers
-
Control Flow (Von Neumann): Execution driven by program counter (PC). Instructions executed sequentially (unless branch). Control dependencies dominate. "What to do next?" determined by PC.
-
Data Flow: Execution driven by availability of data. Instructions fire when their input operands are ready. No PC. Data dependencies determine order. Better for fine-grained parallelism but complex scheduling. Used in some signal processing, research architectures.
B. Multiprocessor Systems
Uniprocessor vs Multiprocessor Systems
-
Uniprocessor: Single CPU. Uses pipelining, superscalar, caching for performance. Limited by instruction-level parallelism (ILP).
-
Multiprocessor: Two or more CPUs in close communication, sharing memory/I/O. Exploits task-level parallelism (TLP). Goal: higher throughput, fault tolerance, scalability.
Characteristics of MIMD Multiprocessors (vs multiple computer systems/networks)
-
Shared Physical Memory: Processors can access a common main memory (UMA/NUMA). Networks of computers (clusters) typically have distributed memory (each node has private memory).
-
High Bandwidth, Low Latency Interconnect: Specialized buses (e.g., Intel UPI, AMD Infinity Fabric) or crossbars. LANs (Ethernet) have higher latency.
-
Tight Coupling: Processors synchronize frequently (locks, barriers). Clusters have looser coupling (message passing).
-
Single System Image: OS sees system as one machine (single kernel). Clusters often run separate OS instances (though middleware can provide single image).
-
Coherence Required: With shared memory, cache coherence is mandatory. Clusters don't share memory → no coherence problem.
Distributed-Memory vs Shared-Memory Systems
| Shared-Memory | Distributed-Memory | |
|---|---|---|
| Memory | Global address space (physically shared or logically distributed). | Each processor has private local memory. |
| Communication | Implicit via loads/stores to shared addresses. | Explicit via message passing (MPI). |
| Programming | Easier (shared variables), but need synchronization (locks, atomics). | Harder (explicit send/receive), but no false sharing/coherence overhead. |
| Scalability | Limited by memory bus/coherence overhead (typically < 64 cores). | Highly scalable (1000s of nodes). |
| Examples | Multicore CPUs, SMP, NUMA systems. | Clusters, supercomputers (Cray, IBM Blue Gene). |
C. Memory Hierarchy and Coherence
Cache Memory
-
Purpose: Bridge speed gap between CPU and main memory.
-
Virtual vs Physical Addressing:
-
Physically Indexed, Physically Tagged (PIPT): Index cache using physical address. Simple, but requires address translation (TLB lookup) before cache access → slower access time.
-
Virtually Indexed, Physically Tagged (VIPT): Index using virtual address, tag compare with physical tag. Allows parallel TLB+cache access. Limitation: Cache size ≤ page size * associativity to avoid aliasing (same virtual address mapping to different physical pages).
-
Virtually Indexed, Virtually Tagged (VIVT): Simplest, but has aliasing (same data at different virtual addresses cached separately) and synonym (same physical data cached in multiple lines) problems. Disadvantages: OS must flush cache on context switch, complex on page migration.
-
-
Disadvantages of Virtual Address Caches:
-
Alias Problem: Two different virtual addresses map to same physical address → cached twice → coherence/consistency issues.
-
Synonym Problem: Same physical address mapped to multiple virtual addresses → updates via one VA not seen via other VA.
-
Context Switch Overhead: May require cache flushes on process switch (if ASID not used).
-
Complex Coherence: Requires extra hardware (process ID tags) or software management.
-
Shared Cache Architectures
-
Centralized (Shared) Cache: Single cache (L3) shared by all cores. On-chip, high bandwidth. Simplifies coherence (snooping possible). Contention can be bottleneck.
-
Distributed Cache: Each core/node has its own cache (L2/L3). Lower contention, but coherence harder (need directory or snooping across caches). Common in NUMA systems (e.g., AMD Zen: each CCX has shared L3, but separate L3s across CCX).
Cache Coherence Problem (Multicache Coherence)
In systems with multiple caches holding copies of the same memory block, ensure that all caches see a consistent view of memory. Specifically:
-
Write Propagation: Writes to a cache line must eventually propagate to other caches.
-
Write Serialization: All processors must agree on order of writes to same location.
-
Two Key Inconsistencies:
-
Stale Read: Processor reads old value because another's write hasn't propagated.
-
Non-Atomic Write: Two processors write concurrently → final value unpredictable.
-
Coherence Protocols
-
Snooping-Based (Bus-Based Systems):
-
All caches monitor (snoop) a shared bus for memory transactions.
-
Write-Invalidate: On write, bus transaction invalidates copies in other caches. (e.g., MESI protocol: Modified, Exclusive, Shared, Invalid states).
-
Write-Update (Write-Broadcast): On write, bus transaction updates all other caches' copies. More bus traffic.
-
Advantages: Simple, fast for small systems (≤ 8-16 cores).
-
Disadvantages: Bus traffic scales with $O(p)$; doesn't scale to many processors.
-
-
Directory-Based (Scalable Systems):
-
Centralized or distributed directory tracks which caches have a copy of each block.
-
On access, query directory → get list of sharers → send invalidates/updates only to those caches.
-
Advantages: Traffic scales with number of sharers, not total processors. Scales to 1000s of cores.
-
Disadvantages: Directory access latency, storage overhead, complexity.
-
[!DIAGRAM]
DiagramCANVAS: Draw a 4-core system with a shared bus. Show each core with L1 cache. Show one core writing to address X, and arrows showing invalidate messages to other caches via bus (snooping).
DiagramCANVAS: Draw a distributed system with 4 nodes, each with private cache and memory. Show a central directory node. Show Node 1 writing to address X, querying directory for sharers, then sending invalidate messages only to Node 2 and 3 (sharers).
D. Parallel Programming Models
Steps to Design Parallel Programs
-
Decomposition: Break problem into concurrent tasks.
-
Task Parallelism: Different tasks (different functions) run concurrently.
-
Data Parallelism: Same task applied to different data elements (e.g., array operations).
-
-
Mapping: Assign tasks to processors (static/dynamic). Consider load balancing, data locality.
-
Orchestration: Manage interaction/synchronization between tasks (communication, synchronization).
-
Execution: Run on target hardware (threads, processes, MPI ranks).
-
Optimization: Tune for load balance, reduce communication, improve locality.
Distributed-Memory Programming with MPI (Message Passing Interface)
-
Model: Processes with private address spaces. Communicate via explicit
send/receivecalls. -
Key Concepts:
-
Communicator: Group of processes (e.g.,
MPI_COMM_WORLD). -
Rank: Unique ID within communicator ($0$ to $p-1$).
-
Point-to-Point:
MPI_Send,MPI_Recv(blocking/non-blocking). -
Collective:
MPI_Bcast(broadcast),MPI_Reduce(sum/max/etc.),MPI_Scatter/MPI_Gather,MPI_Barrier. -
Datatypes: Define custom data layouts for non-contiguous data.
-
-
Typical Pattern: Master-worker, SPMD (Single Program Multiple Data).
-
Example (Vector Add):
MPI_Scatter(A, n/p, MPI_DOUBLE, local_A, n/p, MPI_DOUBLE, 0, MPI_COMM_WORLD); for(i=0; i<n/p; i++) local_C[i] = local_A[i] + B[i]; // B is broadcasted MPI_Gather(local_C, n/p, MPI_DOUBLE, C, n/p, MPI_DOUBLE, 0, MPI_COMM_WORLD);
Thread Management
-
Creation:
pthread_create(POSIX),std::thread(C++11), OpenMP#pragma omp parallel. -
Synchronization:
-
Mutex/Lock:
pthread_mutex_lockfor critical sections. -
Semaphore: Counting semaphore for resource pools.
-
Condition Variable:
pthread_cond_wait/signalfor event-based sync. -
Barrier:
pthread_barrier_waitor OpenMPomp barrier.
-
-
Scheduling: OS scheduler maps threads to CPU cores. Can be:
-
Coarse-grained: Long-running threads, minimal sync.
-
Fine-grained: Many short tasks, frequent sync (overhead).
-
Work Stealing: Idle threads steal tasks from busy threads' queues (e.g., OpenMP, TBB).
-
Parallel Data Management
-
Distribution: How to partition data across processors.
-
Block: Contiguous chunks (good for spatial locality).
-
Cyclic/Stride: Round-robin assignment (good for load balance if uneven work).
-
Block-Cyclic: Hybrid.
-
-
Partitioning: Align data layout with computation to minimize communication.
-
Owner-Computes Rule: Compute where data resides.
-
Ghost/Shadow Regions: For stencil computations, replicate border data to avoid communication in inner loops.
-
-
Aggregation: Combine small messages into larger ones to reduce communication overhead.
I/O Handling in Parallel Programming Languages
-
Challenge: Multiple processes/threads writing to same file → race conditions, interleaved output.
-
Solutions:
-
Collective I/O: Processes coordinate (via MPI-IO) to perform large, contiguous reads/writes. OS/parallel file system (Lustre, GPFS) optimizes.
-
Independent I/O: Each process opens/writes its own file (e.g.,
output.rank000). Simple, no coordination. -
Synchronized I/O: Use locks/barriers so only one process writes at a time. High overhead.
-
Two-Phase I/O: Aggregation phase (processes combine requests), I/O phase (subset of processes perform actual I/O).
-
Map and Scan Operations (Parallel Patterns)
-
Map: Apply same function $f$ independently to each element of a collection. Embarrassingly parallel.
result = [f(x) for x in data] # Parallelizable across elements -
Scan (Prefix Sum): Given binary operator $\oplus$, compute all prefixes: $$\displaystyle [a_0, a_0\oplus a_1, a_0\oplus a_1\oplus a_2, ...] $$.
-
Example: Exclusive scan: $$\displaystyle [0, a_0, a_0\oplus a_1, ...] $$.
-
Uses: Histogram, sorting, stream compaction, polynomial evaluation.
-
Parallel Algorithm: Hillis-Steele (incl. scan, $O(\log n)$ steps, $O(n \log n)$ work) or Blelloch (excl. scan, $O(\log n)$ steps, $O(n)$ work).
-
-
Fusing Map and Scan: Combine operations to reduce passes over data.
-
Example:
map-reduceis a fused pattern: map → local results → reduce (tree-based) → global result. -
Benefit: Reduces memory traffic, synchronization overhead.
-
E. Performance Evaluation
Performance Metrics
-
Execution Time: $$\displaystyle T_p $$ (parallel), $$\displaystyle T_s $$ (serial). Wall-clock time.
-
Speedup: $$\displaystyle S_p = T_s / T_p $$. Linear speedup if $$\displaystyle S_p = p $$.
-
Efficiency: $$\displaystyle E_p = S_p / p = T_s / (p T_p) $$. Fraction of time processors are usefully busy.
-
Scalability: How $$\displaystyle S_p $$ or $$\displaystyle E_p $$ changes with $p$. Strong scaling: Fixed problem size, increase $p$. Weak scaling: Problem size per processor fixed, increase $p$.
-
Cost: $$\displaystyle C_p = p \times T_p $$. Optimal if $$\displaystyle C_p \approx T_s $$ (cost of parallel solution ≈ serial solution).
Amdahl's Law
-
Statement: Maximum speedup from parallelization is limited by sequential fraction of the program.
-
Formula: Let $f$ be fraction of code that is inherently sequential ($0 \leq f \leq 1$). Then:
$$S_p \leq \frac{1}{f + \frac{(1-f)}{p}}$$
-
Implications:
-
As $p \to \infty$, $$\displaystyle S_p \to 1/f $$.
-
If $$\displaystyle f=0.1 $$ (10% sequential), max speedup = 10×, no matter how many processors.
-
Critique: Assumes fixed problem size. Often Gustafson's Law is more realistic (scaled problem size).
-
Gustafson's Law (Scaled Speedup)
-
Assumes problem size scales with number of processors to maintain constant per-processor work.
-
Formula: Let $s$ be sequential fraction of scaled problem. Then:
$$S_p = p - s(p-1)$$
- Implication: If $s$ is small, near-linear speedup possible with increasing $p$ and problem size.
Pipeline Speedup Theory (Proof)
-
Non-pipelined (Serial) Processor: Execute $n$ instructions. Each takes $$\displaystyle t_{inst} $$ time. $$\displaystyle T_s = n \cdot t_{inst} $$.
-
$k$-stage Pipeline: First instruction takes $k \cdot \Delta t$ (stage time). Subsequent instructions complete every $\Delta t$. $$\displaystyle T_p = k\Delta t + (n-1)\Delta t $$.
-
Speedup:
$$S = \frac{T_s}{T_p} = \frac{n \cdot t_{inst}}{k\Delta t + (n-1)\Delta t}$$
- Ideal Case: $$\displaystyle t_{inst} = k\Delta t $$ (each instruction perfectly balanced across $k$ stages).
$$S_{ideal} = \frac{n \cdot k\Delta t}{k\Delta t + (n-1)\Delta t} = \frac{nk}{k + n - 1}$$
- As $n \to \infty$:
$$\lim_{n\to\infty} S_{ideal} = \frac{nk}{n} = k$$
- Conclusion: A $k$-stage linear pipeline can be at most $k$ times faster than a non-pipelined processor, for large $n$. Achieved when pipeline is perfectly balanced and $n \gg k$.
F. Advanced Architectures and Technologies
GPU Architecture
-
Design Philosophy: Throughput-oriented. Thousands of simple cores to handle massive data-parallel workloads (graphics, AI, scientific computing).
-
Key Components:
-
Streaming Multiprocessors (SMs): Each SM contains many CUDA cores (ALUs), registers, shared memory, schedulers.
-
Massive Parallelism: 1000s of threads (warps of 32 threads) scheduled in SIMT (Single Instruction, Multiple Threads) fashion.
-
Memory Hierarchy: Registers (per thread) → Shared Memory (per SM, low latency) → L1/L2 Caches → Global Memory (DRAM, high bandwidth, high latency).
-
High Bandwidth Memory: GDDR6, HBM2e (1024-bit+ bus).
-
-
Programming Model: Offload parallel kernels (SIMD-like) to GPU. Host (CPU) manages data transfer and kernel launch.
-
Applications:
-
Graphics Rendering: Vertex/fragment shading.
-
AI/ML: Dense matrix multiplies (tensor cores), neural network training/inference.
-
Scientific Computing: Molecular dynamics, fluid dynamics (stencil codes).
-
Cryptocurrency Mining: Hash computations (historically).
-
Transactional Memory (TM) vs Traditional Database Transactions
| Transactional Memory (TM) | Database Transactions (ACID) | |
|---|---|---|
| Scope | In-memory, within a single program/process (shared memory). | Persistent storage (disk), across processes/systems. |
| Granularity | Memory accesses (loads/stores). | Database operations (SQL statements). |
| Isolation Level | Typically serializable (conflict detection/avoidance). | Configurable (Read Uncommitted to Serializable). |
| Durability | Not required. Crash loses transactional updates. | Required. Write-ahead logging (WAL) ensures durability. |
| Conflict Detection | At runtime (read/write sets). | At commit time (locking, timestamps). |
| Goal | Simplify concurrent programming (replace locks). | Ensure data integrity in persistent storage. |
- TM Types: Hardware TM (HTM) (Intel TSX), Software TM (STM), Hybrid TM.
Architectural Characteristics of Future Systems
-
Heterogeneous Integration: CPUs + GPUs + FPGAs + AI accelerators (TPUs, NPUs) on package/chip (e.g., AMD Ryzen with RDNA graphics, Intel Foveros).
-
Memory-Centric Computing: Process-in-memory (PIM), near-data processing to reduce data movement (the "memory wall").
-
Chiplet-Based Design: Disaggregated dies (CPU, I/O, memory) connected via high-speed interconnects (UCIe). Improves yield, cost, flexibility.
-
Advanced Packaging: 2.5D/3D stacking (HBM, L4 cache on logic die).
-
Security-First: Hardware-enforced isolation (Intel SGX, AMD SEV, ARM TrustZone), side-channel attack mitigations.
-
Energy Efficiency: Dark silicon, power gating, near-threshold computing.
-
Quantum-Classical Hybrid: Classical processors controlling quantum co-processors (NISQ era).
-
Neuromorphic Computing: Event-driven, brain-inspired architectures for sparse, asynchronous workloads.
-
Challenges: Programming Model Complexity (heterogeneity), Memory Bandwidth Wall, Thermal Density, Security vs Performance Trade-offs, Software Ecosystem Maturity.
[!TIP] Exam Focus: For "architectural characteristics," focus on heterogeneity, memory-centric, chiplets, security, and energy. Be ready to contrast TM (volatile, in-memory) vs DB transactions (persistent).