Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

🏠 Back to Blog

Notes based on Operating Systems: Three Easy Pieces by Remzi & Andrea Arpaci-Dusseau

Concurrency: Multi-Tasking & Synchronization

This section explores another important aspect of computer systems: concurrency. Concurrency generally means “multiple things going on at the same time”. We could also call it multi-tasking. With the support for virtualization, we have already seen how the OS runs multiple processes on a machine with shared CPU cores and memory space, which is one form of concurrency. However, processes generally represent separate entities and the OS tries its best to ensure isolation between processes. Are there alternatives?

Abstraction of Thread

We introduce a new abstraction of a running entity called a thread. For a single process, it could have multiple threads sharing the same address space (hence, the same code, data, etc.). Different threads differentiate with each other mainly in:

Different PC register values: each thread has its own PC, meaning they can (and probably will) be executing different pieces of code of the program Separate stacks and SP: each thread has its own function call stack space (because they will be executing different routines) and the corresponding stack pointer registers, etc.; the function stacks in this case may be called thread-local We call the action of utilizing multiple threads as multi-threading. A program spawning threads is a multi-threaded program. A default process with just the main thread can be viewed as a single-threaded process. Multi-threading brings at least two benefits:

Parallelism if we have multiple CPU cores and/or certain excessive resources: multiple threads can run on different cores at the same time; this action is called parallelization Multiprogramming: enables overlapping of I/O (or other blocking behaviors) with other activities of the program: when one thread issues an I/O on behalf of the program, other threads can fill this gap and do useful computation on the core

Multi-Threaded Address Space

A multi-threaded process’s address space looks like:

The OS and the threading library could have different ways to layout the thread-local stacks.

Thread Control Block (TCB)

Since threads have their own PC and other registers, switching between threads on a CPU core is also a context switch. The difference from process context switches is that thread context switches are more lightweight: just registers, no page table pointer switches, no TLB flushes, no open file state changing, etc.

Apart from the registers context, each thread could have different scheduling properties, e.g., priorities, time-slice lengths, etc. A thread control block (TCB) is the data structure for keeping these states of a thread. TCBs of a process are typically part of the PCB of the process. For example, we could modify our xv6 PCB structure to something like:

struct context {

int eip;
int esp;
int ebx;
int ecx;
int edx;
int esi;
int edi;
int ebp;

}; ​ /** The TCB structure. */ struct thread {

int tid;                    // Thread ID
struct context context;     // Register values context
struct sched_properties sp; // Scheduling properties

}; ​ /** The PCB structure. */ struct proc {

...
struct thread threads[MAX_THREADS];     // TCBs

};

Thread APIs

The POSIX specification defines a set of pthread APIs:

pthread_create() - create a new thread with attributes and start its execution at the entry of a routine (often called thread function); the creator is the parent thread and the created one is a child thread; thread creations could form a deep tree like in process forks pthread_join() - wait for a specific thread to exit (called joining the thread) pthread_exit() - for a thread to exit itself (returning from thread function implicitly means exit) Synchronization primitives: pthread_mutex_(), pthread_cond_(), … These interfaces are either raw system calls or wrappers provided by a pthread library (should be explicitly linked against) if the OS provides slightly different threading syscalls in the low level.

The portable operating system interface (POSIX) is a set of specifications that define the interfaces an operating system should provide to applications. There can be many different OS designs and implementations, but as long as they comply to a (sub)set of POSIX interfaces, applications can be developed on one POSIX system platform and ported to run on another POSIX platform with ease. Please see the linked wikipedia page for what interfaces the POSIX specification includes and which operating systems are POSIX-compliant.

Synchronization

Threads are powerful because they share the same address space , but such sharing leads to a critical problem: how to prevent race conditions if they operate on shared data and how to express a “waiting-on-something” logic? In general, the term synchronization refers to mechanisms + policies to ensure

