~/techniques/concurrency

Concurrency

Threads that share data can lose updates. Locks, condition variables and a bounded queue make them take turns safely.

what

Guard shared state with a lock, sleep on a condition variable inside a while loop until the state you need is true, and pass work between threads through a bounded queue.

use when

Many slow I/O calls to overlap, CPU work to spread over cores, or a design question about a thread-safe counter, cache, rate limiter or job queue.

time

O(1)

space

O(n)

You’ll recognise it when

  • Several threads read and write the same data: a counter, a cache, a bank balance, a rate limiter.
  • One part of a program makes work and another part uses it: a crawler that finds links and workers that fetch them, a pool of workers taking jobs.
  • A thread has to wait until something is true: the queue has an item, there’s a free slot, every worker has finished.
  • The question says “thread-safe”, “concurrent” or “blocking”, or asks you to make many slow network calls finish sooner.
  • Something works on your laptop and fails once in a thousand runs under load.

Concurrency is easy to mix up with parallelism. Concurrent tasks are all in progress at once and take turns; parallel tasks run at the same instant on different cores. Threads always give you the first, and the second only in some languages (see the GIL below).

The idea

Two cashiers share one notebook with the day’s total. To record a sale, a cashier reads the total, adds the sale in their head, and writes the new total down. If both read 100 before either writes, each writes 100 plus their own sale, and one sale disappears. The fix is a single pen: only the cashier holding the pen may read or write the total. That pen is a lock.

Every concurrency bug on this page comes from two things together: shared data that changes, and switches you don’t control. Between any two steps of one thread, another thread can run. So you make each read-modify-write of shared data one indivisible step, with a lock. When a thread must wait for something, it sleeps on a condition and re-checks when it wakes, instead of spinning. Better still, you share less: threads hand data to each other through a queue, so the queue is the only thing that needs a lock.

Threads, processes or asyncio

There are three common ways to have several things in progress at once. They differ in what they share and in what they speed up.

Shares memory? Best for
Threads yes, all of it waiting on the network or disk
Processes no, data is copied CPU-heavy pure Python
asyncio yes, in one thread thousands of network calls
  • Threads are several flows of control inside one process. They see the same memory, which is handy and is also where races come from. In C++ and Java they run on all your cores, so they speed up CPU work too. In standard Python only one runs Python code at a time (see the GIL below).
  • Processes are separate programs, each with its own memory, so they can’t race on shared variables. Anything you send between them gets copied, and they’re slower to start.
  • asyncio runs everything on one thread. An event loop switches between coroutines at each await, so thousands of them can wait on the network at once. One blocking call, like time.sleep or a long loop, stalls all of them.

A rule of thumb for Python: a few dozen downloads at once, use a thread pool (ThreadPoolExecutor). Thousands of open connections, use asyncio. Number crunching in pure Python, use processes (ProcessPoolExecutor), or NumPy, which does its heavy loops outside the interpreter.

import asyncio
async def fetch(i):
await asyncio.sleep(1) # a stand-in for a network call
return i * i
async def main():
results = await asyncio.gather(*(fetch(i) for i in range(100)))
print(sum(results)) # about 1 second in total, not 100
asyncio.run(main())

asyncio has fewer races, because nothing else runs between two awaits. A read-modify-write that spans an await can still race, though, and asyncio has its own asyncio.Lock, asyncio.Queue and asyncio.Semaphore for that.

The GIL in plain words

The standard Python interpreter, CPython, has a Global Interpreter Lock. A thread must hold it to run Python code, so only one thread runs Python code at any moment. The threads take turns: the interpreter asks the running thread to let go every 5 milliseconds or so, and a thread always lets go while it waits for the network, the disk or a sleep.

  • Threads don’t make pure-Python number crunching faster: only one of them computes at a time. Use processes for that.
  • Threads do help with I/O. While one waits for a reply, the others run. Libraries like NumPy, hashlib and zlib also let go of the GIL during long computations.
  • The GIL does not make your code thread-safe. It protects the interpreter’s own bookkeeping, one bytecode at a time. count += 1 is several bytecodes (load, add, store), and a switch can land between them.

Python 3.13 added an optional build without the GIL, where threads really run in parallel. Code that leaned on the GIL for safety breaks there first. Java and C++ have no GIL at all, so their threads run on every core and races show up much more often.

