skip to content

In Luigi, why can a crashed task leave an output that makes the next run skip it?

level: seniorimportance: nice to knowfreq 35%

answer

  1. existence is not correctness
  2. the crash still left something behind
  3. write somewhere else first
  4. the move into place is the commit

basics

~20 s

Because Luigi decides completion from Target existence, a half-written file at the output path is indistinguishable from a finished one. The fix is atomic writes: produce the data at a temporary path and move it into place only on clean exit.

solid answer

~50 s

Luigi's completion check asks the Target whether it exists, not whether the previous attempt ended well. So if `run()` opens the final path, streams half the rows and then the process dies, the file is there and the next invocation happily skips the task while downstream tasks consume truncated data. The defence is to make the appearance of the Target the *commit point*. `luigi.LocalTarget.open('w')` already does this — it writes to a temporary file and renames it into place when the stream closes — which is why writing to `self.output().path` directly is discouraged. For work done by an external process, `self.output().temporary_path()` hands you a scratch path and moves it to the final one only if the block exits cleanly. For directory or object-store outputs, stage under a temporary prefix and move, or point the Target at a `_SUCCESS` marker written last.

code

python · 5 lines
python
# WRONG: streams straight to the final path
def run(self):
    with open(self.output().path, "w") as f:
        for row in rows():
            f.write(row)   # a crash here leaves a truncated, "complete" file

go deeper

for a junior

Learn the habit before the theory: write with self.output().open('w'), never by opening the output path yourself, because the Target's writer publishes the file only when it is finished.

for a middle

Explain why the framework cannot detect the problem — completion is existence, and a truncated file exists — and name temporary_path() as the tool for work done by an external process.

for a senior

Diagnose it live: trace bad downstream numbers back to a skipped task, prove the file was published mid-write, and prescribe stage-and-move or a marker file rather than a manual cleanup ritual.

for a principal

Make atomic publication a pipeline-wide convention with a shared base task, and be clear about what it does not buy you — complete-but-wrong output still needs validation and a reprocessing policy.

## The failure mode A task runs at 02:00, opens `/data/agg/orders_2026-08-20.csv`, writes 400,000 of 900,000 rows, and the worker is OOM-killed. At 02:30 the pipeline is re-invoked. Luigi rebuilds the graph, calls `complete()` on this task, the default implementation asks the Target whether it exists, the filesystem says yes — and the task is skipped. Every downstream task then runs happily on a truncated file. Nothing errors; the numbers are just wrong. This is the single most common production complaint against target-based completion, and it is a favourite senior interview scenario precisely because the framework is behaving exactly as designed. ## Why Luigi cannot tell Luigi keeps no run history to consult, so it has no memory that the earlier attempt died. Its only question is whether the promised artifact is present. Existence is a *proxy* for correctness, and a partially written file breaks the proxy: it is present but not correct. The framework cannot fix this on your behalf, because only the task knows what a complete result looks like. What it can do — and does — is give you the tools to make presence mean completeness. ## Make the appearance of the Target the commit point The rule is that data must become visible at its final location only when it is whole. Everything else follows from it. **Use the Target's own writer.** `luigi.LocalTarget.open('w')` does not open the final path. It writes to a temporary file and moves it into place when the stream is closed cleanly; if the process dies first, the final path never appears and the task is correctly still incomplete. This is why idiomatic Luigi code reads `with self.output().open('w') as f:` rather than `open(self.output().path, 'w')` — the latter throws away the protection you were given for free. **Use temporary_path() for external processes.** Plenty of tasks do not write bytes themselves: they shell out to a tool, hand a path to a library, or run a job that produces a file. `FileSystemTarget.temporary_path()` is a context manager that yields a scratch path and moves it to the Target's real path only if the block exits without an exception: ```python def run(self): with self.output().temporary_path() as tmp: subprocess.check_call(["render-report", "--out", tmp]) ``` **Stage-and-move for directories and object stores.** When the output is a directory of part files, or a prefix in an object store where there is no cheap atomic rename, write under a temporary prefix and move or copy into place at the end. Where even that is expensive, invert the check: write the data first, write a small `_SUCCESS` marker last, and point the Target — or an overridden `complete()` — at the marker. The marker is small enough that its write is effectively atomic, and it appears only after the data is whole. **Databases have their own version.** Write to a staging table and swap, or wrap the load in a transaction and let the commit be the visible event. A Target that checks for rows in a table being loaded row-by-row has the same partial-visibility bug as the file case. ## Strengthen the completion check where the shape is unusual When atomic publication genuinely is not available, override `complete()` to check something that partial output cannot fake: the presence of a marker, a recorded row count matching a manifest, a checksum. Keep it cheap and side-effect free, because Luigi calls it on every node on every scheduling pass. ## Operational habits around it Even with atomic writes, keep two habits. First, make outputs easy to delete at a known granularity — one directory per parameter value beats a single append-mostly file — because deleting the Target is Luigi's only rerun lever. Second, resist the temptation to "fix" incidents by hand-deleting files as a routine: if you find yourself doing it weekly, the pipeline is publishing non-atomically somewhere and the deletion is a symptom, not a remedy. ## Interview framing The tell of a weak answer is blaming the scheduler, or proposing that Luigi should validate contents before marking a task done. The strong answer names the mechanism — completion equals existence, so publication must be atomic — then names the concrete tools in order: the Target's own `open()`, `temporary_path()` for external writers, staging plus move or a `_SUCCESS` marker for directories and object stores, and a stricter `complete()` as the last resort. It also acknowledges the honest limitation: atomic publication protects against a crash mid-write, not against a task that runs to completion and produces confidently wrong numbers.

  • Your output is a directory of part files in object storage, where rename is not cheap. How do you keep completion honest?
    Write the parts under a temporary prefix and publish by moving or by writing a small `_SUCCESS` marker last, then point the Target — or an overridden `complete()` — at that marker rather than at the data prefix. The marker write is small enough to be effectively atomic, so it can only appear after the payload is whole.
  • Does atomic publication also protect you from a task that finishes but produces wrong numbers?
    No, and saying so is part of the answer. Atomicity only guarantees that whatever is visible was produced by a run that reached the end; it says nothing about correctness of the logic or completeness of the inputs. Guarding against bad-but-complete output needs validation inside `run()` before publication, or an assertion task downstream that fails the graph.

Print the report to a scratch tray and only slide it into the out-tray when the last page lands; anything in the out-tray is then trustworthy by construction.

saying these in an interview costs you the question

  • Blames the scheduler rather than the non-atomic write
  • Says Luigi validates output contents before marking a task done
  • Proposes hand-deleting outputs as the permanent fix
  • Writes directly to self.output().path and sees nothing wrong
  • Assumes a retry alone will overwrite the partial file safely

context