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

APIStyleUse when
multiprocessingLow-level, Pools and ProcessesYou need fine control
concurrent.futuresHigh-level, futuresDefault 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()
MechanismSpeedSupports
Value, ArrayFastNumeric scalars, typed arrays
QueueFastPicklable objects
ManagerSlowdict, list, set, arbitrary
Files or DatabaseVariesAnything

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.

WorkloadRight tool
Network I/O (API calls, scraping)asyncio or threads
Disk I/O (reading files)Threads
CPU-heavy computationMultiprocessing
Numpy-heavy number crunchingNeither — Numpy releases the GIL
Short tasks with many itemsSerial, or chunked multiprocessing
Mixed I/O and CPUasyncio 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:

  1. 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))))
  1. 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.