~/techniques/concurrency
Concurrency
Threads that share data can lose updates. Locks, condition variables and a bounded queue make them take turns safely.
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.
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.
O(1)
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, liketime.sleepor 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,
hashlibandzlibalso 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 += 1is 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.
- Start:
count = 0, and each thread has its owntmp. The CPU runs one step at a time and may switch threads after any step. - A reads
countinto itstmp: 0. - The scheduler switches to B before A has written anything back. B reads
counttoo, and gets the same 0. - Each thread adds 1 to its own
tmp. Both now hold 1, andcountis still 0. - A writes its 1 to
count. So far, so good. - 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.
- Two increments ran, and
count = 1. That’s a race condition: the result depends on timing, and here one update is lost. - The same schedule, now with a
lock. A takes it before reading. - B asks for the
lockwhile A holds it, so B sleeps. It can’t readcountuntil A is completely done. - A writes 1 and releases the
lock. B wakes up, takes thelock, and reads the new value, 1. count = 2. The read, the add and the write are now one step that no other thread can cut into: a critical section.
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
whileloop. 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, wherewait()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.
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
ReentrantReadWriteLockand C++ hasstd::shared_mutex. Python’s standard library has none, but you can build one from aConditionand a count of active readers. - Barriers. Each thread stops at the barrier until all
nhave arrived, then they all carry on. Simulations use one between phases. Seethreading.Barrier(n), Java’sCyclicBarrierand C++20’sstd::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.ThreadPoolExecutorand Java’sExecutorServiceare producer and consumer with the queue built in. - Ring buffers. A bounded queue stored in a fixed array, with
headandtailindexes 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++‘sstd::atomic<int>withfetch_adddo 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
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?
The work is CPU-bound Python code, and the GIL lets only one thread run Python code at a time, so threads take turns on one core. Separate processes each have their own interpreter and GIL, so 8 of them really use 8 cores. asyncio is one thread too, and only helps while coroutines wait on I/O.
-
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 = 5a = count # thread A readsb = count # thread B readsa = a + 1count = a # A writesc = count # thread C readsb = b + 1count = b # B writesc = c + 1count = c # C writesprint(count)A writes 6. C then reads that 6. B writes its 6, worked out from the 5 it read at the start, so B’s update is lost. C writes 6 + 1 = 7. Three increments ran but
countonly went up by 2. With a lock around each read-add-write, it would be 8. -
3
This
takeworks 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()Between the notify and this thread getting the lock back, another consumer can take the item, and some platforms wake threads with no notify at all.
while not self.items:re-checks every time.notify_all()would make it worse: more threads wake for the same single item. -
4
Thread 1 calls
transfer(x, y, 5)while thread 2 callstransfer(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 -= amountb.balance += amountThread 1 holds
x.lockand wantsy.lock; thread 2 holdsy.lockand wantsx.lock: a deadlock. If every thread takes locks in one global order, such as by id, no cycle of waiting can form. Updates aren’t lost here: both balances are only changed while both locks are held. -
5
Which statement about CPython’s GIL is true?
The GIL guards the interpreter one bytecode at a time, and
count += 1is several bytecodes: load, add, store. A switch between them loses an update. Threads waiting on the network release the GIL, so I/O does overlap, and each process has its own GIL. -
6
What does this print?
import threadings = 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))The semaphore starts with 2 permits. The first two acquires take them; the third finds none left and, because it doesn’t block, returns
Falseat once.release()gives one permit back, so the last acquire succeeds. -
7
Five consumers are waiting on
not_empty. A producer puts one item and callsnot_empty.notify_all(). Assuming the consumers use awhileloop, what happens?notify_all()wakes every waiter, and they take turns getting the lock back. The first one takes the item; the rest see an empty queue in theirwhilecheck and go back to sleep. It’s correct, but it wastes four wakeups, which is why the queue usesnotify()for one item. -
8
About how many seconds does this print?
import asyncio, timeasync def handle(i):time.sleep(1) # a blocking call inside a coroutinereturn iasync def main():start = time.perf_counter()await asyncio.gather(handle(1), handle(2), handle(3))print(round(time.perf_counter() - start))asyncio.run(main())time.sleepblocks the one thread the event loop runs on, and a coroutine only gives up control at anawait. So the three handlers run one after another: 3 seconds. Withawait asyncio.sleep(1)they would overlap and finish in about 1 second.
Practice problems
Solve these right here, in Python, C++ or Java. Tests run as you go.