Atomicity of accesses Ordering of execution It applies to any case where there are multiple running entities sharing something, but its exceptionally critical to OSes because the problem of sharing happens so often.

Race Conditions

If we do not apply any restrictions on the scheduling of threads, the result of execution of a multi-threaded program is likely to be indeterministic. Each thread is executing its own stream of instructions, yet every time which one executes its next instruction could be arbitrary. Consider an example where two threads executing the same “plus-one” routine to a shared variable with original value 50:

; increments the variable at addr by 1 mov addr, %eax add $0x1, %eax mov %eax, addr You might expect a correct program to always produce 52 as the result. However, the following sequence could happen and the result could be 51 which is incorrect:

This is what we call a race condition (or, more specifically, a data race): unprotected timing of execution. The result of such a program would be indeterminate: sometimes it produces the correct result but sometimes the results are different and are likely wrong. We call such piece of code a critical section: code routine that accesses a shared resource/data and, if unprotected, might yield incorrect result.

Atomicity & Mutex Locks

What we want is that whenever a thread starts executing a critical section, it executes until the completion of the critical section without anyone else interrupting in the middle to execute any conflicting critical section (executing unrelated code would be fine). In other words, we want the entire critical section to be atomic (we assume every single instruction is already guaranteed atomic).

At a higher level, one big, atomic action that groups a sequence of small actions is logically called a transaction (a name widely used in database systems, though the problem setup is quite different).

Atomicity is guaranteed if we have mutual exclusion (mutex) among concurrent executions of critical sections. We introduce a powerful synchronization primitive called a lock to enforce mutual exclusion. A lock is a data structure that supports at least the following two operations:

acquire - when entering a critical section, grab the lock and mark it as acquired (or locked, held); if there are concurrent attempts on acquiring a lock, must ensure that only one thread succeeds and proceeds to execute the critical section; others must somehow wait until its release to compete again release - when leaving the critical section, release the lock so it turns available (or unlocked, free) pthread_mutex_t balance_lock = PTHREAD_MUTEX_INITIALIZER; ​ void deposit_one() {

pthread_mutex_lock(&balance_lock);
balance++;        // Critical section
pthread_mutex_unlock(&balance_lock);

} ​ void withdraw_one() {

pthread_mutex_lock(&balance_lock);
if (balance > 0)  // Critical section with the same interest on
    balance--;    // the shared balance variable, so protected
                  // with the same lock.
pthread_mutex_unlock(&balance_lock);

} We will talk about how does the OS implement locks in a dedicated section below.

Ordering & Condition Variables

Apart from atomicity of critical sections, we also desire the ability to enforce certain ordering of executions so that a routine cannot proceed until some condition is made true by some other thread. Though constantly spinning on a variable might work, it is very inefficient. We introduce a second type of synchronization primitive called a condition variable. A condition variable is a data structure that supports at least the following two operations:

wait - block and wait until some condition becomes true; waiter typically gets added to a queue signal - notify, wake up one waiter on a condition (typically head of wait queue) (optional) broadcast - wake up all waiters on the condition variable (assuming Mesa semantic, see below); a condition variable using broadcast is called a covering condition A condition variable is often coupled with a mutex lock to protect the shared resource to which the condition is related:

volatile int done = 0; pthread_cond_t done_cond = PTHREAD_COND_INITIALIZER; pthread_cond_t done_lock = PTHREAD_MUTEX_INITIALIZER; ​ void thread_exit() {

pthread_mutex_lock(&done_lock);
done = 1;
pthread_cond_signal(&done_cond);
pthread_mutex_unlock(&done_lock);

} ​ void thread_join() {

pthread_mutex_lock(&done_lock);
// Use a while loop in case something reverted the condition
// right after I wake up.
while (done == 0) {
    // When entering wait, done_lock will be released
    pthread_cond_wait(&done_cond, &done_lock);
    // Upon waking up, done_lock will be re-acquired
}
pthread_mutex_unlock(&done_lock);

} Another great example of demonstrating the usage of condition variables is the famous bounded buffer producer/consumer problem introduced by Dijkstra, where we have a fixed-size buffer array, producer(s) trying to write to empty slots of the buffer, and consumer(s) trying to grab things from the buffer. Use cases of bounded buffer include a web server request dispatching queue, or piping the output of one command to the input of another. A correct implementation looks like:

