A distributed task can run more than once, so its candidate files and publication step must prevent duplicate or partial output.
Task retries and atomic partition output
Assume repeated attempts
A worker may write a file and lose its completion acknowledgement. The scheduler retries on another worker; both attempts can now exist. If a reader lists every file in a directory, it may count the partition twice. Assign each attempt a private path and publish only one accepted attempt through a manifest or transactional table commit.
Separate compute from commit
A task result is a candidate until the coordinator validates it. Capture task identity, attempt ID, input snapshot, row count and digest. The commit layer selects one accepted attempt per task and publishes a single generation pointer. Snapshot publication keeps readers on the last complete generation while the new one is checked.
Make cleanup conservative
An abandoned attempt file is not necessarily safe to delete while its writer or a delayed coordinator is still active. Use an age threshold beyond the maximum attempt lifetime, then compare files against all live manifests. Log the proposed deletions before applying them. Retention rules must cover retry and rollback windows.
Reconcile at the output grain
For a daily invoice partition, check one row per invoice line or one row per account-day according to the declared target grain. Validate row count and total cents against the pinned input. A successful task state does not prove all expected partitions committed. Quality gates should precede the final pointer update.
Inject the uncertain outcome
Test a failure after the file write and before acknowledgement. Two candidate attempts should appear, yet the published manifest should contain exactly one for that task. Rerun from the same input snapshot and compare the canonical output digest; a new source snapshot is a new run, not an idempotent retry.
Implementation
task_attempts = [
{"task": "day-47-part-2", "attempt": 1, "file": "candidate-a", "valid": True},
{"task": "day-47-part-2", "attempt": 2, "file": "candidate-b", "valid": True},
{"task": "day-47-part-3", "attempt": 1, "file": "candidate-c", "valid": True},
]
def accepted_manifest(attempts):
accepted = {}
for attempt in sorted(attempts, key=lambda row: (row["task"], row["attempt"])):
if attempt["valid"]:
accepted.setdefault(attempt["task"], attempt["file"])
return accepted
manifest = accepted_manifest(task_attempts)
assert manifest == {"day-47-part-2": "candidate-a", "day-47-part-3": "candidate-c"}
assert "candidate-b" not in manifest.values()Performance and operating cost
Sorting A attempt records costs O(A log A) and the manifest uses O(T) space for T logical tasks. A production coordinator may use keyed storage instead, but still needs durable conflict handling. Candidate files temporarily increase storage; that is preferable to silently publishing two attempts.
Common Mistakes
- Do not expose candidate directories as committed datasets.
- Do not equate one successful attempt with one accepted output.
- Do not delete unreferenced files before in-flight attempts are ruled out.
Read next
- Table snapshots and atomic publication
- DAG intervals and idempotent task outputs
- Data quality gates: quarantine bad rows and reconcile complete batches
- Project: release a distributed invoice rollup
- Replay manifests and audit trails
Continue the workflow: Partition registration and manifest gates.
