TT Lab
Get started
Learn Learning paths Courses

Working With Customer Data

It Worked on My Laptop, It Dies in the Pod

Continue in TT Lab

Goal

You build stream.py, a tool that handles large files without loading them into memory, and prove with an exit code, by applying a memory limit, that it really is streaming. You design aggregation, top N, distinct counting, external sorting, and verification all by the number of things you hold in your hand.

Why it matters

Code that reads a whole file is shorter and faster on a small file. So when you first write it there is no reason to suspect it, and it is not us but the customer who decides that the file grows. When a container exceeds its memory limit, the kernel kills the process, and that death is not caught as a Python exception and nothing is left on standard output. The criterion for the design is not the file size but the number of things you hold at any moment. An aggregate takes one line and the top N takes N items. But with a distinct count, memory is proportional to the number of distinct values even when streaming — if you do not know this difference, it becomes "I wrote it as streaming, so why does it die?" And "I ran it and it worked" is not proof. It only means it worked on that file at that time. Only passing with a limit applied is proof, and the result must remain as an exit code, not a sentence, so that the next person does not ask again. The grader does not trust your wording. It sets up logs it made in a temporary directory and actually runs your tool under a memory limit, checking the answers. The counts and amounts change on every run.

Steps

  1. Create and run /root/stream/gen_events.py to produce /root/stream/data/events.csv. It has 300,000 lines.
  2. Build agg <파일> (the placeholder is the file) in /root/stream/stream.py so that it reads one line at a time and outputs the count, sum, mean, minimum, and maximum.
  3. Add top <파일> <N> so that it keeps only the top N items with a heap.
  4. Add distinct <파일> <칼럼> so that it counts distinct values and outputs the peak memory at that time, and compare a column with few distinct values against one with many, and write it in /root/stream/distinct.json.
  5. Create /root/stream/cap.py and /root/stream/slurp.py, run the two versions side by side under a memory limit of 64MiB, and write the result in /root/stream/limit.json.
  6. Add sortmerge <파일> <칼럼> <청크행수> so that it does an external sort by split, sort, and merge.
  7. Add verify <파일> <칼럼> so that it outputs whether it is sorted and an order-independent fingerprint, and cross-check the original and the sorted copy and write it in /root/stream/verify.json.
  8. Summarize everything on one sheet and produce /root/stream/stream_report.json and /root/stream/stream_report.md.

Notes

Build a file that cannot be read in one go

Create and run /root/stream/gen_events.py to produce /root/stream/data/events.csv. The header is event_id,shop_id,kind,amount,ts and it has 300,000 lines.

Even when building it, do not gather into a list; write one line at a time to the file. Make shop_id have few distinct values and event_id differ on every line, so that you can later see the difference in memory for distinct counting.

Read one line at a time and aggregate

Build agg <파일> (the placeholder is the file) in /root/stream/stream.py so that it outputs rows, sum_amount, mean_amount, min_amount, max_amount, and peak_kib.

If you put the file object straight into a for loop, it is read one line at a time. The moment you call read() or readlines(), that property is lost. For the minimum and maximum you only need to compare with the value you just read, and the mean is produced at the end from the sum and the count.

Hold only the top N in your hand

Add top <파일> <N> so that it outputs the top N items by amount as [[event_id, amount], ...]. It is in descending order of amount, and when amounts are equal, the larger event_id goes on top.

If you sort everything and cut off the front, everything goes into memory. Keep a min-heap of size N, and once the heap has N items, push a new value in only when it is larger than the bottom. heappushpop of heapq does it in one step. If you make the values to put in the heap (amount, event_id) pairs, tie-breaking is handled along with it.

What is the memory of distinct counting proportional to

Add distinct <파일> <칼럼> so that it outputs column, rows, exact, and peak_kib. And write the results for the two columns shop_id and event_id in /root/stream/distinct.json as rows, low, and high.

Even when you read one line at a time, you have to remember the values already seen, so memory is proportional to the number of distinct values. If you measure a column with few distinct values and a column that differs on every line, the difference shows up in numbers. Put the result of the side with fewer distinct values in low and the side with more in high.

Prove it by applying a limit

Create /root/stream/cap.py and /root/stream/slurp.py, run stream.py agg and slurp.py each under a limit of 64MiB, and write limit_mib, file_bytes, rows, stream_exit, and slurp_exit in /root/stream/limit.json.

Linux has a limit on the address space a process can take, and in Python you can set it with the resource module. When you launch a child with subprocess, the limit has to be set on the child's side so that the parent does not die with it. cap.py itself ends with 0 even if the child dies, and only reports the child's exit code as JSON.

Do an external sort by split, sort, and merge

Add sortmerge <파일> <칼럼> <청크행수> so that it sorts each chunk and writes it to chunks/, and merges them and writes the result, with the header, to sorted.csv in the same directory. The response is rows, chunks, chunk_rows, output, and peak_kib.

Repeat gathering as many rows as the chunk row count, sorting them, writing them out to a file, and emptying the buffer. When merging, keep all the chunk files open but take only one line from each file, and heapq.merge takes a sort key and does that for you. Delete the chunk files made earlier at the start — if you do not, the results get mixed up on the next run.

Verify without intermediate products

Add verify <파일> <칼럼> so that it outputs rows, ordered, sum_amount, digest, and peak_kib, and cross-check the original and the sorted copy and write column, source, sorted, and same_multiset in /root/stream/verify.json.

If you load both the original and the result into lists to compare, the payoff of having sorted by streaming disappears. If you hash each line and add them all up, the same value comes out even if the order differs, and it differs if even one line is dropped or duplicated. Whether it is sorted can be known by remembering only the key of the previous line.

Leave the proof on one sheet

Write file_bytes, rows, sum_amount, chunks, limit_mib, stream_exit, slurp_exit, and digest_match in /root/stream/stream_report.json, and write /root/stream/stream_report.md in four sections: ## 무엇을 받았나 ## 어떻게 처리했나 ## 정말로 스트리밍인가 ## 남은 한계 (the Korean headings mean "What we received", "How we processed it", "Is it really streaming", and "Remaining limits").

You can just read and combine the limit.json, verify.json, and distinct.json you left in the earlier steps. In the report, write the limit value and the two exit codes as numbers — that is the proof, and in the remaining-limits section, write the fact that the memory of distinct counting is proportional to the number of distinct values.