volatile int buffer[NUM_SLOTS]; volatile int fill_idx = 0, grab_idx = 0; volatile int count = 0; ​ pthread_cond_t buf_not_empty = PTHREAD_COND_INITIALIZER; pthread_cond_t buf_not_full = PTHREAD_COND_INITIALIZER; pthread_mtuex_t buf_lock = PTHREAD_MUTEX_INITIALIZER; ​ void producer() {

while (1) {
    pthread_mutex_lock(&buf_lock);
    while (count == NUM_SLOTS)
        pthread_cond_wait(&buf_not_full, &buf_lock);
    buffer[fill_idx] = something;
    fill_idx = (fill_idx + 1) % NUM_SLOTS;
    count++;
    pthread_cond_signal(&buf_not_empty);
    pthread_mutex_unlock(&buf_lock);
}

} ​ void consumer() {

while (1) {
    pthread_mutex_lock(&buf_lock);
    while (count == 0)
        pthread_cond_wait(&buf_not_empty, &buf_lock);
    int grabbed = buffer[grab_idx];
    grab_idx = (grab_idx + 1) % NUM_SLOTS;
    count--;
    pthread_cond_signal(&buf_not_full);
    pthread_mutex_unlock(&buf_lock)
}

} Using a while loop on the condition is something related to the semantic of the conditional variable.

Mesa semantic: first explored in the Mesa system and later deployed by virtually every system, this semantic says a condition variable signal is just a hint that the state of the condition has changed, but it is not guaranteed that the state remains true when the woken up thread gets scheduled and runs Hoare semantic: described by Hoare, this stronger semantic requires a condition variable signal to immediately wake up and schedule the woken up thread to run, which is less practical in building actual systems

Semaphores

The third type of synchronization primitive we will investigate is a semaphore. Dijkstra and colleagues invented semaphores as a more general form of synchronization primitive which can be used to implement the semantics of both locks and condition variables, unifying the two conceptually. A semaphore is an integer value that we can manipulate with two routines:

wait (or down, or V) - decrement the value of semaphore by 1, and then wait if the value is now negative post (or up, or P) - increment the value of semaphore by 1, and if there are one or more threads waiting, wake up one We can use semaphores to implement mutex locks and condition variables. A binary semaphore is equivalent to a mutex lock:

sem_t sem; sem_init(&sem, 1); // Initialize to 1 ​ sem_wait(&sem); … // Critical section sem_post(&sem); And semaphores could be used to enforce ordering as well, like what condition variables do:

sem_t sem; sem_init(&sem, 0); // Initialize to 0 ​ void thread_exit() {

sem_post(&sem);

} ​ void thread_join() {

sem_wait(&sem);

} sem_t vacant, filled; sem_init(&vacant, NUM_SLOTS); // all vacant sem_init(&filled, 0); // none filled ​ sem_t mutex; sem_init(&mutex, 1); // Lock initialized to 1 ​ void producer() {

while (1) {
    sem_wait(&vacant);
    sem_wait(&mutex);
    put_something();
    sem_post(&mutex);
    sem_post(&filled);
}

} ​ void consumer() {

while (1) {
    sem_wait(&filled);
    sem_wait(&mutex);
    grab_something();
    sem_post(&mutex);
    sem_post(&vacant);
}

} Generally, it is the initial value of a semaphore that determines how it behaves. To set the initial value of a semaphore, Kivolowitz has a general rule that you set it to the number of resources you are able to give away immediately after initialization. For example, mutex lock 1, waiting for child thread done 0, buffer vacant slots size of buffer, etc. However, generality is not always good and easy - given that we already have good primitives like locks and condition variables, it is mostly not necessary to use semaphores in practice.

