Why does passing a large list as a multiprocessing task argument cost so much?
answer
- Objects cannot leave an address space
- Serialize, pipe, rebuild — then back again
- Cost is per task, not per worker
- Peak memory holds several copies at once
- Send an index, or share the buffer
basics
~20 sWorker processes share no objects, so every argument is pickled in the parent, pushed through a pipe, and unpickled in the child, and the result makes the same trip back. A large list pays that cost on every task.
solid answer
~50 sNothing is shared between processes, so `multiprocessing` serializes each argument with `pickle`, writes the bytes to a pipe, and rebuilds the object in the worker; the return value repeats the journey. That is three costs, not one: CPU to traverse and rebuild the object graph, transient memory because the original, the byte string and the reconstructed copy can all be live at once, and allocator pressure from rebuilding millions of small objects one at a time. For a 2.4 GB working set the transfer dwarfs the work. The fixes are to stop sending it: send indices or slices instead of the whole structure, load the data once per worker in an initializer, chunk the work so the per-task overhead is amortized, use pickle protocol 5 out-of-band buffers for buffer-backed data, or put the bytes in a `multiprocessing.shared_memory.SharedMemory` block and send only its name.
code
python · 8 linesimport pickle, time
tickets = [{"id": i, "subject": "printer offline", "body": "x" * 80} for i in range(200_000)]
blob = pickle.dumps(tickets, protocol=pickle.HIGHEST_PROTOCOL)
start = time.perf_counter()
pickle.loads(blob)
print(f"{len(blob) / 1e6:.1f} MB, unpickle took {time.perf_counter() - start:.3f}s")go deeper
Remember that arguments to a worker process are copied, not shared, and that the copy is made by pickling. Knowing that a lambda cannot be sent to a worker is a good concrete detail to have ready.
Explain the full round trip and its three costs — CPU for graph traversal, transient peak memory from simultaneous copies, and allocator pressure in the child — and show that the cost repeats per task, not per worker.
Demonstrate diagnosis and redesign: measure serialization separately from the work, shrink arguments to identifiers, preload per worker, chunk to amortize overhead, and move bulk payloads into a shared block so only a name crosses.
Own the boundary as an architectural constraint. Decide when process parallelism is worth its data-movement tax at all against threads around GIL-releasing code, asyncio for IO, or simply a better algorithm, and set the team's convention for what is allowed to cross.
## What actually happens to an argument When you submit work to a worker process, the arguments do not travel as objects — objects cannot leave an address space. `multiprocessing` pickles them: `pickle` walks the object graph in the parent, produces a `bytes` payload, and that payload is written to a pipe or socket. The child reads the bytes and unpickles them, allocating a brand-new object graph in its own heap. When the function returns, its result is pickled in the child and unpickled in the parent. Under the `spawn` and `forkserver` start methods the callable itself is pickled too, which is why a lambda or a locally defined closure raises where a module-level function works. ## Three separate costs **CPU.** Pickling is a graph traversal with a per-object cost. A 2.4 GB working set made of a few large buffers pickles fast; the same 2.4 GB made of millions of small dicts and strings is dominated by per-object dispatch, and both directions pay it. Measure it before blaming the worker. **Peak memory.** During a transfer the parent can be holding the original object, the serialized `bytes`, and the pipe's buffered copy simultaneously, while the child holds the incoming bytes and the graph it is rebuilding. Sending a 2.4 GB structure to eight workers does not cost 2.4 GB; it can cost a multiple of it, and this is a classic way to be killed by the out-of-memory reaper on a box that "clearly had enough RAM". **Allocator pressure and GC.** Unpickling millions of small objects allocates millions of objects, which the child's cyclic collector then has to track and traverse. The child's memory is also permanently larger, because it now owns a full private copy. ## The failure mode people miss: per-task, not per-worker The cost is paid once *per task*, not once per worker. Mapping a function over 50,000 items while passing the whole reference dataset as a second argument serializes that dataset 50,000 times. This is why naive process parallelism is so often slower than the single-threaded loop it replaced. ## Fixes, roughly in order of preference 1. **Do not send it.** Send an index range, a key, or a slice, and let the worker fetch or compute what it needs. The smallest argument that identifies the work is the right argument. 2. **Load once per worker.** Pass an initializer that populates a module-level structure in the child at startup, so the payload crosses once per process instead of once per task. Pool and executor plumbing for that is `lang-python-concurrency`'s subject; the point here is that the data crosses N times instead of T times. 3. **Chunk.** Coarser tasks amortize a fixed per-task overhead. Thousands of microsecond tasks are dominated entirely by serialization and dispatch. 4. **Use out-of-band buffers.** Pickle protocol 5 (the default since 3.8, and the value of `pickle.DEFAULT_PROTOCOL`) can hand large buffer-backed payloads to the transport as `pickle.PickleBuffer` objects instead of copying them into the pickle stream. Any object exposing the buffer protocol benefits. 5. **Share the bytes.** Put the payload in a `multiprocessing.shared_memory.SharedMemory` block or a memory-mapped file and send only the name or path. The argument shrinks to a few dozen bytes and the copy disappears entirely. ## Things that are not the fix Raising the pickle protocol beyond the default rarely matters — the default is already the newest. Switching to a `multiprocessing.Manager` proxy usually makes things *worse*: it replaces one bulk copy with an IPC round trip per attribute access. And switching the start method changes what the child inherits, not what arguments cost: arguments are pickled under `fork` too. ## Diagnosing it The measurement that settles the argument is cheap: time `pickle.dumps` on one representative argument and `pickle.loads` on the resulting bytes, then compare that against the time the worker function takes on the same input. If serialization is a meaningful fraction of the task, no amount of adding workers helps — you are paying the tax more times, in parallel, on the same pipes. Watch the parent's memory during a burst of submissions too: a parent that grows while dispatching is buffering serialized payloads faster than the children drain them, which is a queue-depth problem dressed up as a memory problem. Also check for arguments that cannot be pickled at all. Open files, sockets, locks, database handles, lambdas and locally defined closures all raise on the way out, and the error surfaces in the parent at submission time rather than in the worker. ## How to tell in an interview The strong answer names the round trip in both directions, distinguishes CPU cost from peak-memory cost, and notices that the cost scales with the number of *tasks*. The weak answer says "pickle is slow" and stops.
- Why does a lambda work as a worker argument under fork but fail under spawn?Under `spawn` and `forkserver` the callable is pickled and sent to a fresh interpreter, and `pickle` stores functions by qualified name — a lambda or a closure defined inside another function has no importable name, so it raises. Under `fork` the child already has the function object in its inherited memory, so nothing needs pickling. Moving the function to module level fixes it everywhere.
- You cut per-task payload to a few bytes and the job is still slower than a single-threaded loop. What next?Suspect task granularity and the work itself. If each task is microseconds, dispatch and result marshalling dominate no matter how small the argument, so chunk more coarsely. Also check whether the work is genuinely CPU-bound in Python: if it already runs in a C extension that releases the GIL, threads avoid the process boundary entirely.
- Does returning a large result cost the same as sending a large argument?Yes, symmetrically — the result is pickled in the child and unpickled in the parent. Workers that return a full transformed dataset are a common hidden cost. Return an aggregate, a count, or a path to where the worker wrote its output, and keep bulk results in a shared buffer or on disk.
saying these in an interview costs you the question
- Says pickling is slow without naming the round trip both ways
- Thinks the payload crosses once per worker rather than per task
- Believes fork avoids pickling of task arguments
- Ignores peak memory: several copies are live at once
- Suggests a Manager proxy, which adds a round trip per access
- Blames the pickle protocol version when the default is already newest