Overview
I once spent a day optimizing a Python data pipeline, cutting its runtime from 40 minutes to 35 by rewriting loops as comprehensions and caching function results. Then I split it across four processes and it ran in 9 minutes. The lesson wasn't subtle: the GIL was the bottleneck, and no amount of single-threaded optimization was going to fix it.
This is what I've learned about multiprocessing since — when to use it, how to avoid the traps, and why it's not always the answer even when the GIL is in the way.
Why threading doesn't help
Python's Global Interpreter Lock means only one thread executes Python bytecode at a time. For I/O-bound work — waiting on network or disk — threads are fine, because the GIL releases during I/O. For CPU-bound work, threads give you concurrency without parallelism: the total work takes the same time or slightly longer due to context switching overhead.
import time
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
def cpu_work(n):
total = 0
for i in range(n):
total += i * i
return total
WORK = 50_000_000
# Serial
start = time.perf_counter()
for _ in range(4):
cpu_work(WORK)
print(f"Serial: {time.perf_counter() - start:.2f}s")
# Threads — no speedup
start = time.perf_counter()
with ThreadPoolExecutor(4) as ex:
list(ex.map(cpu_work, [WORK] * 4))
print(f"Threads: {time.perf_counter() - start:.2f}s")
# Processes — real speedup
start = time.perf_counter()
with ProcessPoolExecutor(4) as ex:
list(ex.map(cpu_work, [WORK] * 4))
print(f"Processes: {time.perf_counter() - start:.2f}s")
On a 4-core machine, the numbers come out roughly 12s / 12s / 3.5s. Threads buy nothing for CPU work; processes buy almost linear speedup.
The two APIs
| API | Style | Use when |
|---|---|---|
multiprocessing | Low-level, Pools and Processes | You need fine control |
concurrent.futures | High-level, futures | Default choice for most work |
concurrent.futures.ProcessPoolExecutor is what I reach for first. The API mirrors ThreadPoolExecutor, so swapping between them is a one-line change, and the code is easier to read than raw multiprocessing.Pool.
from concurrent.futures import ProcessPoolExecutor, as_completed
def process_chunk(chunk):
return sum(x * x for x in chunk)
chunks = [[1, 2, 3], [4, 5, 6], [7, 8, 9]]
with ProcessPoolExecutor(max_workers=4) as ex:
results = list(ex.map(process_chunk, chunks))
print(results) # [14, 77, 194]
map preserves order and blocks until all tasks finish. If you want results as they complete, use submit and as_completed:
with ProcessPoolExecutor(max_workers=4) as ex:
futures = {ex.submit(process_chunk, c): c for c in chunks}
for future in as_completed(futures):
chunk = futures[future]
try:
result = future.result()
print(f"Chunk {chunk} -> {result}")
except Exception as e:
print(f"Chunk {chunk} failed: {e}")
The pickling constraint
Arguments and return values are pickled to cross process boundaries. This is the source of most multiprocessing surprises.
# Works — plain function, plain data
def add(a, b):
return a + b
with ProcessPoolExecutor() as ex:
print(ex.submit(add, 1, 2).result()) # 3
# Fails — lambda can't be pickled
with ProcessPoolExecutor() as ex:
ex.submit(lambda x: x + 1, 5) # PicklingError
Lambdas, closures, locally-defined functions, and instances of most classes can't be pickled. The worker process needs to import the function by name, which means it has to exist at module level in a module that the worker can import.
# Don't do this
def make_processor(multiplier):
def process(x):
return x * multiplier
return process
with ProcessPoolExecutor() as ex:
# PicklingError: can't pickle local function 'process'
ex.map(make_processor(3), [1, 2, 3])
# Do this instead
def process_with_multiplier(args):
x, multiplier = args
return x * multiplier
with ProcessPoolExecutor() as ex:
ex.map(process_with_multiplier, [(1, 3), (2, 3), (3, 3)])
If you need state, pass it as an argument or use an initializer:
_GLOBAL_CONFIG = None
def init_worker(config):
global _GLOBAL_CONFIG
_GLOBAL_CONFIG = config
def worker(x):
return x * _GLOBAL_CONFIG["multiplier"]
with ProcessPoolExecutor(
max_workers=4,
initializer=init_worker,
initargs=({"multiplier": 3},),
) as ex:
print(list(ex.map(worker, [1, 2, 3]))) # [3, 6, 9]
The initializer runs once per worker process, which is much more efficient than passing the same config with every task.
Sharing state is hard by design
Each process has its own memory space. Unlike threads, you can't just share a variable and expect the change to propagate.
from multiprocessing import Process
counter = 0
def increment():
global counter
for _ in range(1000):
counter += 1
processes = [Process(target=increment) for _ in range(4)]
for p in processes:
p.start()
for p in processes:
p.join()
print(counter) # 0 — the child processes modified their own copies
This is correct behavior and it's the whole point — the isolation is what prevents the GIL from being a bottleneck. To actually share state, you need explicit mechanisms:
from multiprocessing import Value, Array, Manager
# Shared memory — fast
counter = Value("i", 0)
counter.value += 1
# Manager — slower, but supports arbitrary types
manager = Manager()
shared_dict = manager.dict()
shared_list = manager.list()
| Mechanism | Speed | Supports |
|---|---|---|
Value, Array | Fast | Numeric scalars, typed arrays |
Queue | Fast | Picklable objects |
Manager | Slow | dict, list, set, arbitrary |
| Files or Database | Varies | Anything |
Manager is convenient and slow. Each access is a network round trip to the manager process. If you're doing 100,000 operations on a shared dict, this dominates the runtime and negates the speedup.
The pattern that actually works: have each worker accumulate its own results locally, and combine them at the end.
def count_words(chunk):
counts = {}
for word in chunk:
counts[word] = counts.get(word, 0) + 1
return counts
def merge_counts(counts_list):
merged = {}
for counts in counts_list:
for word, n in counts.items():
merged[word] = merged.get(word, 0) + n
return merged
with ProcessPoolExecutor() as ex:
partial = list(ex.map(count_words, text_chunks))
final = merge_counts(partial)
No shared state at all. Each worker is independent, and the merge happens once. This scales linearly and avoids every locking problem.
Chunking: the setting that changes performance
ProcessPoolExecutor.map takes a chunksize argument, and it matters more than people expect.
# Default chunksize 1 — one task per item, lots of IPC overhead
ex.map(func, range(1_000_000))
# Chunksize 10000 — tasks are batched, less overhead
ex.map(func, range(1_000_000), chunksize=10000)
Every task has pickling and IPC overhead. If each task takes 0.1ms to run, the overhead dominates and multiprocessing is slower than serial. Chunking amortizes the overhead across many items.
The rule: chunksize ≈ total_items / (num_workers × 4). For 1,000,000 items and 8 workers, that's about 30,000. Benchmark before and after — the difference can be 10x.
When multiprocessing is the wrong answer
The overhead is real. Starting a process pool takes 50–200ms. Sending arguments, receiving results, and managing the pool adds up. For short tasks, the overhead exceeds the time saved.
| Workload | Right tool |
|---|---|
| Network I/O (API calls, scraping) | asyncio or threads |
| Disk I/O (reading files) | Threads |
| CPU-heavy computation | Multiprocessing |
| Numpy-heavy number crunching | Neither — Numpy releases the GIL |
| Short tasks with many items | Serial, or chunked multiprocessing |
| Mixed I/O and CPU | asyncio with run_in_executor for the CPU part |
The Numpy point is important. If your work is Numpy operations on large arrays, threads will parallelize it because Numpy's C code releases the GIL. Multiprocessing adds IPC overhead for no benefit. I've seen people rewrite Numpy-heavy code as multiprocessing and make it slower.
Windows and macOS gotchas
On Linux, the default start method is fork — the child process inherits the parent's memory. On Windows and macOS (since Python 3.8), the default is spawn, which starts a fresh Python interpreter and re-imports your modules.
Two consequences:
- The
if __name__ == "__main__":guard is required. Without it, spawn re-imports the main module and recursively spawns processes until the OS runs out of resources.
from concurrent.futures import ProcessPoolExecutor
def work(x):
return x * 2
if __name__ == "__main__":
with ProcessPoolExecutor() as ex:
print(list(ex.map(work, range(10))))
- Startup is slower on spawn. Each worker imports the module, which means module-level code runs again. Keep module-level code minimal, or the pool will take seconds to start.
The __main__ guard costs one line and prevents a class of confusing failures. Put it in from the start.
Choosing the executor
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
import os
def get_executor(workload_type):
if workload_type == "cpu":
return ProcessPoolExecutor(max_workers=os.cpu_count())
elif workload_type == "io":
return ThreadPoolExecutor(max_workers=32)
else:
raise ValueError(f"Unknown workload type: {workload_type}")
I use os.cpu_count() for CPU work and something like 32 for I/O work. For CPU, more workers than cores is counterproductive — context switching replaces actual work. For I/O, more workers than cores is fine, because most of them are waiting.
If you want one API for both, use asyncio.to_thread for I/O and run_in_executor with a process pool for CPU. That's the pattern I've settled on for mixed workloads.