How it works

Two threads, A and B, each run count += 1 once, starting from count = 0. Each increment is three steps: read count into a private tmp, add 1, and write tmp back. Scroll through the steps and the graphic follows along. The switch at its top left re-runs the same schedule with or without the lock, and edit lets you choose which thread runs when.

  1. Start: count = 0, and each thread has its own tmp. The CPU runs one step at a time and may switch threads after any step.
  2. A reads count into its tmp: 0.
  3. The scheduler switches to B before A has written anything back. B reads count too, and gets the same 0.
  4. Each thread adds 1 to its own tmp. Both now hold 1, and count is still 0.
  5. A writes its 1 to count. So far, so good.
  6. B writes its 1 on top. It doesn’t add to what’s there: it replaces it with a value worked out from an old read.
  7. Two increments ran, and count = 1. That’s a race condition: the result depends on timing, and here one update is lost.
  8. The same schedule, now with a lock. A takes it before reading.
  9. B asks for the lock while A holds it, so B sleeps. It can’t read count until A is completely done.
  10. A writes 1 and releases the lock. B wakes up, takes the lock, and reads the new value, 1.
  11. count = 2. The read, the add and the write are now one step that no other thread can cut into: a critical section.
loading concurrency…

Why it works: only one thread can hold lock at a time, and every access to count happens while holding it. So the steps of two increments can’t interleave, and each increment sees the result of the one before. A lock only protects data if every access goes through it, reads included.

import threading
class Counter:
def __init__(self):
self.count = 0
self.lock = threading.Lock() # shared by every thread
def increment(self):
with self.lock:
self.count += 1
# leaving the block releases the lock, even on an error
def run(threads, per_thread):
counter = Counter()
def work():
for _ in range(per_thread):
counter.increment()
workers = [threading.Thread(target=work)
for _ in range(threads)]
for t in workers:
t.start()
for t in workers:
t.join() # wait for each to finish
return counter.count
#include <mutex>
#include <thread>
#include <vector>
class Counter {
public:
long long count = 0;
void increment() {
std::lock_guard<std::mutex> guard(lock);
count += 1;
}
private:
std::mutex lock; // shared by every thread
};
long long run(int threads, int perThread) {
Counter counter;
std::vector<std::thread> workers;
for (int t = 0; t < threads; t++) {
workers.emplace_back([&counter, perThread] {
for (int i = 0; i < perThread; i++) counter.increment();
});
}
for (auto& w : workers) w.join(); // wait for each to finish
return counter.count;
}
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.locks.ReentrantLock;
class Counter {
long count = 0;
// one lock, shared by every thread
private final ReentrantLock lock = new ReentrantLock();
void increment() {
lock.lock();
try {
count += 1;
} finally {
lock.unlock();
}
}
static long run(int threads, int perThread)
throws InterruptedException {
Counter counter = new Counter();
List<Thread> workers = new ArrayList<>();
for (int t = 0; t < threads; t++) {
Thread w = new Thread(() -> {
for (int i = 0; i < perThread; i++) counter.increment();
});
workers.add(w);
w.start();
}
for (Thread w : workers) w.join(); // wait for each to finish
return counter.count;
}
}

with self.lock: releases the lock however the block ends, even when it raises an exception. C++’s std::lock_guard does the same when it goes out of scope, and in Java the unlock() goes in a finally. A lock that is never released freezes every thread that asks for it later.

Waiting for a condition

Sometimes a thread has to wait until something is true: an item is in the queue, or there is room for one. Spinning in a loop (while not items: pass) burns a whole core, and if the spinning thread holds the lock, the thread that would add an item can never get in. A condition variable lets a thread sleep until another thread says the state has changed.

The pattern is always the same:

with cond: # cond = threading.Condition(lock)
while not ready(): # while, not if
cond.wait() # unlock, sleep, relock before returning
use_the_state()
# on the other side
with cond:
change_the_state()
cond.notify() # wake one waiter (notify_all: every one)
  • Call wait() with the lock held. It releases the lock while the thread sleeps, otherwise nobody could change the state, and takes it back before returning.
  • Re-check in a while loop. Between the notify and the moment the woken thread gets the lock back, another thread can grab the item. Java and POSIX threads also allow spurious wakeups, where wait() returns with no notify at all.
  • Change the state and notify while holding the same lock. Then a waiter can’t miss the change in the gap between its check and its wait().