Implementing Locks

We introduced the lock interface but haven’t talked about how are locks implemented. We would like a lock implementation that provides the following four properties:

Correctness: mutual exclusion must be enforced Fairness: each thread contending for a lock should get a fair shot at acquiring it, possibly considering their incoming order; there should be no starvation: a lock waiter thread constantly losing the competition, thus never obtaining it Performance: the acquire and release operations themselves should not incur too much overhead, and the lock should be against contention: multiple threads trying to acquire but none succeeds in bounded-time Scalability: with more CPU cores and/or an increasing number of lock competitors, performance should not drop badly Certainly, a simple spin-wait on a shared variable is neither correct nor performant (this is not saying spinning locks are bad, as we will discuss in a section below):

void acquire(mutex_t *lock) {

while (mutex->flag == 1) {}  // spin-wait
mutex->flag = 1;

} ​ void release(mutex_t *lock) {

mutex->flag = 0;

} This implementation does not guarantee mutual exclusion, for example:

Controlling Interrupts

Assume we only have a single-processor machine and we only want one lock variable, a naive approach would be to just disable interrupts for critical sections:

void acquire() {

DisableInterrupts();

} ​ void release() {

EnableInterrupts();

} Though simple, this approach has significant drawbacks:

Controlling interrupts are privileged operations and we do not want any thread to abuse them Does not work on multiprocessors as it shuts down the whole interrupt mechanism if anyone enters critical section Cannot do fine-grained locking: there is effectively only one lock variable because disabling interrupts is a global effect, so we cannot assign different locks to unrelated critical sections (in other words, it is very coarse-grained locking) Disabling interrupts could lead to bad things because the CPU misses everything from any external device

Hardware Atomic Instructions

Building locks out of pure load/store instructions are possible (e.g., see Peterson’s algorithm and Lamport’s bakery algorithm), yet with a little bit of extra hardware support, things can get much easier and much more efficient. Here we demand the hardware to provide certain atomic instructions: do something more than just a single load/store in one instruction, guaranteed to be atomic by the hardware. Classic examples include:

Test-and-Set (TAS): write a 1 to a memory location and return the old value on this location “simultaneously”

TEST_AND_SET(addr) -> old_val // old_val = *addr; // *addr = 1; // return old_val; This slightly more powerful TAS instruction enables this lock implementation:

void acquire() {

while (TEST_AND_SET(&flag) == 1) {}

} ​ void release() {

flag = 0;

} Compare-and-Swap (CAS): compare the value on a memory location with a given value, and if they are the same, write a new value into it, again “simultaneously”

COMPARE_AND_SWAP(addr, val, new_val) -> old_val // old_val = *addr; // if (old_val == val) // *addr = new_val; // return old_val; Building a lock out of CAS:

void acquire() {

while (COMPARE_AND_SWAP(&flag, 0, 1) == 1) {}

} ​ void release() {

flag = 0;

} Load-Linked (LL) & Store-Conditional (SC): a pair of instructions used together; LL is just like a normal load; SC tries to store a value to the location and succeeds only if there’s no LL going on at the same time, otherwise it returns failure

LOAD_LINKED(addr) -> val // return *addr; STORE_CONDITIONAL(addr, val) -> success? // if (no LL to addr happening) { // *addr = val; // return 1; // success // } else // return 0; // failed Building a lock out of LL/SC:

void acquire() {

while (1) {
    while (LOAD_LINKED(&flag) == 1) {}
    if (STORE_CONDITIONAL(&flag, 1) == 1)
        return;
}

} ​ void release() {

flag = 0;

} Fetch-and-Add (FAA): increment a value while returning the old value at the location

