Introduction to Concurrent Programming with GPUs

Tasks

Assignments & Tests

Concurrent Programming Problems

Dining Philosophers Problem

Premise

Solution

import random
import threading
import time


class Philosopher(threading.Thread):

    running = True

    def __init__(self,
                 xname: str,
				 fork_on_left: threading.Lock,
			     fork_on_right: threading.Lock):
        threading.Thread.__init__(self)
        self.name: str = xname
        self.fork_on_left: threading.Lock = fork_on_left
        self.fork_on_right: threading.Lock = fork_on_right

    def run(self):
        while self.running:
            #  Philosopher is thinking (but really is sleeping).
            time.sleep(random.uniform(3, 13))
            print(f'{self.name} is hungry.')
            self.dine()

    def dine(self):
        fork1, fork2 = self.fork_on_left, self.fork_on_right

        while self.running:
            fork1.acquire(True)
            locked = fork2.acquire(False)
            if locked:
                break
            fork1.release()
            print(f'{self.name} swaps forks')
            fork1, fork2 = fork2, fork1
        else:
            return

        self.dining()
        fork2.release()
        fork1.release()

    def dining(self):
        print(f'{self.name} starts eating')
        time.sleep(random.uniform(1, 10))
        print(f'{self.name} finishes eating and leaves to think.')


def dining_philosophers():

    forks = [threading.Lock() for n in range(5)]
    philosopher_names = ('Aristotle',
						 'Kant',
						 'Buddha',
						 'Marx',
						 'Russel')

    philosophers = [Philosopher(philosopher_names[i],
							    forks[i % 5],
							    forks[(i+1) % 5]) for i in range(5)]

    random.seed(507129)
    Philosopher.running = True
    for p in philosophers:
        p.start()
    time.sleep(100)
    Philosopher.running = False
    print("Now we're finishing.")

dining_philosophers()

Explanation

Deadlock Avoidance Strategy

The key to avoiding deadlock is in the dine method:

Producer-Consumer Problem

Premise

Solution

#include <iostream>
#include <thread>
#include <deque>
#include <mutex>
#include <chrono>
#include <condition_variable>

using std::deque;
std::mutex mu,cout_mu;
std::condition_variable cond;

class Buffer
{
public:
    void add(int num) {
        while (true) {
            std::unique_lock<std::mutex> locker(mu);
            cond.wait(locker, [this](){return buffer_.size() < size_;});
            buffer_.push_back(num);
            locker.unlock();
            cond.notify_all();
            return;
        }
    }
    int remove() {
        while (true)
        {
            std::unique_lock<std::mutex> locker(mu);
            cond.wait(locker, [this](){return buffer_.size() > 0;});
            int back = buffer_.back();
            buffer_.pop_back();
            locker.unlock();
            cond.notify_all();
            return back;
        }
    }
    Buffer() {}
private:
    deque<int> buffer_;
    const unsigned int size_ = 10;
};

class Producer
{
public:
    Producer(Buffer* buffer, std::string name)
    {
        this->buffer_ = buffer;
        this->name_ = name;
    }
    void run() {
        while (true) {
            int num = std::rand() % 100;
            buffer_->add(num);
            cout_mu.lock();
            int sleep_time = rand() % 100;
            std::cout << "Name: " << name_ << "   Produced: " << num << "   Sleep time: " << sleep_time << std::endl;
            std::this_thread::sleep_formilliseconds(sleep_time);
            cout_mu.unlock();
        }
    }
private:
    Buffer *buffer_;
    std::string name_;
};

class Consumer
{
public:
    Consumer(Buffer* buffer, std::string name)
    {
        this->buffer_ = buffer;
        this->name_ = name;
    }
    void run() {
        while (true) {
            int num = buffer_->remove();
            cout_mu.lock();
            int sleep_time = rand() % 100;
            std::cout << "Name: " << name_ << "   Consumed: " << num << "   Sleep time: " << sleep_time << std::endl;
            std::this_thread::sleep_formilliseconds(sleep_time);
            cout_mu.unlock();
        }
    }
private:
    Buffer *buffer_;
    std::string name_;
};

int main() {
    Buffer b;
    Producer p1(&b, "producer1");
    Producer p2(&b, "producer2");
    Producer p3(&b, "producer3");
    Consumer c1(&b, "consumer1");
    Consumer c2(&b, "consumer2");
    Consumer c3(&b, "consumer3");

    std::thread producer_thread1run, &p1;
    std::thread producer_thread2run, &p2;
    std::thread producer_thread3run, &p3;

    std::thread consumer_thread1run, &c1;
    std::thread consumer_thread2run, &c2;
    std::thread consumer_thread3run, &c3;

    producer_thread1.join();
    producer_thread2.join();
    producer_thread3.join();
    consumer_thread1.join();
    consumer_thread2.join();
    consumer_thread3.join();

    getchar();
    return 0;
}