Producer and consumer

This pattern ties everything together. Producers put work into a queue and consumers take it out, and the queue is the only thing they share. It has a fixed capacity, so a fast producer can’t fill all the memory: when the queue is full, put waits. When it’s empty, take waits. One lock guards items, and two condition variables share that lock: producers wait on not_full, and consumers wait on not_empty.

In the graphic, each turn lets one thread run one call. Watch C1 in the default run: P1’s notify wakes it, but C2 takes the item first, so C1’s while loop sends it back to sleep. Later P1 fills the queue and has to wait for space.

loading bounded-queue…
import threading
from collections import deque
class BoundedQueue:
def __init__(self, capacity):
self.items = deque()
self.capacity = capacity
self.lock = threading.Lock()
# two conditions, one lock: both guard the same items
self.not_full = threading.Condition(self.lock)
self.not_empty = threading.Condition(self.lock)
def put(self, item):
with self.lock:
while len(self.items) == self.capacity:
self.not_full.wait()
self.items.append(item)
self.not_empty.notify()
def take(self):
with self.lock:
while not self.items:
self.not_empty.wait()
item = self.items.popleft()
self.not_full.notify()
return item
STOP = None # a "poison pill": tells one consumer to quit
def run(capacity, producers, consumers, per_producer):
q = BoundedQueue(capacity)
taken = [[] for _ in range(consumers)]
def produce(p):
for i in range(per_producer):
q.put(p * per_producer + i)
def consume(c):
while (item := q.take()) is not STOP:
taken[c].append(item)
makers = [threading.Thread(target=produce, args=(p,))
for p in range(producers)]
takers = [threading.Thread(target=consume, args=(c,))
for c in range(consumers)]
for t in makers + takers:
t.start()
for t in makers:
t.join()
for _ in takers:
q.put(STOP) # one pill per consumer, after all real items
for t in takers:
t.join()
return taken
#include <condition_variable>
#include <deque>
#include <mutex>
#include <thread>
#include <vector>
class BoundedQueue {
public:
explicit BoundedQueue(size_t capacity) : capacity(capacity) {}
void put(int item) {
std::unique_lock<std::mutex> guard(lock);
while (items.size() == capacity) {
notFull.wait(guard);
}
items.push_back(item);
notEmpty.notify_one();
}
int take() {
std::unique_lock<std::mutex> guard(lock);
while (items.empty()) {
notEmpty.wait(guard);
}
int item = items.front();
items.pop_front();
notFull.notify_one();
return item;
}
private:
std::deque<int> items;
size_t capacity;
std::mutex lock;
std::condition_variable notFull, notEmpty; // both use `lock`
};
const int STOP = -1; // a "poison pill": tells a consumer to quit
std::vector<std::vector<int>> run(int capacity, int producers,
int consumers, int perProducer) {
BoundedQueue q(capacity);
std::vector<std::vector<int>> taken(consumers);
std::vector<std::thread> makers, takers;
for (int p = 0; p < producers; p++) {
makers.emplace_back([&q, p, perProducer] {
for (int i = 0; i < perProducer; i++) q.put(p * perProducer + i);
});
}
for (int c = 0; c < consumers; c++) {
takers.emplace_back([&q, &taken, c] {
for (int item; (item = q.take()) != STOP;)
taken[c].push_back(item);
});
}
for (auto& t : makers) t.join();
// one pill per consumer, after all real items
for (int c = 0; c < consumers; c++) q.put(STOP);
for (auto& t : takers) t.join();
return taken;
}
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
class BoundedQueue {
private final ArrayDeque<Integer> items = new ArrayDeque<>();
private final int capacity;
private final ReentrantLock lock = new ReentrantLock();
// two conditions, one lock: both guard the same items
private final Condition notFull = lock.newCondition();
private final Condition notEmpty = lock.newCondition();
BoundedQueue(int capacity) { this.capacity = capacity; }
void put(int item) throws InterruptedException {
lock.lock();
try {
while (items.size() == capacity) {
notFull.await();
}
items.addLast(item);
notEmpty.signal();
} finally {
lock.unlock();
}
}
int take() throws InterruptedException {
lock.lock();
try {
while (items.isEmpty()) {
notEmpty.await();
}
int item = items.pollFirst();
notFull.signal();
return item;
} finally {
lock.unlock();
}
}
static final int STOP = -1; // a "poison pill": stops a consumer
static List<List<Integer>> run(int capacity, int producers,
int consumers, int perProducer) throws InterruptedException {
BoundedQueue q = new BoundedQueue(capacity);
List<List<Integer>> taken = new ArrayList<>();
List<Thread> makers = new ArrayList<>(), takers = new ArrayList<>();
for (int p = 0; p < producers; p++) {
int id = p;
makers.add(new Thread(() -> {
try {
for (int i = 0; i < perProducer; i++)
q.put(id * perProducer + i);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}));
}
for (int c = 0; c < consumers; c++) {
List<Integer> got = new ArrayList<>();
taken.add(got);
takers.add(new Thread(() -> {
try {
for (int item; (item = q.take()) != STOP; ) got.add(item);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}));
}
for (Thread t : makers) t.start();
for (Thread t : takers) t.start();
for (Thread t : makers) t.join();
// one pill per consumer, after all real items
for (int c = 0; c < consumers; c++) q.put(STOP);
for (Thread t : takers) t.join();
return taken;
}
}

To shut the consumers down, the main thread puts one STOP value, a “poison pill”, per consumer after all the real work. A consumer quits when it takes one. The queue is first in, first out, so every real item is taken before the pills.

In real code you’d use the library’s queue: Python’s queue.Queue(maxsize) or Java’s ArrayBlockingQueue. C++ has no standard one, so a class like this is what you’d write. Interviewers like asking for it because it uses every tool on this page.

Semaphores

A semaphore is a counter of permits. acquire() takes a permit, waiting while there are none left, and release() gives one back. A lock is a semaphore with a single permit. Use one to cap how many threads do something at the same time: open database connections, downloads in flight.

import threading
slots = threading.Semaphore(3) # at most 3 downloads at a time
def download(url):
with slots: # waits while 3 are already running
fetch(url)

Java has Semaphore and C++20 has std::counting_semaphore. A bounded queue can be built from two semaphores as well: one counts free slots, the other counts items, plus a lock for the list itself.

Deadlock

Moving money between two accounts needs both of their locks. Thread 1 runs transfer(x, y): it locks x, then y. At the same moment thread 2 runs transfer(y, x): it locks y, then x. Each now holds one lock and waits forever for the other. That’s a deadlock: a cycle of threads, each waiting for a lock the next one holds.

The standard fix is to take locks in one global order, for example by account id. Every thread then asks for the lower-numbered lock first, so no cycle can form.

def transfer(a, b, amount):
# every thread locks the lower id first
first, second = (a, b) if a.id < b.id else (b, a)
with first.lock:
with second.lock:
a.balance -= amount
b.balance += amount

Other ways out: one coarse lock for everything (simple, but only one transfer runs at a time), or acquire(timeout=...) and back off on failure. In C++, std::scoped_lock guard(m1, m2); locks several mutexes at once without deadlocking.

Why it’s O(1) per operation

Taking a free lock, a put and a take are each a few steps of work, plus waking one thread when someone is waiting. Big-O isn’t where concurrency gets expensive, though. The cost is waiting:

  • Contention. Only one thread can be inside a critical section. If each increment holds the lock for a microsecond, all your threads together manage at most a million increments a second, however many cores you have. Keep critical sections short: do the work outside, and lock only to read or publish the result.
  • Blocking. A thread that has to wait is put to sleep and later woken by the operating system. That costs microseconds, compared with tens of nanoseconds for taking a lock nobody else holds.
  • The serial part. If 10% of a job must run one thread at a time, no number of threads makes it more than 10 times faster (Amdahl’s law).
operation work can it wait?
lock.acquire() O(1) yes, while it’s held
put(item) O(1) yes, while full
take() O(1) yes, while empty
notify() O(1) no
notify_all() O(waiters) no, but all of them wake

The queue holds at most capacity items, so it takes O(n) space for a capacity of n.

Common mistakes

Waiting with if instead of while

A woken thread isn’t guaranteed that the condition still holds: another thread may have taken the item first, or the wakeup may be spurious. With if, it goes ahead anyway and pops from an empty queue.

if not self.items: # ✗ checks once, then trusts it
self.not_empty.wait()
while not self.items: # ✓ re-checks after every wakeup
self.not_empty.wait()

Trusting the GIL

The GIL runs one bytecode at a time, but count += 1 is several. Threads can still lose updates, and on Python’s build without a GIL they lose them often.

self.count += 1 # ✗ load, add, store: a switch fits between
with self.lock: # ✓ one thread at a time
self.count += 1

Checking, then acting, in two steps

Each call can be thread-safe on its own and the pair can still race: between the check and the action, another thread changes things.

if not q.full(): # ✗ another producer can fill it right here,
q.put_nowait(item) # and then this raises queue.Full
q.put(item) # ✓ checks and waits under one lock

A new lock every time

A lock only keeps out threads that use the same lock object. A lock created inside the method is a fresh one on every call, so nobody ever waits on it.

with threading.Lock(): # ✗ a new lock per call protects nothing
self.count += 1
with self.lock: # ✓ made once in __init__, shared
self.count += 1

Variations

  • Read-write locks. Many readers at once, or one writer alone. They pay off when reads far outnumber writes. Java has ReentrantReadWriteLock and C++ has std::shared_mutex. Python’s standard library has none, but you can build one from a Condition and a count of active readers.
  • Barriers. Each thread stops at the barrier until all n have arrived, then they all carry on. Simulations use one between phases. See threading.Barrier(n), Java’s CyclicBarrier and C++20’s std::barrier.
  • Thread pools. Start a fixed set of workers once and feed them tasks through a queue, instead of starting a thread per task. concurrent.futures.ThreadPoolExecutor and Java’s ExecutorService are producer and consumer with the queue built in.
  • Ring buffers. A bounded queue stored in a fixed array, with head and tail indexes that wrap around. With exactly one producer and one consumer, it can even work without a lock, using atomic indexes.
  • Atomics. For a single counter or flag, Java’s AtomicInteger.incrementAndGet() and C++‘s std::atomic<int> with fetch_add do the read-modify-write as one hardware instruction, with no lock. Python’s standard library has no atomic integers, so use a lock.

Check yourself

8 quick questions. Pick an answer to see why it's right or wrong.

  1. 1

    You must run a slow pure-Python function on 2,000 images, using an 8-core machine and standard CPython. What gives the biggest speedup?

  2. 2

    Three threads each add 1 to count. The lines below are the order their steps happened to run in. What does it print?

    count = 5
    a = count # thread A reads
    b = count # thread B reads
    a = a + 1
    count = a # A writes
    c = count # thread C reads
    b = b + 1
    count = b # B writes
    c = c + 1
    count = c # C writes
    print(count)
  3. 3

    This take works with one consumer but sometimes crashes with several. Why?

    def take(self):
    with self.lock:
    if not self.items:
    self.not_empty.wait()
    return self.items.popleft()
  4. 4

    Thread 1 calls transfer(x, y, 5) while thread 2 calls transfer(y, x, 3). What can go wrong, and what’s the usual fix?

    def transfer(a, b, amount):
    with a.lock:
    with b.lock:
    a.balance -= amount
    b.balance += amount
  5. 5

    Which statement about CPython’s GIL is true?

  6. 6

    What does this print?

    import threading
    s = threading.Semaphore(2)
    a = s.acquire(blocking=False)
    b = s.acquire(blocking=False)
    c = s.acquire(blocking=False)
    print(a, b, c)
    s.release()
    print(s.acquire(blocking=False))
  7. 7

    Five consumers are waiting on not_empty. A producer puts one item and calls not_empty.notify_all(). Assuming the consumers use a while loop, what happens?

  8. 8

    About how many seconds does this print?

    import asyncio, time
    async def handle(i):
    time.sleep(1) # a blocking call inside a coroutine
    return i
    async def main():
    start = time.perf_counter()
    await asyncio.gather(handle(1), handle(2), handle(3))
    print(round(time.perf_counter() - start))
    asyncio.run(main())

Practice problems

Solve these right here, in Python, C++ or Java. Tests run as you go.

All 39 problems on this topic

Further reading

esc