FETCH_AND_ADD(addr) -> old_val // old_val = *addr; // *addr += 1; // return old_val; FAA enables us to build a ticket lock which takes fairness into consideration and ensures progress for all threads, thus preventing starvation:

volatile int ticket = 0; volatile int turn = 0; ​ void acquire() {

int myturn = FETCH_AND_ADD(&ticket);
while (turn != myturn) {}

} ​ void release() {

turn++;

} In modern multicore systems, the scalability of locks becomes critical. There are many more advanced lock design & implementations trying to avoid cache contention and trying to have NUMA-awareness. Please see this post if interested.

Spinning vs. Blocking

A spinning lock (or spinlock, non-blocking lock) is a lock implementation where lock waiters will spinning in a loop checking for some condition. The examples given above are basic spinlocks. Spinlocks are typically used for low-level critical sections that are short, small, but invoked very frequently, e.g., in device drivers.

Advantage: low latency for lock acquirement as there is no scheduling stuff kicking in – value changes reflect almost immediately

Disadvantage:

Spinning occupies the whole CPU core and wastes CPU power if the waiting time is long that could have been used for scheduling another free thread in to do useful work Spinning also introduces the cache invalidation traffic throttling problem if not handled properly Spinning locks make sense only if the scheduler is preemptive, otherwise there is no way to break out of an infinite loop spin Spinning also worsens certain cases of priority inversion, if we are using a priority-based scheduler A blocking lock is a lock implementation where a lock waiter yields the core to the scheduler when the lock is currently taken. A lock waiter thread adds itself to the lock’s wait queue and blocks the execution of itself (called parking, or yielding, or descheduling) to let some other free thread run on the core, until it gets woken up (typically by the previous lock holder) and scheduled back. It is designed for higher-level critical sections. The pros and cons are exactly the opposite of a spinlock.

Advantage: not occupying full core during the waiting period, good for long critical sections Disadvantage: switching back and forth from/to the scheduler and doing scheduling stuff takes significant time, so if the critical sections are fast and invoked frequently, better just do spinning For example, our previous TAS-based lock could be modified to do blocking and queueing instead:

volatile queue_t *queue; volatile int flag = 0; volatile int guard = 0; // essentially an internal “lock”

                     // protecting the flag variable

​ void acquire() {

while (TEST_AND_SET(&guard, 1) == 1) {}
if (flag == 0) {
    flag = 1;
    guard = 0;
} else {
    queue_add(queue, self_tid);
    guard = 0;
    park();
}

} ​ void release() {

while (TEST_AND_SET(&guard, 1) == 1) {}
if (queue_empty(queue))
    flag = 0;
else
    unpark(queue_pop(queue));
guard = 0;

} The above code has one subtle problem called a wakeup/waiting race: what if a context switch happens right before the acquirer’s park(), and the switched thread happens to release the lock? The releaser will try to unpack the not-yet-parked acquirer, so the acquirer could end up parking forever. One solution is to add a setpark() call before releasing guard, and let the park() call wake up immediately if any unpark() is made after setpark().

} else {

queue_add(queue, self_tid);
setpark();
guard = 0;
park();

} There are also hybrid approaches mixing spinlocks with blocking locks, by first spinning for a while in case the lock is about to be released, otherwise go to park. It is referred to as a two-phase lock. For example, the Linux locks based on its futex syscall support is one such approach.

Advanced Concurrency

We have yet more to discuss around concurrency and synchronization.

Lock-Optimized Data Structures

Adding locks to protect a data structure makes it thread-safe: safe to be used by concurrent running entities (not only for multi-threaded applications, but also for OS kernel data structures). Simply acquiring & releasing a coarse-grained lock around every operation on an object could lead to performance and scalability issues.

In general, the more fine-grained locks (smaller, fewer logics in critical sections), the better scalability. Here we list some examples on optimized data structures that reduce lock contention:

Approximate counters: on multicore machine, for a logical counter variable, have one local counter per core as well as one global counter; local counters are each protected by a core-local lock, and the global counter is protected by a global lock; a counter update issued by a thread just grabs the local lock on the core and applies to the local counter; whenever a local counter reaches a value/time threshold, the core grabs the global lock and updates the global counter. Move memory allocation & deallocation out of locks, and only lock necessary critical sections; take extra care about failure paths. Hand-over-hand locking (lock coupling) on linked lists: instead of locking the whole list on every operation, have one lock per node; when traversing, first grabs the next node’s lock and then releases the current nodes lock; however, acquiring & releasing locks at every step of a traversal probably leads to poor performance in practice. Dual-lock queues: use a head lock for enqueue and a tail lock for dequeue on a queue with a dummy node Bucket-lock hash tables: use one lock per bucket on a (closed-addressing) hash table Reader-writer lock: a lock that allows, at any time, either one writer and nobody else holding it, or multiple readers but no writer holding it; see Chapter 31.5 of the book for more Many more concurrent data structures designs exist since this has been a long line of research. There are also a category of data structures named lock-free (or non-blocking) data structures which don’t use locks at all with the help of clever memory operations & ordering. See read-copy-update (RCU) for an example.

Concurrency Bugs

The biggest source of bugs in concurrency code is the deadlock problem: all entities end up waiting for one another and hence making no progress - the system halts. Examples of deadlocks are everywhere:

Deadlocks stem from having resource dependency cycles in the dependency graph. The simplest form of such cycle:

thread T1: thread T2:

acquire(lock_A);        acquire(lock_B);
acquire(lock_B);        acquire(lock_A);
...                     ...

T1’s code order results in lock_B depending on lock_A while T2’s code order results in the opposite, forming a cycle. If T2 gets scheduled and acquires lock_B right after T1 acquires lock_A, they are stuck. In real-world code, the code logic will be much more complicated and encapsulated, so deadlock problems are way harder to identify.

The famous dinning philosophers problem is a great thought experiment for demonstrating deadlocks: philosophers sit around a table with forks between them, each philosopher spends time to think or eat; when one attempts to eat, must grab the fork on their left and on their right.

If every philosopher always grabs the fork on the left before grabbing the fork on the right, there could be deadlocks: each philosopher grabs the fork on their left at roughly the same time and none is able to grab the fork on their right. A classic solution is to break the dependency: let one specific philosopher always grab the fork on the right before grabbing the fork on the left.

There has also been a long line of research on deadlocks. Theoretically, four conditions need to hold for a deadlock to occur:

Mutual exclusion: threads claim exclusive control of resources that they have acquired Hold-and-wait: threads hold resources allocated to them while waiting for additional resources No preemption: resources cannot be forcibly removed from holders Circular wait: there exists a circular chain of threads s.t. each holds some resources being requested by the next thread in chain Accordingly, to tackle the deadlocks problem, we could take the following approaches, trying to break some of the conditions (see Chapter 32 of the book for detailed examples):

Prevention (in advance)

Avoid circular wait: sort the locks in a partial order, and whenever the code wants to acquire locks, adhere to this order Avoid hold-and-wait: acquire all locks at once, atomically, through e.g. always acquiring a meta-lock Allow preemption: use a trylock() interface on lock_B and revert lock_A if B is not available at the time Avoid mutual exclusion: use lock-free data structures with the help of powerful hardware atomic instructions Avoidance (at run-time) via scheduling: don’t schedule threads that could deadlock with each other at the same time

Recovery (after detection): allow deadlocks to occur but reserve the ability to revert system state (rolling back or just rebooting)

Apart from deadlocks, most of the other bugs are due to forgetting to apply synchronization, according to Lu et al.’s study:

Atomicity violation: a code region is intended to be atomic, but the atomicity is not enforced (e.g., by locking) Order violation: logic A should always be executed before B, but the order is not enforced (e.g., by condition variables)

