skip to content

A 6-hour nightly webhook-log replay dies on memory - how do you restructure it to stream?

level: seniorimportance: should knowfreq 42%

answer

  1. Find what holds the whole run
  2. Memory should not scale with window length
  3. Every step becomes a stage, not a list
  4. Watch the last incomplete group

basics

~20 s

Find the stage that materializes - readlines(), a parsed list, a results accumulator, a sort - and make each one a generator stage so a single record flows end to end. Batch with itertools.batched to keep the final short batch.

solid answer

~50 s

Work backwards from peak memory to the first stage that holds everything. The usual culprits are reading the whole file at once instead of iterating the open file object line by line, parsing into a list before filtering, and appending every result to a list before writing. Rewrite each as a generator stage - lines, then parsed records, then filtered records - so a single delivery record crosses the whole chain and is written out before the next is read; memory then tracks one record, not six hours of them. Where the sink needs groups, batch with `itertools.batched` (3.12+) rather than a hand-rolled counter, because the classic bug is dropping the final partial batch when the record count is not a multiple of the batch size. Anything that genuinely needs to see everything - a sort, a dedupe set - must be bounded explicitly or pushed into storage.

code

python · 9 lines
python
import io
import itertools

log = io.StringIO("".join(f"delivery-{n}\n" for n in range(9)))

lines = (line.rstrip("\n") for line in log)          # one line alive at a time
retries = (name for name in lines if not name.endswith(("0", "4")))
for batch in itertools.batched(retries, 3):          # 3.12+, keeps the short last batch
    print(batch)

go deeper

for a junior

Know the core habit: iterate an open file object line by line instead of calling readlines(), and avoid building a list of every parsed record. That one change is usually the difference between flat and unbounded memory.

for a middle

Explain how to convert each step into a generator stage so a single record crosses the whole chain, and name the accumulators that break it - a results list, a sort, a dedupe set, or a stray list() added while debugging.

for a senior

Diagnose it in production terms: locate the materializing stage, prove memory no longer scales with the window length, and handle the batch boundary correctly so the final partial batch is neither dropped nor sent twice.

for a principal

Own the wider design: which operations genuinely require global state, whether to bound them by keys or push them into storage, and whether the sink should be idempotent so a resumed or re-run replay is safe rather than merely memory-efficient.

A nightly replay job walks six hours of webhook delivery records, re-sends the failures, and dies partway through with the machine out of memory. The fix is almost never a bigger box; it is finding the stage that accumulates. ## Locate the materializing stage Read the job backwards from the point of peak memory and look for anything that holds the whole run at once: - `rows = handle.readlines()` - the entire file as a list of strings, up front. - `records = [json.loads(line) for line in rows]` - now the same data again, as parsed objects, which are several times larger than the text they came from. - `failed = [r for r in records if r["status"] >= 500]` - a third copy. - `results = []` with `results.append(...)` in the loop, written only after the loop ends. - A `sorted(...)` or a growing `set` used for deduplication. Any one of these turns a job whose *working set* is one record into a job whose working set is the whole night. Six hours of deliveries at a modest rate is millions of records, and parsed objects carry per-object overhead that raw text does not. ## Make every stage a generator The streaming shape keeps the same logical steps but changes each from a collection into a stage: ```python def records(handle): for line in handle: # the file object yields one line at a time yield json.loads(line) def failures(rows): for row in rows: if row["status"] >= 500: yield row with open(path, encoding="utf-8") as handle: for row in failures(records(handle)): resend(row) ``` Now exactly one line, one parsed record and one in-flight send exist at a time. Peak memory becomes a property of the largest single record, not of the run length - which is the whole point: the job's memory no longer scales with how long the window is, so extending it from six hours to a day changes nothing. Two habits keep it that way. Write results out inside the loop instead of collecting them, and never insert a `list()` "just to check something" - a debugging line like `print(len(list(rows)))` both re-materializes the stream and consumes it. ## Batching, and the boundary bug Downstream sinks usually want groups, not single records - a batch endpoint, a bulk insert, a chunked upload. Hand-rolled batching is where this job breaks in the most annoying way: ```python batch = [] for row in rows: batch.append(row) if len(batch) == 500: send(batch) batch = [] # if the loop ends here, the last partial batch is silently dropped ``` If the record count is not an exact multiple of 500, the final short batch never ships. Nothing raises, nothing logs, and the loss is invisible until someone reconciles counts weeks later. The mirror-image bug is flushing inside the loop *and* after it without clearing, which sends the boundary batch twice - duplicate deliveries, which for a webhook replay is worse than a drop. Since Python 3.12 the stdlib does this correctly: `itertools.batched(rows, 500)` yields tuples and keeps the short final batch, and passing `strict=True` raises instead if you require exact multiples. If you must hand-roll, the flush after the loop guarded by `if batch:` is the line that must exist. ## What cannot stream, and how to bound it Some requirements genuinely need to see everything: global ordering, deduplication across the whole window, "the last event per key". These are barriers, and pretending otherwise just moves the memory spike. The options are to bound the state instead of the data - a dedupe set holding fixed-width keys rather than whole records is often two orders of magnitude smaller - to sort within batches when global order is not truly required, or to push the operation into a database or an external sort. Aggregates should be running values: a `collections.Counter` of status codes costs one entry per distinct code, regardless of how many records were counted. ## Confirming the fix Once every stage is a generator, memory should be flat and roughly independent of input size - so a run over ten minutes of logs and a run over six hours should look the same in RSS. If it still grows, the accumulator survived somewhere: a cache, a list built in an inner loop, or a stage that sorts. Checkpointing progress also becomes possible in the streaming shape - because records are processed one at a time, a crashed run can resume from the last acknowledged offset rather than starting the six hours again.

  • How do you tell whether the job is really streaming after the rewrite?
    Memory should be flat and independent of how much input you feed it: a run over ten minutes of records and a run over six hours should show the same steady-state footprint. If it still climbs with input size, an accumulator survived - a results list, a growing dedupe set, a cache, or a sorting stage acting as a barrier.
  • The replay needs deduplication across the whole six-hour window. How do you keep memory bounded?
    Do not hold records to dedupe them - hold keys. A set of fixed-width delivery ids is a fraction of the size of the parsed records, and often small enough to be fine. If even that is too large, push the check into storage with a unique constraint or a bounded structure, or make the sink idempotent so a duplicate send is harmless, which is usually the more robust answer for a replay.
  • What does the streaming shape buy beyond memory for a long-running nightly job?
    Restartability and earlier feedback. Because records are handled one at a time, the job can checkpoint the last acknowledged offset and resume there after a crash rather than repeating hours of work, and results appear from the first record instead of only after the whole file is parsed. It also short-circuits cleanly when a run is cancelled.

saying these in an interview costs you the question

  • Answers with more RAM or a larger machine instead of streaming
  • Calls readlines() and considers the job a streaming pipeline
  • Appends every result to a list and writes only at the end
  • Hand-rolls batching without flushing the final partial batch
  • Flushes inside and after the loop, sending the boundary batch twice
  • Claims a sorting stage preserves the pipeline's flat memory

context