Concurrency and Parallelism
A web server handling a thousand connections on one core and a video encoder using 16 cores are both "doing many things at once", and they are opposite problems. Concurrency is a matter of structure: many tasks in progress, interleaved. Parallelism is a matter of execution: many instructions running at the same instant on different cores.
This topic separates the two and states the law, from 1967, that caps how much faster more cores can make anything. The cap is set not by the hardware but by the part of the work that cannot be split, and that part is usually hiding somewhere nobody looked.
Interleaving and Simultaneity
Concurrency is making progress on several tasks by switching between them. It needs only one core: the core runs a slice of task A, then a slice of task B, then A again, and both are in progress the whole time. Parallelism is several tasks executing at the same instant, which needs several cores, the many-core hardware of Chapter 3.
The two are independent. An event loop, the last topic of this chapter, is concurrent and not parallel: one thread, thousands of connections in progress. A loop over a large array split across cores is parallel and not concurrent: one task, many hands. A server with a pool of threads on a many-core machine is both.
Concurrency Without Parallelism
Most server work is waiting. A request that takes 50 milliseconds may use 2 milliseconds of processor time and spend the other 48 waiting on a database and the network. While it waits, the core has nothing to do for it.
Interleaving fills that gap. One core can keep about 25 such requests in progress at once, starting the next while the previous one waits, and serve 25 times the traffic it could serve one request at a time. That win needs no second core at all. It needs only a way to switch to other work while one task waits, which is what threads and event loops provide.
Threads and Processes as the Two Vehicles
Processes, from Chapter 10, share nothing by default. They communicate by sending messages or through memory they deliberately arrange to share. Threads share everything: every thread in a process sees the same heap and can write the same variables. Communicating by writing memory is fast and needs no copying, and it is where the next topic's lost updates begin.
The vehicle decides the blast radius. A bug in one process corrupts that process. A bug in one thread can corrupt the memory of every thread beside it, and a crash ends them all.
Amdahl's Law
If part of a job must run serially, one step after another on a single core, that part caps the speed-up however many cores run the rest. Gene Amdahl made the argument in 1967. Take a job that is 10% serial. Even with unlimited cores the parallel 90% shrinks to nothing, and the serial 10% still takes a tenth of the original time, so the job runs at most 10 times faster. On 16 cores it runs 6.4 times faster, and a 17th core buys almost nothing.
The figure plots the same law for jobs that are 50, 90, 95 and 99% parallel. Each curve rises quickly over the first cores and then flattens toward its ceiling: 2 times, 10, 20 and 100. At 64 cores they have reached about 2, 8.8, 15.4 and 39 times. The half-serial job gains almost nothing past four cores.
The Serial Fraction Hides Everywhere
The serial part is rarely a line of code labelled "serial". It is a lock that every worker takes, the subject of the third topic in this chapter. It is one database row that every request updates. It is the interpreter lock of CPython's default build, which lets only one thread run Python code at a time, covered in the last topic. It is the final merge step after a parallel sort, and the log file with one writer.
Each of these is serial code in Amdahl's sense, whatever the architecture diagram says. The speed-up ends at the first one. Finding it is most of the work of scaling, and adding cores before finding it is paying for hardware the serial fraction will not let you use.
Where Scaling Stops
A continuous-integration pipeline with 32 parallel test shards and a 4-minute serial build step never finishes in under 4 minutes, however many shards it adds. A service scaled from 4 instances to 40 that all update one counter row scales until the row and then stops, with most of the 40 waiting their turn at that row. The hardware bill grows and the throughput graph goes flat.
For Python the choice follows from the workload. An input-and-output-bound service gets its concurrency from threads or an event loop on one core, and more cores add little. A CPU-bound one needs processes, or the free-threaded build of the last topic, to get parallelism at all. Backend Deep Dive covers choosing between them in practice.
- "A program with threads uses all my cores." Threads give concurrency. Whether they run in parallel depends on the runtime and on whether the threads have work to do or are waiting, and CPython's default build runs one thread's Python code at a time.
- "Doubling the cores halves the time." Only the parallel part shrinks. With 10% serial work, going from 8 cores to 16 takes the run from about 21% of the single-core time to about 16%, nowhere near half.
- "More threads is always faster." Past the core count for CPU work, or past the downstream limit, such as a database pool of 20 connections, for waiting work, extra threads add switching and contention and no throughput.
- "Concurrency bugs only happen on multicore machines." Interleaving on one core produces every race in the next topic. A single-core machine switches threads at arbitrary instructions and loses updates the same way.
- "Parallel speed-up is limited by the number of cores." It is limited first by the serial fraction. A 32-core machine running a job that is 25% serial tops out below 4 times faster, at about 3.7.
- Find the serial fraction before buying cores. Amdahl's ceiling, not the core count, bounds the speed-up.
- Use concurrency for waiting and parallelism for computing. Interleave waiting work on few cores, and spread CPU-bound work across processes or free threads.
- Size a worker pool by the resource it waits on. A thread pool larger than the database pool it feeds is a queue in disguise.
- Remove shared write points before adding workers. One hot row, lock or file turns every added worker into a waiter.
Knowledge Check
A job is 10% serial. Roughly how much faster can it run on 16 cores, and on unlimited cores?
- 16 times and unlimited, since each core takes an equal share of the work
- 9 times and 90 times, since 90% of the work gains from every added core
- About 6.4 times, and at most 10 times however many cores are added
- About 14 times, and at most 16 times, the number of cores in the machine
Which workload is concurrent but not parallel?
- One thread serving 2,000 idle connections from an event loop
- A matrix multiplication split into 16 equal parts on 16 cores
- A thread pool of 32 threads serving requests on a 32-core server
- Sixteen worker processes compressing sixteen files at the same time
A single-core virtual machine runs a threaded service. Can two threads still lose an update to a shared counter?
- No, because only one thread can run at a time on a single core
- No, because the kernel never preempts threads that are holding shared data
- Only if the virtual machine is later moved to a many-core host
- Yes, because a switch can happen between one thread's read and write
A CI pipeline runs 32 test shards in parallel after a 4-minute serial build. The team doubles the shards to 64. What happens to the total time?
- It roughly halves, because twice the shards run twice as many tests at once
- It shrinks a little at most, and never goes below the 4-minute build
- It stays exactly the same, because test shards never help build time
- It grows, because 64 shards always cost more than 32 to schedule
A service scaled from 4 to 40 instances handles barely more traffic than before, and all instances update one counter row per request. What is the most likely limit?
- The load balancer, which can spread traffic evenly over only 4 instances
- The instances' cores, which are now split between too many processes
- The one counter row, which every request must update in turn
- The network, which slows down as more instances join the service
You got correct