Event-Based Model

Many modern applications, in particular GUI-based applications and web servers, explore a different type of concurrency based on events instead of threads. To be more specific:

Thread-based (work dispatching) model: what we have discussed so far are based on programs with multiple threads (say with a dynamic thread pool); multi-threading allows us to exploit parallelism (as well as multiprogramming)

Data parallelism: multiple threads doing the same processing, each on a different chunk of input data Task parallelism: multiple threads each handling a different sub-task; could be in the form of pipelining This model is intuitive, but managing locks could be error-prone. Dispatching work could have overheads. It also gives programmers little control on scheduling.

Event-based concurrency model: using only one thread, have it running in an event loop; concurrent events happen in the form of incoming requests, and the thread runs corresponding handler logic for requests one by one

while (1) {

events = getEvents();
for (e in events)
    handleEvent(e);

} This model sacrifices parallelism but provides a cleaner way to do programming (in some cases) and yields shorter latency (again, in some cases).

See the select() and epoll() syscalls for examples. Modern applications typically use multi-threading with event-based polling on each thread in some form to maximize performance and programmability.

There are some other popular terminologies used in concurrent programming that I’d like to clarify:

Polling vs. Interrupts

Polling: actively check the status of something; useful for reducing latency if the thing is going to happen very fast (common for modern ultra-fast external devices), but wastes CPU resource on doing so Interrupts (general): passively wait for signals; has overhead on context switching to/from the handler Synchronous vs. Asynchronous

Synchronous (blocking, sync): interfaces that do all of the work upon invocation, before returning to the caller Asynchronous (non-blocking, async): interfaces that kick off some work but return to the caller immediately, leaving the work in the background and the status to be checked by the caller later on; asynchronous interfaces are crucial to event-based programming since a synchronous (blocking) interface will halt all progress Inter-process communication (IPC): so far we talked about threads sharing the same address space, but there must also be mechanisms for multiple isolated processes to communicate with each other. IPC could happen in two forms

Shared-memory: use memory map to share some memory pages (or through files) across processes, and do communication (& correct synchronization) beyond Message-passing: sending messages to each other, see UNIX signals, domain sockets, pipes, etc.

Multi-CPU Scheduling

With the knowledge about the memory hierarchy and concurrency, let’s now go back to scheduling and add back a missing piece of the “virtualizing the CPU” section: scheduling on multiprocessor machines (multiple CPUs, or multiple cores packed in one CPU). Multi-CPU scheduling is challenging due to these reasons:

Cache coherence traffic can cause thrashing, as we have mentioned Threads are unlikely independent so we need to take care of synchronization Cache affinity: a process, when run on a particular CPU, builds up state in that CPU’s cache and thus has affinity to be scheduled on that CPU later NUMA-awareness: similar to cache affinity, on NUMA architectures it is way cheaper to access socket-local memory than memory on a remote socket Here are two general multiprocessor scheduling policies:

Single-Queue Multiprocessor Scheduling (SQMS)

It is simple to adapt a single-CPU scheduling policy to SQMS, but the downside is that SQMS typically requires some form of locking over the queue, limiting scalability. Also, it is common to enhance SQMS with cache affinity- and NUMA-awareness, like shown in the example figure. It migrates process E across CPUs while preserving others on their core. Fairness could be ensured by choosing a different one to migrate for the next window.

Multi-Queue Multiprocessor Scheduling (MQMS)

MQMS has one local queue per CPU, and uses single-CPU scheduling policies for each queue. MQMS is inherently more scalable, but has a problem of load imbalance: if one CPU has a shorter queue than another, processes assigned to that CPU effectively get more time slices scheduled. A general load balancing technique is do to work stealing across cores based on some fairness measure.

The Linux multiprocessor scheduler adopts MQMS (O(1), CFS) and SQMS (BFS) at different times.