skip to content

A step rejects a small fraction of records in every run. What do you arrange in advance so the rejected ones can be examined?

level: seniorimportance: should knowfreq 50%

answer

  1. count everything, keep a few
  2. the bound is per worker, not global
  3. write it to shared storage
  4. project out the sensitive fields
  5. dropping without a threshold is loss

basics

~20 s

Count every rejection, but keep only a bounded sample of the records themselves: a fixed cap per worker process per run, written to shared storage under the run's identity, with the reason attached and sensitive fields left out.

solid answer

~50 s

Handle the failure per record inside the step rather than letting it end the unit of work. Do three things at once: increment a **carried counter** named for the reason, so the true volume survives the worker; append the record — or an identifying projection of it — to a sample that is **capped per worker process per run**, so the worst case is workers times the cap rather than a share of the data; and write that sample to **shared storage** under the run's identity, never to the worker's local disk, which vanishes with the machine. Capture only what reproduces the record: identifiers, the offending field, a hash of anything sensitive. Some engines offer a secondary output channel from a step; where they do not, it is simply a second output the program writes. Put a threshold on the counter — a step that drops records nobody is watching is silent data loss.

code

python · 29 lines
python
MAX_CAPTURED = 100          # per worker process, per job run

captured = []
rejected = 0                # carried out of the run and summed
capture_write_failures = 0

def handle(record):
    global rejected
    try:
        return transform(record)
    except Exception as exc:
        rejected += 1                       # counts every rejection
        if len(captured) < MAX_CAPTURED:    # keeps only the first few
            captured.append({
                "source": record.get("source_file"),
                "offset": record.get("offset"),
                "reason": type(exc).__name__,
                "field": record.get("amount"),
            })
        return None                         # drop the record, not the run

def on_worker_finish(run_id, worker_id, write_shared, report_counter):
    global capture_write_failures
    try:
        write_shared("rejects/" + run_id + "/" + worker_id + ".json", captured)
    except Exception:
        capture_write_failures += 1         # evidence must never end the run
    report_counter("rejected", rejected)
    report_counter("capture_write_failures", capture_write_failures)

go deeper

for a junior

Know the pair: a count for how many records were rejected and a small sample for what they looked like. One number and a few examples answer most questions about bad data.

for a middle

Explain why the sample must be bounded and why the bound is per worker process, and why the evidence goes to shared storage rather than the machine that produced it.

for a senior

Show you arrange this before the incident: reason-named counters, a cap, a projection that leaves sensitive values out, a threshold on the rate and a decided response when it is crossed.

for a principal

Own the dropping policy itself. Decide organisationally whether a pipeline may continue past bad data at all, who is told, and what the consumers of the numbers were promised about completeness.

## Why capturing everything is the wrong instinct The natural reaction to a rejected record is to print it. On a run of ten billion records, a rejection rate of one in a thousand is ten million printed records — the same order as the data, written by machines that are about to be released, through a path nobody sized for it. The capture becomes a second workload, and the evidence still disappears with its host. The working shape separates two questions that feel like one: - **How much?** Answered exactly, and cheaply, by a count. - **What did they look like?** Answered well enough by a handful of examples. Volume needs a counter; diagnosis needs a sample. Confusing them produces either a number you cannot act on or a pile you cannot afford. ## The shape of a bounded capture 1. **Catch per record, inside the step.** One malformed record should not end a unit of work, which the engine would then retry, fail again and eventually give up on. 2. **Count by reason.** A missing field, a failed lookup and an out-of-range value are three different facts and deserve three counters. A single `rejected` total hides which of them moved. 3. **Cap the sample per worker process per run.** A per-worker cap needs no coordination on a hot path; a global cap would make workers agree before every append. The worst case is bounded, and known in advance. 4. **Project before you keep.** Store the identifiers needed to find the record again, the reason, and the offending field. Where the payload carries personal data, keep a hash rather than the value. 5. **Write it where it survives.** Shared storage, partitioned by the run's identity and the worker's, so two runs never overwrite each other and finding yesterday's sample is a listing rather than an excavation. 6. **Never let the capture kill the run.** A failure while writing evidence must be swallowed and counted, not propagated. ## What it costs, and why the bound matters | Design | Worst case written | Coordination needed | Volume still known | |---|---|---|---| | Print every rejected record | Proportional to the data | None | Yes, if you can read it all | | Sample at a fixed rate | Proportional to the data | None | Yes, from the counter | | Cap per worker process per run | Workers times the cap | None | Yes, from the counter | | Cap globally across the run | The cap | Agreement on every append | Yes, from the counter | The third row is the default for a reason: the cost is fixed by cluster width rather than by how bad the day was, and a run that rejects a hundred records and a run that rejects a hundred million write the same amount of evidence. Only the counter tells them apart, which is exactly what the counter is for. ## The trap: capture without a threshold is silent data loss A step that catches, samples and continues is a step that drops data on purpose. That is a legitimate design — the alternative is a pipeline that stops dead on one bad row — but it is only safe while someone is told how much is being dropped. Two things make it safe: - an alarm on the rejection counter, on a **rate** rather than an absolute, so a change in the input is noticed rather than a long-standing background level; - a rule for what the job does when the rate crosses a line: stop, or write the whole batch aside, or continue and raise. Decide it before it happens, because at the time it happens the decision is made under pressure. The failure mode to picture is not a loud one. It is a schema change upstream that turns one field into a string, a rejection rate that goes from 0.1% to 40% overnight, and a dashboard that stays green because the job still finishes on time and nobody wired the counter to anything. ## Where engines differ Do not assume one mechanism. Some engines give a step a first-class secondary output that flows on as its own branch of the graph, which is the tidiest home for rejects; others have nothing of the kind, and you write the sample yourself as an additional output of the same program. Where a runtime is record-at-a-time, the write can be done as records arrive; in a continuous job built from repeated small finite runs, the natural boundary is the end of each slice. The shape above holds in all three cases, but the place you hang it does not.

  • Why cap per worker process rather than across the whole run?
    Because a global cap has to be agreed before every append, which puts coordination on the hottest path in the job for no diagnostic gain. A per-worker cap needs no agreement at all, and the worst case is still bounded and known: cluster width times the cap. If that product is too large, lower the cap.
  • What stops this pattern from becoming silent data loss?
    A threshold on the rejection rate and a decided response. Dropping a known-bad fraction is a fair design as long as the volume is visible and a sudden change reaches a person. Without that, an upstream schema change can push rejection from a fraction of a percent to nearly everything while the run still finishes on time and looks healthy.
  • The rejected records carry personal data. Does that kill the pattern?
    No, it narrows it. Keep what lets you find and reproduce the record — source identity, offset, reason, the shape of the offending field — and replace the sensitive value with a hash or a type-and-length description. In most investigations the value itself is not what you need; what you need is which parser assumption it broke.

saying these in an interview costs you the question

  • Printing every rejected record on a run that rejects millions
  • Writing the captured sample to the worker's own local disk
  • Swallowing the failure with no counter and no threshold on it
  • Storing whole records without thinking about the sensitive fields
  • Sizing the capture as a share of the data instead of a fixed cap
  • Letting a failed write of the evidence abort the job run