This C++ code demonstrates a classic Producer-Consumer pattern using multithreading.

Explanation

  1. Global Variables:
    • std::mutex mu: A mutex to protect shared access to the Buffer.
    • std::mutex cout_mu: A mutex to protect std::cout, ensuring console output from different threads doesn't get mixed up.
    • std::condition_variable cond: Used to signal between threads. Producers wait on it if the buffer is full, and consumers wait if it's empty.
  2. Buffer Class:
    • Manages a shared std::deque<int> buffer_ (a double-ended queue) with a fixed maximum size_ of 10.
    • void add(int num) (Producer's method):
      • Acquires a std::unique_lock<std::mutex> locker(mu) on the buffer's mutex.
      • cond.wait(locker, [this](){return buffer_.size() < size_;});:
        • The thread waits (releases the lock mu and sleeps) until the condition buffer_.size() < size_ (buffer is not full) is true.
        • The lock is reacquired automatically before checking the condition and upon waking up.
      • Adds num to the buffer_.
      • locker.unlock(): Releases the lock explicitly.
      • cond.notify_all(): Notifies all waiting threads (both producers and consumers) that the buffer state has changed.
      • The outer while(true) loop and return ensure the method completes after one successful addition.
    • int remove() (Consumer's method):
      • Acquires a std::unique_lock<std::mutex> locker(mu).
      • cond.wait(locker, [this](){return buffer_.size() > 0;});:
        • The thread waits until buffer_.size() > 0 (buffer is not empty).
      • Retrieves and removes an item from the buffer_.
      • locker.unlock(): Releases the lock.
      • cond.notify_all(): Notifies all waiting threads.
      • Returns the removed item.
      • The outer while(true) loop and return ensure the method completes after one successful removal.
  3. Producer Class:
    • Stores a pointer to the shared Buffer and a name_.
    • void run():
      • Enters an infinite loop.
      • Generates a random number (num).
      • Calls buffer_->add(num) to put the item into the buffer.
      • Locks cout_mu to print producer information and a random sleep duration.
      • Simulates work by sleeping for a random time.
      • Unlocks cout_mu.
  4. Consumer Class:
    • Stores a pointer to the shared Buffer and a name_.
    • void run():
      • Enters an infinite loop.
      • Calls buffer_->remove() to get an item from the buffer.
      • Locks cout_mu to print consumer information and a random sleep duration.
      • Simulates work by sleeping for a random time.
      • Unlocks cout_mu.
  5. main() Function:
    • Creates a single Buffer instance b.
    • Creates three Producer instances (p1, p2, p3) and three Consumer instances (c1, c2, c3), all sharing the same buffer b.
    • Launches six threads: one for each producer's run() method and one for each consumer's run() method.
    • producer_threadX.join() and consumer_threadX.join(): The main thread waits for these threads to finish. However, because the run() methods contain infinite while(true) loops, these threads will never actually finish on their own. The program will run indefinitely until manually terminated (e.g., Ctrl+C).
    • getchar(): This line is effectively unreachable due to the indefinite joins.

Synchronization Explained

Sleeping Barber Problem

Premise

Solution

# Based on code from https://github.com/Nohclu/Sleeping-Barber-Python-3.6-/blob/master/barber.py
import time
import random
import threading
from queue import Queue

CUSTOMERS_SEATS = 15        # Number of seats in BarberShop
BARBERS = 3                # Number of Barbers working
EVENT = threading.Event()   # Event flag, keeps track of Barber/Customer interactions
Earnings = 0
SHOP_OPEN = False

class Customer(threading.Thread):       # Producer Thread
    def __init__(self, queue):          # Constructor passes Global Queue (all_customers) to Class
        threading.Thread.__init__(self)
        self.queue = queue
        self.rate = self.what_customer()

    @staticmethod
    def what_customer():
        customer_types = ["adult", "senior", "student", "child"]
        customer_rates = {"adult": 16,
                          "senior": 7,
                          "student": 10,
                          "child": 7}
        t = random.choice(customer_types)
        print(t + " rate.")
        return customer_rates[t]

    def run(self):
        if not self.queue.full():  # Check queue size
            EVENT.set()  # Sets EVENT flag to True i.e. Customer available in the Queue
            EVENT.clear()  # A lerts Barber that their is a Customer available in the Queue
        else:
            # If Queue is full, Customer leaves.
            print("Queue full, customer has left.")

    def trim(self):
        global Earnings
        print("Customer haircut started.")
        a = 3 * random.random()  # Retrieves random number.
        time.sleep(a)
        payment = self.rate
        # Barber finished haircut.
        print("Haircut finished. Haircut took {}".format(a))
        Earnings += payment


class Barber(threading.Thread):     # Consumer Thread
    def __init__(self, queue):      # Constructor passes Global Queue (all_customers) to Class
        threading.Thread.__init__(self)
        # TODO set this class's queue property to the passed value
        self.queue = queue
        self.sleep = True   # No Customers in Queue therefore Barber sleeps by default

    def is_empty(self):  # Simple function that checks if there is a customer in the Queue and if so
        if self.queue.empty():
            self.sleep = True   # If nobody in the Queue Barber sleeps.
        else:
            self.sleep = False  # Else he wakes up.
        print("------------------\nBarber sleep {}\n------------------".format(self.sleep))

    def run(self):
        global SHOP_OPEN
        while SHOP_OPEN:
            while self.queue.empty():
                # Waits for the Event flag to be set, Can be seen as the Barber Actually sleeping.
                EVENT.wait()
                print("Barber is sleeping...")
            print("Barber is awake.")
            customer = self.queue
            self.is_empty()
            # FIFO Queue So first customer added is gotten.
            customer = customer.get()
            customer.trim()  # Customers Hair is being cut
            customer = self.queue
            customer.task_done()
            print(self.name)    # Which Barber served the Customer


def wait():
    time.sleep(1 * random.random())


if __name__ == '__main__':
    Earnings = 0
    SHOP_OPEN = True
    barbers = []
    all_customers = Queue(CUSTOMERS_SEATS)  # A queue of size Customer Seats

    for b in range(BARBERS):
        # TODO Pass the all_customers Queue to the Barber constructor
        b = Barber(all_customers)
        # Makes the Thread a super low priority thread allowing it to be terminated easier
        b.daemon = True
        b.start()   # Invokes the run method in the Barber Class
        # Adding the Barber Thread to an array for easy referencing later on.
        barbers.append(b)
    for c in range(10):  # Loop that creates infinite Customers
        print("----")
        # Simple Tracker too see the qsize (NOT RELIABLE!)
        print(all_customers.qsize())
        wait()
        c = Customer(all_customers)  # Passing Queue object to Customer class
        all_customers.put(c)    # Puts the Customer Thread in the Queue
        c.start()
    all_customers.join()    # Terminates all Customer Threads
    print("Barbers payment total:" + str(Earnings))
    SHOP_OPEN = False
    for i in barbers:
        i.join(timeout=4)    # Terminates all Barbers
        # Program hangs due to infinite loop in Barber Class, use ctrl-z to exit.

Explanation

  1. Global Variables & Constants:
    • CUSTOMERS_SEATS = 15: Maximum number of customers that can wait in the shop.
    • BARBERS = 3: Number of barbers working.
    • EVENT = threading.Event(): A synchronization primitive. Customers set it to signal their arrival (waking up a potentially sleeping barber). Barbers wait on this event.
    • Earnings = 0: Tracks the total money collected.
    • SHOP_OPEN = False (initially, set to True in main): Controls the main loop of the barber threads.
  2. Customer(threading.Thread) Class (Producer-like):
    • Represents a customer arriving at the barbershop.
    • __init__(self, queue):
      • Takes the shared all_customers queue.
      • self.rate: Randomly determines the price for the customer's haircut using what_customer().
    • what_customer() (static method):
      • Randomly assigns a customer type ("adult", "senior", "student", "child") and returns the corresponding price.
    • run(self):
      • This method is called when a customer thread starts.
      • Checks if the self.queue (waiting room) is not full.
        • If not full: It calls EVENT.set() (to signal availability) and then immediately EVENT.clear(). This is a very brief signal meant to wake up a barber.
      • If the queue is full: Prints that the customer has left.
      • Note: The customer object itself is added to the all_customers queue in the main part of the script, not within its own run method.
    • trim(self):
      • Simulates the haircut process: prints messages, sleeps for a random duration (3 * random.random()), and adds self.rate to global Earnings.
  3. Barber(threading.Thread) Class (Consumer-like):
    • Represents a barber who serves customers.
    • __init__(self, queue):
      • Takes the shared all_customers queue.
      • self.sleep = True: An attribute indicating if the barber is sleeping (though EVENT.wait() primarily controls this).
    • is_empty(self):
      • Checks if the customer queue is empty and updates self.sleep for printing status.
    • run(self):
      • The main loop for the barber, continues as long as SHOP_OPEN is True.
      • Waiting for Customer:
        • while self.queue.empty(): Enters a loop if no customers are in the queue.
          • EVENT.wait(): The barber "sleeps" here, waiting until another thread (a customer) calls EVENT.set().
      • Once woken (or if the queue wasn't empty):
        • customer = self.queue.get(): Retrieves a Customer object from the queue (FIFO). This call will block if the queue is empty (acting as a secondary wait if the event was set but another barber got the customer).
        • customer.trim(): Calls the customer's trim method to simulate the haircut.
        • self.queue.task_done(): Signals to the queue that the processed item is complete (used with queue.join()).
  4. wait() function:
    • A simple helper function to pause for a short, random duration.
  5. if __name__ == '__main__': (Main Execution Block):
    • Initializes Earnings to 0 and SHOP_OPEN to True.
    • Creates all_customers = Queue(CUSTOMERS_SEATS), a thread-safe queue to hold waiting customers.
    • Barber Creation:
      • Creates BARBERS number of Barber threads, passing all_customers to each.
      • Sets each barber thread as a daemon thread (so they exit when the main program exits).
      • Starts each barber thread (b.start()).
    • Customer Generation Loop:
      • A loop runs 10 times (not infinite as the comment suggests in that line, but the overall pattern can be made infinite).
      • c = Customer(all_customers): Creates a new customer.
      • all_customers.put(c): The main thread adds the customer object to the waiting queue.
      • c.start(): Starts the customer thread (its run method executes, which signals the EVENT).
    • all_customers.join(): The main thread waits until every customer put into the queue has had task_done() called for it (i.e., all 10 customers are served).
    • Prints final Earnings.
    • SHOP_OPEN = False: Signals barber threads to stop their while SHOP_OPEN: loops.
    • Attempts to join() barber threads with a timeout. The comment # Program hangs due to infinite loop in Barber Class, use ctrl-z to exit. indicates a potential issue where barbers might get stuck on EVENT.wait() if SHOP_OPEN becomes false while they are waiting and no more customer events are triggered to wake them.

Synchronization and Flow

  1. Barbers start and, if the all_customers queue is empty, they wait on EVENT.wait().
  2. The main loop creates a Customer object, adds this object to the all_customers queue, and then starts the Customer thread.
  3. The Customer thread's run() method executes. If there's space, it calls EVENT.set(), briefly waking up a barber waiting on EVENT.wait(). It then immediately calls EVENT.clear().
  4. The woken barber (or one that wasn't sleeping) attempts self.queue.get(). This retrieves the Customer object that the main thread put on the queue.
  5. The barber calls the trim() method on the retrieved Customer object.
  6. After the haircut, the barber calls self.queue.task_done().
  7. This continues until 10 customers are processed. Then SHOP_OPEN is set to False, and the program attempts to shut down.

Data and Code Synchronisation

Concurrent Programming Patterns

Divide and Conquer

Map and Reduce

Pipelines/Workflows

Recursion

Repository Pattern

Flynn's Taxonomy

Definition

Flynn's Taxonomy is a classification system for parallel computer architectures, proposed by Michael J. Flynn in 1966 and extended in 1972. It categorizes architectures based on the number of concurrent instruction streams and data streams they can process.
202505112145-flynns_taxonomy.png

Categories

SISD (Single Instruction, Single Data Stream)

SIMD (Single Instruction, Multiple Data Streams)

MISD (Multiple Instruction, Single Data Stream)

MIMD (Multiple Instruction, Multiple Data Streams)

Significance

Python Parallel Programming

Options for Concurrency

threading

multiprocessing

asyncio

Quick Comparison Table

Feature threading multiprocessing asyncio
Mechanism Threads (shared memory) Processes (separate memory) Coroutines (single thread, event loop)
Parallelism Concurrent (GIL limits CPU-bound) True Parallel (multi-core) Concurrent (single-core)
Use Case I/O-bound CPU-bound High-volume I/O-bound
GIL Limits true parallelism Bypassed (per process) Operates within it (single thread)
Data Sharing Easy (direct, needs sync) Harder (IPC) Easy (within its process)
Overhead Moderate High Low (for I/O)