TT Lab
Get started
Learn Learning paths Courses

Working With Customer Data

The File Does Not Fit in Memory

Continue in TT Lab

In one line

"Read it all and put it in a list" is a design that is right only when the file is small, and whether that assumption has broken can be known only by putting a limit on it.

Why this was needed

The customer sent a month of event logs. A script that ran fine on the laptop dies in the Pod. There is nothing in the log and only the process has disappeared.

The cause is usually one line. It is code that loads the whole file, like lines = open(path).read().splitlines(). On a small file this line is even shorter and faster, 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.

What is more troublesome is that this failure is silent. When a container exceeds its memory limit, the kernel kills the process. Python's exception handling cannot catch that death, and nothing is left on standard output. All you see the next morning is "the job didn't finish."

How it works

There is a single criterion for the design. At any moment, how many things do you hold in your hand?

An aggregate takes one line. Count, sum, minimum, and maximum can be updated with just the line you just read. Even if the file becomes 10 times larger, memory stays the same. If you iterate over a Python file object as it is, it is read one line at a time — the moment you call readlines() or read(), that property is lost.

The top N takes N items. If you sort everything and cut off the front, everything goes into memory, but if you keep a min-heap of size N and push a new value in only when it is larger than the bottom of the heap, what you hold is always N items. The heappushpop of heapq does this in one step.

Counting distinct values is different. 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. Counting 17 region codes and counting 500,000 session ids are the same streaming code, but the memory is completely different. So when designing a distinct count, the first thing to ask is not the file size but the number of distinct values.

Sorting is done in pieces. External sorting is an old method. You read only as much as fits, sort it, and write it to a temporary file (split and sort), then take lines out of those files one at a time and merge them (merge). What you hold in the merge stage is only one line per file. In Python, heapq.merge does this merging for you. Unix sort does the same thing — POSIX sort is specified on the premise that it uses temporary files, and in the GNU implementation -S sets the buffer size and -T sets the temporary directory. If temporary space runs short, the sort fails.

What it looks like in the field

First, you use "I ran it and it worked" as proof. That statement means it worked on that file at that time. Proof is passing with a limit applied. Linux has a limit on the address space a process can take, and in Python you set RLIMIT_AS with setrlimit of the resource module. Under the limit, the streaming version passes and the version that reads everything dies with MemoryError. Those two exit codes are better than ten sentences.

Second, you pile up intermediate products and verify them. If you read both the original and the result into lists to compare "to see whether the sort went well", the payoff of having sorted by streaming disappears. Instead, compute an order-independent fingerprint while streaming. If you hash each line and add them all up (discarding the overflow), the same value comes out even if the order differs, and it differs if even one line is dropped or duplicated. Looking at the sum and the count together is enough.

Third, you forget the temporary space. An external sort is a trade that moves memory to disk. The Pod's temporary disk is 6Gi, and that space is shared too. If you do not delete the chunk files, the disk fills up first in the middle of sorting.

Fourth, you pick the chunk size by feel. If the chunks are too large, it dies in memory, and if they are too small, the file handles reach hundreds and the merge is slow. Decide the chunk size by confirming "this size will surely fit" with a limit test.

What really matters in practice

What you will do in the next lab

You build a 300,000-line event log yourself and build a tool that reads it one line at a time and aggregates. You keep only the top N with a heap and measure for yourself how the peak memory of a distinct count changes with the number of distinct values. Then you apply a memory limit and run the streaming version and the version that reads everything side by side, leaving the difference as exit codes, and after sorting a large file with an external sort, you prove with an order-independent fingerprint that not a single line was lost. The grader runs the files it made with your tool under the limit and checks the answers.