A finite job's final step holds 2,000 pieces and writes to a directory. How many files does it leave, and what set that number?
answer
- the write belongs to the piece
- nobody gathers the records first
- count at the end, not the read
- one thread, one output
- piece count becomes file count
basics
~20 sRoughly one file per piece, so about 2,000. Each worker thread writes the records it already holds, where it is; nothing gathers them first. The piece count at the final step, not the machine count, set the file count.
solid answer
~40 sWriting is the last thing a worker thread does with its piece, not a separate phase: the thread that computed those records opens an output at the destination and streams its own records into it. So the file count follows the piece count at the **final** step — about 2,000 here. It is the count at the end, not at the read: a step that regroups records by key sets a new count, and the last such change is what the write sees. Adding machines changes nothing, because the machine count only decides how many worker threads exist, not how many pieces there are. "About" rather than "exactly" because empty pieces, size-rolled outputs and grouped writes all bend the one-to-one relation.
go deeper
Recall the one-to-one shape: each worker thread writes the piece it holds, so the file count tracks the piece count. More machines do not change it.
Explain that it is the piece count at the final step, and name the cases that bend the one-to-one relation: empty pieces, size-rolled outputs, grouped writes, and a continuous job's commit intervals.
Show that you treat the file count as a decision, not an accident: the number chosen upstream is the layout every downstream reader inherits, and you can say who pays for it.
The angle is standardisation: a house expectation for output shape across many jobs, weighed against the per-job tuning it removes and the cost of retrofitting jobs that already run.
## One thread, one piece, one output A **piece of the input** is one slice of the stored input that a single worker thread reads and processes from start to finish. The **piece count** is how many such slices the run is divided into, and it is the ceiling on how many threads can be busy at once. Writing is not a phase bolted on at the end: it is the last thing each worker thread does with the piece it is already holding. The thread opens an output at the destination, streams its records into it, and closes it. Nothing collects the records first. Most engines do offer a way to pull a whole result back to **the coordinating process** — the single process that turns the program into a graph, decides the pieces and hands them out — and that is one way to get a single file; it is also bounded by that one process's memory, which makes it a debugging convenience rather than a write strategy. The consequence is arithmetic, not policy. Two thousand pieces at the final step means roughly two thousand data files. Nobody chose "2,000 files". Somebody chose a piece count, and the file count fell out of it. ## Which count, exactly Four different numbers all get called "the count", and answers routinely swap them: - **the piece count** — how many slices the run is divided into; - **the worker-thread count** — how many lanes exist to run them, where one piece occupies one lane for its whole life; - **the machine count** — how many machines the cluster has; - **the output-file count** — what the write leaves behind. The piece count sets how many threads can be busy; the machine count only sets how many threads exist; and only the piece count reaches the file count. Doubling the cluster does not double the files. One refinement carries most of the mark: it is the piece count **at the final step**, not at the read. A step that regroups records by key establishes a new count, and the last such change before the write is the one the file count follows. A job that reads ten thousand pieces and regroups to two hundred before writing leaves about two hundred files. How that number should be chosen is a separate subject, owned by `Choosing How Many Pieces`. ## Where "about one" stops being "exactly one" - a piece that produced no records leaves an empty file on some runtimes and nothing at all on others; - a runtime that rolls to a new output once the current one passes a size or row threshold leaves several files for one large piece; - a write grouped into **column-value directories** — files grouped into directories named for one column's value — emits one file per directory a piece touched, which multiplies the count; - destinations written through a staging location leave marker or temporary entries beside the data until the run is accepted; - some runtimes measure what the previous step actually produced and merge pieces before writing, while others write exactly what they hold. ## The continuous case is different in kind A **continuous job** is a run over an input with no end, so nothing about the input can be measured in advance and the division has to be declared: the author states a **declared operator width** — how many copies of each operator run — and that number stands until the job is restarted. There is no final piece count and no single write. Each instance emits the output for the interval it has just finished, commits it, and starts another. | | finite job | continuous job | |---|---|---| | what sets the file count | the piece count at the final step | declared width x number of commit intervals | | when it is known | once the last step's count is fixed | never final; it grows with elapsed time | | how it gets large | one very high piece count | a short commit interval running for weeks | That second column is why a continuous job writing to files is the reliable way to produce a million tiny files without anyone having typed a large number anywhere. ## Why an interviewer asks it Because the file count is a bill someone else pays. The next job that reads the directory generally derives one piece per file, so scheduling overhead can outweigh reading — that inheritance is owned by `Layout the Reader Inherits`. Listing a directory of hundreds of thousands of entries is slow on **an object store** (remote blob storage reached over the network, with no machine a piece of work could be placed on) and costly in metadata on **a cluster file system** (storage running on the same machines as the compute). Repairing the result afterwards belongs to `lakehouse table compaction`. A candidate who can say "the count I chose upstream is the count I left on disk" already has the whole mechanism.
- Does the piece count at read time decide the file count?No — only the count the final step holds does. A step that regroups records by key sets a new count, and the last such change before the write is what the file count follows. A job that reads ten thousand pieces and regroups to two hundred writes about two hundred files.
- Why might a single piece leave more than one file?Some runtimes roll to a new output once the current one passes a size or row threshold, so a large piece leaves several. A write grouped into column-value directories leaves one file per directory that piece touched. Staged destinations also leave marker entries beside the data until the run is accepted.
- Who pays for the file count this job leaves?The next reader, mostly: a job reading that directory generally derives one piece per file, so scheduling can outweigh reading, and listing the directory is itself slow at high entry counts. That inheritance is owned by `Layout the Reader Inherits`, and repairing the layout by `lakehouse table compaction`.
Two thousand people are each asked to post their own envelope, and two thousand envelopes arrive. Nobody decided on two thousand envelopes; they decided how many people would write. Deciding the number of writers is the only way to decide the number of envelopes.
saying these in an interview costs you the question
- Thinks results are gathered into one process and written as a single file.
- Says the machine count decides how many files the job leaves.
- Believes the file count is fixed by how many files were read.
- Claims every runtime merges small outputs automatically before writing.
- Thinks a continuous job keeps one fixed file per operator instance.