A per-record function passes when the whole job runs in one test process but fails on a real cluster - why?
answer
- one process, one memory
- the function has to travel
- whatever it captured travels too
- each worker mutates its own copy
- built once per worker, not once
basics
~20 sBecause the function has to reach worker processes on other machines, and whatever it captured has to go with it. In one process the capture is the same object in the same memory; distributed, it must be transportable, rebuilt per worker, or already present there.
solid answer
~50 sAn in-process run - the whole job executed inside the test's own process on one machine, with no network and no second worker - never moves the function anywhere. On a cluster it must reach separate processes, and the ways runtimes do that differ: some send the function together with everything it captured alongside each unit of work, some require the logic to be named code already installed on every worker, some call it in a separate process in another language runtime with records crossing both ways. Three failure classes follow. The capture drags something untransportable, such as an open connection or a file handle. A variable mutated inside the function is read back afterwards, which works in one memory and silently does nothing across many. A resource built once next to the wiring exists once in-process but must be built once per worker on the cluster.
go deeper
Remember that on a cluster the function does not stay where you wrote it: it has to reach other processes, and anything it captured has to reach them too. That is why a passing one-process run is not proof.
Explain the mechanics: capture and transport, a counter mutated per worker copy, a resource built once per process rather than once per job, and the fact that how the function reaches a worker differs between runtimes.
Show that you read the symptom - failure before the first record means transport, not logic - and that you keep functions free of captured resources as a standing rule rather than fixing them one failure at a time.
The lever is a convention: functions are top-level and take their resources as arguments, and every pipeline gets one cheap multi-process run. It costs little and removes an entire class of production-only failure.
## Two runs with two different memory models In an in-process run - the whole job executed inside the test's own process on one machine, with no network and no second worker - a per-record function is just a function. It closes over variables in the same heap; when it reads the enclosing object, that is the object; when it writes a counter, that is the counter. On a cluster the function has to be **in several processes at once**, usually on several machines. How it gets there is one of the things this product class genuinely disagrees about: - some runtimes package the function together with everything it captured and send that package with each unit of work; - some require the logic to be named code, already installed on every worker and referred to by name; - some call a hand-written function in a separate process running another language, with records crossing a boundary in both directions for every call; - and in the in-process run, nothing crosses at all - which is exactly why this class of defect is invisible there. ## The failure classes, and what each looks like 1. **An untransportable capture.** The function references a logger, an open connection, a file handle, a thread pool, or - more often - the enclosing object that happens to hold one of those. Whatever the runtime's transport, an object holding a live operating-system resource cannot be reconstituted on another machine. The symptom is a failure *before any record is processed*, which is the tell: the rule never ran, so the rule is not the bug. 2. **Accumulation in the enclosing scope.** A count or a list is updated inside the function and read after the run. In one process that is one variable. Across processes each worker mutates its own copy and the original is untouched, so the test's assertion passes and production reports zero. Runtimes provide an accumulator facility for this, and what it guarantees when a unit of work is recomputed after a failure varies - treat its number as approximate unless the runtime says otherwise. 3. **Per-worker initialisation.** A costly resource - a parsed model, a compiled pattern, a client - is built once next to the wiring. In-process that is one construction. Distributed it must be built once **per worker process**, so either the wiring arranges that, or the function builds it lazily and caches it per process. Get it wrong in the other direction and you build it once per record, which is a performance defect no unit test will show. 4. **Identity and singleton assumptions.** A cache, a counter, a deduplication set inside the function is per worker, never global. Logic that dedupes 'the records seen so far' is correct in one process and quietly wrong across twenty. 5. **Code and dependency availability.** The workers must have the classes or modules the function needs, at a compatible version. In-process the test's own dependency set is the answer; on the cluster it is whatever was installed there. 6. **Concurrency.** An in-process run is often single-threaded, or narrower than production. A worker process commonly handles several units of work concurrently, so shared mutable structure inside the function becomes contended - and the test never contended it. ## Reading the symptom | what you see | most likely cause | what would have caught it | |---|---|---| | failure before the first record is handled | an untransportable capture | a run with at least two separate worker processes | | output correct, side counter zero | accumulation in the enclosing scope | asserting on returned values rather than a variable | | enormous slowdown, correct output | a costly resource built per record | a per-worker construction count, or a timed run | | results differ between runs at the same input | per-worker cache treated as global | a run with more than one worker over the same input | | missing class or wrong behaviour on workers only | dependency drift between test and cluster | a run on the cluster's own environment | ## How to stop it happening - Make the function **free**: a top-level function, or a small object holding nothing but plain data. If it must be a method, check what its enclosing object holds. - **Take heavy resources as arguments** that the wiring supplies, or build them lazily and hold them per process, and say which of the two you chose. - **Return, never accumulate.** Rejects, counts and diagnostics come back as output records. - **Assume nothing about the transport.** Whether your function is shipped, named, or invoked across a language boundary is a property of the runtime you are on, and it decides what is allowed to cross. The honest framing for an interview is that this defect class is structural: it is not that the test was careless, it is that a single-process run has no boundary for the function to fail at.
- How should a per-record function report how many records it rejected?By returning the rejects as output records, which the wiring routes somewhere. If you instead use the runtime's counter facility, treat the number as approximate: when a unit of work is recomputed after a failure, some runtimes count its contribution twice, and what is guaranteed differs between them. A variable in the enclosing scope is never an option - each worker increments its own.
- Why can a function that is thread-safe in the test still be contended on a worker?Because an in-process run is often single-threaded or narrower than production, while a worker process commonly handles several units of work at once. Shared mutable structure inside the function - a cache, a buffer, a counter - is exercised by one caller in the test and by several on the cluster, so a race that exists in the code is never triggered until it runs for real.
saying these in an interview costs you the question
- Assumes a variable mutated inside the function can be read after the run.
- Thinks whatever works in one process works once workers are separate processes.
- Treats a cache built inside the function as shared by the whole job.
- Believes every runtime ships the function itself out to the workers.
- Builds a costly resource per record because the test never noticed the cost.