When your data doesn’t fit in memory: the basic techniques
pythonspeed.com
pythonspeed.com
A semi-common beginning programmer's exercise is to write a program that numbers the lines in a text file. The naive solution will use O(n) space, while a bit more thought reveals that this can be done in constant (to be really precise, O(log n) where n is the number of lines in the file) space.
The hard part for me is the "transpose" or "striding" problem. i.e. when the data is stored in a series of (x, y, z) files for a given hour, day, or month and I need a time series at a point.
Many modern performant solutions will use both, but they're not the same thing.
Mass storage is most efficient when doing large sequential reads and writes, so you normally feed your constant-space streaming algorithms from buffers with a large number of input records.
Sometimes you can just tell the OS do efficient chunking prefetch for you.
If your intended audience are data scientists then why didn’t you mention Dask?
In term of RAM this can be done an arbitrarily small buffer in the sense that this can be done by reading the file byte by byte. I.e. O(1).
Edit: For all practical purposes the line counter is also constant space.
printf(“0: “);
ln = 1;
while ((ch = getch()) != EOF) {
putch(ch);
if (ch == EOL) {
printf(“%d: “, ln);
ln++;
}
}
And I believe that original unix implementation of nl is mostly this, although with the stdio functions replaced with raw read/write.Or "database"
That's on point, some of the worst and most inefficient code I had to maintain was written by data scientists. They are for sure leagues above any of my data science knowledge in their respective field but the average engineering knowledge and best practices seems pretty low.
"The simplest and most common way to implement indexing is by naming files in a directory:"
mydata/
2019-Jan.csv
2019-Feb.csv
2019-Mar.csv
2019-Apr.csv
Um. mydata/
2019/
11/
12/
2019111200.csv
2019111201.csv
Where each one of those CSV files could be 10-100 GB in size.Usually you want to process it into a columnar format like Parquet from there, though
Ah, simple :)
cat file | wc -l wc -l file
To hide filename: wc -l < fileAnyway, since he said number we probably want `nl` not `wc`.
Separating pipe-setup from processing is sound engineering to me and makes things much easier. For instance as I iterate, I will have cat piped to head piped to filter but eventually I'll take the head out and run the whole thing. That's a trivial ^W versus moving around arguments and worrying about arg order, etc.
However using the input redirect is fine, especially if you format it that way (which is exactly the same thing):
<file wc -l
Then, adding another step is as natural as with cat: <file grep meh |wc -l <file grep -c mehI know `< file wc -l | ... | ...` is identical because the redirection is done before the command is executed, but what does putting it in front help with?
And not change anything else in the pipeline.
The first iteration looks like this:
<file head -n10 | grep foo | tr baz qux
The second iteration looks like this:
<file gzcat | head -n10 | grep foo | tr baz qux
It makes the change much easier, since the first thing written is the first thing that happens. In `gzcat < file`, the logical first step - reading and streaming file - is now the physical second step. Like a German sentence, whenever you want to prepend the file, you have to maintain an item in second position.
head filename | ...
then just change "head" to "cat" when I get it working, or grep if I want to check parts of a file, etc.In this particular case, people (like me) would wonder, "what's the point" and then go searching for the details of how `cat` and `wc` work, assuming that there's might be a reason the original developer wrote this code, as opposed to the simpler `wc -l` (i.e. Chesterton's Fence).
Also, context is really important. Not that people on the internet ever stop to think about what the context is that the person is writing the code in. I don't write bash scripts that run mission critical workloads. I just ad-hoc type stuff into my shell to produce output I need for stuff at the time.
cat -n file nl fileMy first inclination would be to have a counter at zero, read in a line, write out the counter then write out the line, then increment the counter.
The space complexity here is O(1), is it not?
We have very good asymptotic complexity for multiplication, but don't use the ones that scale the closest to linear because the time constant is enormous.
Also relevant here is the RAM machine model that is often employed, where you can have essentially unlimited memory that can be randomly accessed at fixed time. This again is roughly in line with real practical machines running algorithms that use, and fit within RAM. In reality if you wanted a more physically consistent model eventually your access times must differentiate and increase for increasing memory -- data occupies physical space and must be fetched at increasingly long distances, limited by the speed of light. In principle this means m bits of memory at best can be accessed in O(m^(1/3)) time [1]. Of course, this realism isn't always practically relevant (maybe for datacenter-scale problems, perhaps even larger) and complicates the analysis of algorithms.
[1]: Just for fun, if you want to get really physically accurate, this isn't quite right either I suspect -- that's because with enough data (physical bits) occupying a volume at constant density it will eventually collapse into a black hole, which kinda destroys your computer :). So for planetary-scale computers you eventually need to spread your data across a disk, like a small galaxy (or Discworld, if you prefer :)), so it won't collapse, giving it O(sqrt(m)) access time. Surprise, very large worlds must be flat.
This is related to the BH entropy formula, and the so called Holographic principle, I guess -- the entropy and hence information content of a volume is surprisingly bounded by its area, since a Black Hole's area is proportional to its mass. The weird thing of course is how the universe itself perhaps should collapse to a black hole since it has a lot of stuff in any sufficiently large volume, being young and dynamic it hasn't happened yet. It's alsogiven by its fractal scale, but I guess this digression has grown large enough already :)
(do tell if you want to learn more)
So integers count as constant memory AND can be considered to be arbitrarily large for the purpose of the current problem. Because having more than 2^64 lines in a file will never happen in practice, and even if it does then 2 integers would do the job for the rest of the observable universe.
Be pragmatic, you are on Hacker News.
In my experience, practitioners just have a lot of misconceptions about "Big O".
On a 64-bit CPU, y=x+1 is O(1) time when Y < 2^64, but becomes O(log N) time for larger numbers. With space it’s even more complex, sometimes for space saving it makes sense to keep less than 64 bits of these integers, i.e. it may become log(N) for much smaller values.
For this reason, I’m not a bug fan of big O. Useful at job interviews. Probably useful for people working on extremely generic stuff like C++ STL which needs to scale both ways by many orders of magnitude, it’s not uncommon to have a vector with 1 element or 100M elements, the standard library must be good at both ends. But I don’t often solve such generic problems, I usually have some ideas how large my N-s are, and can act accordingly.
Sure, for files larger than 18.45 EB at the minimum, computing the new line count will take O(log N) time.
I'm all for nitpicking the issues with Big O notation, but this instance doesn't strike me as one.
The lesson to take away here is that O(log N) is REALLY REALLY small, for all sensible values of N. If an O(log N) algorithm is slow, it’s not because of the O(log N) part, it’s because of the hidden factors in O notation.
No it is not. "Big O" is a way of expressing asymptotic growth of some quantity, most often running time. But you need a machine model to even express running time. And sure, you can define your machine model to require O(log n) time for adding two numbers, but that's not particularly useful. In particular, the word RAM model—which is what is normally used—assumes a word size of θ(log n), and that basic operations on those words take a constant amount of time.
We use models because otherwise even the simplest things would be impossible to analyse. Or do you want to include a three-level cache hierarchy, micro-op decoding, the paging system, branch mispredictions and memory timings in your complexity analysis? Why is it always the "bUt AdDinG tWo nuMbERs tAkeS O(log n) tImE" hill that HN commenters are willing to die on?
First you need to settle on a model. Then you can express time, space, or some other quantity in that model using Big O notation.
Arguably, I've found that HN is one of the worst places to be precise about Big O notation. Everyone here seems to think they know what it is, only to lay bare their misconceptions about it a second later.
First of all, there's absolutely nothing that says that we have to assume that the number of lines is less than 2^64. Big O notation does not care what happens below that number, it only cares what happens when you go to infinity, which is a lot larger than 2^64. At that point you need arbitrarily sized integers, which are not fixed-space. Hence O(log n). Word size or RAM model does not matter one whit.
This is how Big O notation works, it has nothing to do with practicality. As another example of a similar thing, consider Sorting Networks. The AKS sorting networks has the best "asymptotic depth" of any known sorting networks (O(n log n) depth, i think), but the hidden constants are so massive that there are no situations where the AKS networks are actually practical, for any reasonable input they are always beaten by other (asymptotically worse) networks. That doesn't mean that they "aren't O(n log n)" or whatever.
Second: if we were to actually write this algorithm in the way which is most conservative with space, we wouldn't use a 64-bit integer in the first place. If we only care about memory, we'd start the count with a single-byte datatype (i.e. unsigned char) to keep the count. The algorithm now uses 1 byte of space. When that overflows (i.e. reaches 256 lines), we would change to a bigger data type, i.e. unsigned short. The algorithm now takes 2 bytes of data.
When THAT overflows, we again promote the datatype to an unsigned int. Now we use 4 bytes of data. When that overflows, we finally promote it to a uint64_t. Now we're on 8 bytes of data. When that overflows... (etc.)
See how the memory usage is growing? And how it's growing exactly with the logarithm of the number of lines?
You're just wrong about this. Yes, you can absolutely say "assuming storage of all numbers is O(1), then the algorithm is O(1)" (essentially what OP said) but that is not a default assumption and it's not accurate to reality (since actually storing numbers requires O(log n) of space). And yes, branch prediction misses and cache misses are far more important factors here, but that doesn't change the fact that it uses O(log n) space. That's just simply true.
Also: we're not dying on a hill. OP made an entirely accurate and true statement, which was met with an avalanche of people criticizing him using incorrect arguments and misunderstandings of Big O notation. We're just saying that he was correct all along.
Arbitrarily large numbers, represented in binary, certainly take time that is linear in the length of their representation to add. Nobody's arguing with that. But using a model that assumes bitwise operations would be needlessly complex, so outside of some really specialized contexts, nobody does it.
> You're just wrong about this. Yes, you can absolutely say "assuming storage of all numbers is O(1), then the algorithm is O(1)" (essentially what OP said) but that is not a default assumption
I don't know the basis on which you argue that the word RAM model isn't the default assumption, but I feel confident in claiming that outside of some niche communities, it most certainly is.
> and it's not accurate to reality (since actually storing numbers requires O(log n) of space).
That's why it's called a model. We use models because reality is too messy to work with, and we usually don't care about the exact details. Just ask a physicist, they make simplifying assumptions all the time because without them, they'd be getting absolutely nowhere. It's pretty much the same in algorithmics.
The issue was about space complexity, not time complexity. Very large numbers take more space than smaller numbers.
Does it make sense when expressed compactly like that?
I thought that big-O does not actually consider the underlying hardware. How can you specify the CPU model / instruction conditions and still use the semantics of big-O analysis? The mathematical definition of O(_) does not involve hardware, it implies that there exists some coefficients and a constant. I get the point you are making but is big-O an appropriate way of measuring in this context?
If you want to talk about something like supporting arbitrarily large integers or defining the memory of the physical computer we're using, we need to be explicit about that, because it changes a lot of things. After all, a finite-memory computer is not Turing complete, and everything it can compute can be computed in constant time. :)
When say that comparison sort is O(n * lg n), we're assuming a computational model where comparing two integers is a constant time operation. That's true in most physical computers when we use the CPU's comparison instructions. It's also true in any theoretical model where we define a (finite) largest possible integer. So it works well enough.
I’m pointing out that, for example, comparison sort is only O(n * lg n) in theory if the theoretical model of computation we’re dealing with can compare any two items in constant time. Comparing arbitrarily large integers can not be done in constant time in theory, so if that’s the problem we’re discussing then the complexity will not be O(n * lg n). But again, we generally are implicitly concerned with models of computation where comparison is a constant time operation.
Yes, Big O notation measures asymptotic growth of functions. But what do those functions express? For them to express running times, you need to define a machine model. It's a model because it's not a real machine – we don't stroll into a computer shop and pick up some computer and use that as our reference. Instead, we define an abstract machine model. Often, that's the word RAM model which assumes that for an input size n, the machine can process words of size θ(log n) in constant time. This is reasonable close to how our actual computers work: 64-bit words, and we can do stuff with them in constant time. But it's not perfect, because the RAM model doesn't have any caches (let alone multiple levels of cache hierarchy), doesn't have disks, doesn't model the paging system, etc. If any of these things are of significant concern to you, pick a different model (e.g. an external-memory model to consider memory hierarchy, or look into cache-oblivious algorithms, ...).
But talking about asymptotic running time without specifying which model you're using is absolutely pointless.
I don't have formal CS education. I'm less interested in math and more oriented towards practical aspects. Infinity doesn't exist in practice (except these IEEE float magic values), but we can reason about shape of functions, e.g. "count of executed machine instructions as a function of N".
Not the same thing as CS-defined big O, but close.
Can be useful to reason about complexity of algorithms. More precisely, about complexity of their implementations.
Usually useless to reason about time: a machine instruction takes variable time to run on a modern CPU, the same instruction can take less than 1 cycle or much more than 1M cycles.
No. Incrementing a 128-bit counter is also O(1) time on amd64/ARM8/... Only when your numbers are arbitrarily large would that become true.
But that's why we use models that make clear what our assumptions are. In particular, the word RAM model (which is normally used to talk about sequential algorithms) assumes a word size of θ(log n), and that some basic operations on words (such as arithmetic) are possible in constant time.
Models obviously have flaws (otherwise they wouldn't be models), but practitioners often make the mistake of overthinking them. How big, exactly, is your file that a constant number of 64-bit integers doesn't suffice? You can easily increment a 512-bit counter in constant time on an 8-bit Arduino. It doesn't matter whether I have to check one hypothetical 512-bit counter, 8 64-bit counters or 64 8-bit counters, it's a constant number and therefore O(1).
Whenever this "Big O is useful for job interviews but not much else" sentiment comes up on HN, it's usually driven by a misunderstanding of what Big O and asymptotic analysis actually mean.
Technically you’re right, I should have put “arbitrarily large” there, but assumed it’s obvious from the context.
> but practitioners often make the mistake of overthinking them
As an experienced practitioner, I happen to know how extremely inaccurate are these models.
When I was inexperienced, the models helped more then they do now, after I have leared about NUMA, cache hierarchy, micro-ops fusion, MMU, prefetcher, DMA IO, interrupts, and more of these hairy implementation details. Before I knew all these things, heuristics based on asymptotic analysis were generally helpful. Now they’re less so, I sometimes deliberately pick an asymptotically slower algorithm or data structure because it’s faster in reality.
I guess this is where I pitch what we call "Algorithm Engineering", which considers both theoretical guarantees and practical performance, including properties of real-world hardware such as the ones you mentioned :) https://en.wikipedia.org/wiki/Algorithm_engineering
Are you sure about that? The most recent source of that article is from 2000. Half of the stuff I have mentioned didn’t exist back then.
As far as I’m aware, there’re very few people in academia who study modern hardware while focused on software performance. I only know about Agner Fog, from technical university of Denmark.
Agner Fog does some cool stuff (and I've used his RNG libraries many times) but as you said he studies modern hardware with a focus on software performance. Algorithm engineering is about finding algorithms that have nice theoretical guarantees (e.g., proving correctness or that there aren't any inputs which would result in substantially worse running time) while keeping in mind that the result would be quite useless if it couldn't be implemented efficiently on the machines we have, so the goal is to come up with algorithms that are also really fast in practice.
- Read byte from input
- Write byte to output
- When byte is "\n": write the counter, write a space, increase the counter by one
Most likely, whatever logic you use to turn the counter into a string will already use more memory than your actual loop.
Technically it's still not O(1), as you will need one more bit for your counter variable every time you double the number of lines; realistically, you will most likely be using an integer anyway, so you can consider it "constant".
Can you elaborate on the O(1) solution?
I guess I'm misinterpreting the task, because I'm having trouble seeing how it's possible to access each line in a file without it being O(n) at a minimum. I understand O(log n) to mean that not every data point is "accessed" due to the way the data structure is organized, lines of a file in this case, but to me it appears that not every line would be prepended with a number. My instinct is that there's a lot going on with new lines/line breaks in files that I'm not aware of, so I'd really appreciate any reading material on the subject to help me better understand this. Thanks!
(My mistake obviously, the OP was very clear.)
sed = 1.txt
Use small buffer. No need to fit all the data into memory.People underestimate how much memory modern computers can actually support if you max them out.
Of course with some planning you can get an AMD system with ECC support on the mainboard; ECC RAM is about the same price as consumer RAM.
0: https://www.cs.toronto.edu/~bianca/papers/sigmetrics09.pdf
1: https://www.cs.virginia.edu/~gurumurthi/papers/asplos15.pdf
The problem is that both the consequences and debugging time of a single random bitflip are potentially unbounded. That time that GMail lost 10% of their accounts and had to restore from tape was because of a single bitflip (due to a software bug, not a cosmic ray). Google Search lost months of engineering time in the early days from cosmic rays, back when they really could've used those engineers on other stuff.
The odds may be low, but they do happen when you have multiple computers, and the consequences are high. Not really odds I'd want to take, when ECC RAM isn't that much more expensive than non-ECC RAM.
However, we have a large number of other processes (testing, type systems, formal verification, code reviews, release processes, etc.) to protect against software bugs. There is no protection against cosmic rays. You don't want to be in a situation where all of the defect-mitigation work that the last 50 years of computer science has accomplished is rendered useless by a random freak occurrence.
(The bug in question was actually in a migration script, and made it into production because people thought that migration scripts were one-off throwaways that didn't need the same amount of testing, code review, verification, and general carefulness that the production code does. Lesson learned. The postmortem for it actually had the lesson of "Treat your migration code as permanent, and apply all the same standards of maintainability and reliability of it that you do to production code.")
Also, hitting by a beam of cosmic ray is not the only way that the bits in RAM can be flipped, dynamic RAM has inherent instabilities like row hammering, or can fail early due to manufacturing defects.
Throw in the likely hood of the flipped bit being consequential and the odds look even better. To potentially do major damage the flipped bit would have to be in an area of memory of something being executed and it has to flip after being read and before being executed, which probably ads at least 2 more orders of magnitude. Even for general bitrot of data has to be in memory and get flipped between reads and writes. Those odds are vanishingly small compared to programmer error.
If you're not buying for your particular specifications, you're doing it wrong anyway... There are plenty of workloads tolerant to ECC errors (e.g. just about any kind of simulation).
They wanted some answer involving using spark engineers you'd effectively pay for at $200/hr for many weeks. I didn't get the job.
OK, but seriously speaking, if the upper bound is beyond 786gb RAM or whatever the current max is, I might want to use dask distributed. https://distributed.dask.org/en/latest/
edit: wrote mb instead of gb
Did you mean gb? In any case you can actually do at least 24TB these days: https://aws.amazon.com/blogs/aws/ec2-high-memory-update-new-...
But then they dont usually like when you demonstrate anything but the answer they want to a highly contrived problem.
Usually when I tell that story, I get a lot of objections about how that solution won't scale and they must not have really had big data from people who are, truth be told, used to working with data on a fraction of the scale that this company did.
That said, it's not a turnkey solution. This company also was more meticulous about data engineering than others, and that certainly had its own cost.
If the only single-machine option you consider is Pandas, which doesn't do streaming well and is built on a platform that makes ad-hoc multiprocessing a chore, you'll hit the ceiling a lot faster than if you had done it in Java, which might in turn be hard to push as far as something like C# (largely comparable to Java, but some platform features make it easier to be frugal with memory and mind your cache lines) or, dare I say it, something native like ocaml or C++.
Alternatively, if you start right off with Spark, you won't be able to push even one node as far as if you hadn't, because Spark is designed from the ground up for running on a cluster, and therefore has a tendency to manage memory the same way a 22-year-old professional basketball player handles money. It makes scale-out something of a self-fulfilling prophecy.
Also, as someone who was doing distributed data processing pipelines well before Hadoop and friends came along, I'm not sure I can swallow "big data" being one and the same as "handling data that is too big to run on one computer." Big data sort of implies a certain culture of handling data at that scale, too.
Because of that, I tend to think of "big data" as describing a culture as much as it describes anything practical. It's a set (not the only set) of technologies for procesing data on multiple machines. Whether you actually need multiple machines to do the job seems to be less relevant than the marketing team at IBM (to pick an easy punching back) would have us believe.
That's because a reasonably sized machine from today is much larger than one from five years ago. And an unreasonably large machine today is also larger but yet more achievable.
A basic dual Epyc system can have 128 cores, and 2TB of ram. Someone mentioned 24 TB of ram, which is probably not a two socket system.
You can do a lot with 2TB of ram.
But I think it's quite safe to say that it's not often because you need to process so much data, but rather that your experiment is a fire hose of data, and you're not sure what you want to keep, and what you can summarize - until after you've looked at the data.
And there might be a reason to keep an archive of the raw data as well.
Another common use case would be seismic data from geological/oil surveys.
But "human generated" data, where you're doing some kind of precise, high value recording, like click streams, card transactions etc might be "dense", but usually quite small compared to such "real world sampling".
Today most commodity servers (which you can buy out of pocket without lengthy preorder) can accomodate from 1.5 to 3 TiB of memory.
That, and the idea of a polyphasic merge sort isn’t taught in ye olde boot camps, I guess.
That's a self-imposed constraint, not a problem constraint.
Computational resources are dirt cheap nowadays. Everyone can get free time in global scale clusters. The only reason anykne is stuck with a laptop to run number/data crunching tasks is because they want to.
Bonus: you get to keep the hardware.
You can have many programs, and they can pass messages around, so you can do jobs that won't fit in one program. It's like coding for a cluster of Arduinos.
I've had to pack data into integers with shift operations. Come up with an efficient maze solving algorithm that only needs 2 bits per cell. Divide tasks into multi-step pipelines. Monitor memory usage and divide the data into smaller chunks that will fit. Plus it's soft real time and has to run properly under overload conditions when it can't get enough CPU time.
It's kind of neat to see autonomous characters running around the virtual world at running speed. This was believed to be impossible, but I got it working out of sheer stubbornness. It was far too much work.
(Second Life has a built-in pathfinding system. It's too buggy to use, and they refuse to fix it. This is a workaround.) (Why such tiny programs? Because they are not only persistent, for years if necessary, but are copied from one machine to another, state and all, as objects move around.)
It pushed me to work on minimum-memory maze solving. It costs a lot to test a cell (this involves ray-casting in the simulated world) so something like A*, which examines most of the cells, is out. An algorithm in Wikipedia turned out not to work; some anon had snuck in a reference to an obscure paper of their own. Had to fix that. What I'm doing is "head for the goal, when you hit something, follow the wall, if it will get closer, head for the goal again". Wall following is both ways simultaneously, so you don't take too long on the long path when a short path is available. After getting a path, the path is tightened up to take out the jaggies. This is not optimal but is usually reasonably good.
The rest of it is more or less routine, and a pain to break into sections. The programs have to communicate with JSON, over a weak IPC system with bad scheduling.
I always liked the 3D "metaverse" concept. Second Life, which has about 30,000 to 50,000 users on line at any one time, comparable to GTA Online, is the biggest virtual world around. Everybody bigger is sharded, but all SL users are in one world. The technology needs a refresh, but every competitor who's tried to build a big virtual world has been unable to get many users. So, you're stuck with outdated tech if you want to get something used in a virtual world. I think they're still in 32 bit mode on the servers, even.
My lesson learnt was that I just needed to reduce how much data I was working with first, instead of trying to stuff a multi-gb file into memory!
Use case: CPU is running a simulation and constantly spitting out data, eventually RAM is not enough to hold said data. Solution: every N simulation steps, store data in RAM onto disk, and then continue simulating. N must be chosen judicially so as to balance time cost of writing to disk (don't want to do it too often).
I figure this is what is referred to as "chunking" in the article? Why not list some packages that can help one chunk?
Overall opinions on this method? Could it be done better?
The benefit of HDF5 is that it allows for very easy slicing of data. So 'chunking' where I load e.g. 10% of one dataset at once is very simple with HDF5.
One classic example is sorting. When your data is too big for RAM, quicksort's performance is horrible and mergesort's performance is fine (even if your data is on magnetic tape).
Another classic example is taking a situation where you build a hash (or dictionary) and using sorted lists instead.
Let's say your task is to take some text and put <b></b> tags around a word but only the first time it occurs. The obvious solution is to scan through the text, building a hash as you encounter words so you can check if you've seen it before. Great until your hash doesn't fit in RAM.
The sort-based solution is to scan through the input file, break it into words, and output { word, byte_offset } tuples into a file. Then sort that file by word using a stable sort (or both fields as sort key). Now that all occurrences of each word are grouped together, make a pass through the sorted data and flag the first occurrence of each word, generating { byte_offset, word, is_first_occurrence } tuples. Then sort that by byte_offset. Finally you can make another pass through your input text and basically merge it with your sorted temp file, and check the is_first_occurrence flags. All of this uses O(1) RAM.
I believe this is basically what databases do with merge joins, but the point is you can apply this general type of thinking to your own programs as well.
The Count-Min Sketch is a data structure consisting of a fixed array of counters.
This is a short and easily accessible paper, which isn't very heavy on math or obscure notation or concepts, so I would recommend it to anyone with even a cursory interest in the subject.
Here's an implementation in Python: https://github.com/barrust/count-min-sketch/blob/master/pyth...
On other problems chunking doesn't work at all and just mmaping or dedicating giant amounts of swap are better strategies. It depends on the problem at hand
Also, SQLite is a first class citizen in this space. Most if not all languages can easily load data to it, virtually any language used for analysis can easily read from it, and it's file-based so there's no reason to spin up a server. Finding 100GB on disk is much easier than 100GB in ram.
They insisted that shouldn't be, because I was doing it on my laptop and they were using a high performance computing cluster. They of course wanted to know how my implementation could be so much faster despite running on only a single machine. I didn't have the heart to suggest that maybe it was because, not despite.
Ironically, I also got the implementation done in a lot fewer person-hours. I just did a straight code-up of the algorithm in the paper, where they had to do a bunch of extra work to figure out how to adapt it to scale-out.
This isn't to say that big data doesn't happen. Just that it's a bit like sex in high school: People talk about it a lot more than they actually have it, perhaps because everyone's afraid their friends will find out they don't have it.
Of course the next problem is "my data doesn't fit on disk", but with xzcat on the command line and lzma.open in python you can work transparently with compressed files.
I usually use LZ4 compression for this purpose, because it's ridiculously fast (700 MB/s per core compression and 5 GB/s decompression). Sure, compression ratios aren't as good, but usually good enough.
Yet, one realizes how well written these basic tools are only after having some bruises with fancier tools.
Since gzipped tar archives contain sequential data, solving it ended up being trivial. With Go, I was able to string together the file download, tar extraction, and forwarding the individual files on, all within the context of a single stream.
curl ... | tar xv | ...
in shell, or any programming language with a good streaming library?It blew my mind when I realized I could add / remove RAM by turning off the instance, dragging a slider, then turning it back on.
I also find h5py really useful for creating massive numpy arrays on disk that are too large to fit in memory. I used it to precompute CNN features for a video classification model (much faster than computing on each gradient descent pass) and it makes it easy to read/write parts of a numpy array when the entire numpy array is too big to fit in memory.
mmaping is a form of uncontrolled chunking, where chunk size and location is determined by operating system's filesytem caching policy, readahead heuristics, and the like. As a result, it can sometimes have much worse performance than an explicit chunking strategy, especially if you just treat it as magic.
Or to put it another way, mmaping is helpful but you still need to understand why it might help and how to use it.
Can you expand on this?
I haven’t heard this before.
Edit: Found a comprehensive discussion on this https://stackoverflow.com/questions/45972/mmap-vs-reading-bl...
I just wish more technical documentation had a human element of “what we intended” to it.
Nothing fits: disk to RAM, RAM to L3, L3 to L2, L2 to L1, L1 to registers. We are just lucky that many programs have spatial and temporal locality.
And L3 is already at 256MB (+32MB L2) in the new AMD Threadripper CPUs.
Seriously, this should be the top of the list and default solution when you need to process large file that does not fit in memory.
The exceptions would be when you can read the file once (then stream might be easier but not cheaper), when you don't want to rely on FS cage (you need your own) or when you need easy interoperability between different OS-es.
It is really magical thing to be able to run just a single function and have your entire terabyte file suddenly present itself as addressable continuous memory space ripe for random access.
Regular random access file I/O seems silly to me, like trying to work with the file through a key hole.
Have you ever seen how ships are built inside a bottle? That's exactly the picture I have in my mind...
Cold store looses big in only scenarios where you have sequential and completely random seek patterns, and there are lots of way to optimise that in read and write heavy workloads. This was the art of running performant multi-terabyte DBs in a world prior to SSDs and ramdisks.
For very big index walks, you want to have data to be more flat, and seeks to be sorted, so there is higher chance that needed records would be accessed without page eviction. Modern DBs, I believe, do something like that internally.
And for write heavy loads, there is no alternative to revising "data structures 101." You can reduce the disk load by may times over with a properly picked tree, graph, or log structure.
This way you can also use clever versioning algorithms to predict the performance (no read or write peaks) and to lower the space consumption.
All using Python's standard library, here's a quick post I wrote about this. I last used it a week ago and got from 950 MB RAM usage down to about 45 megs. https://kokes.github.io/blog/2018/11/25/merging-streams-pyth...
You can easily retain the full version history in a log structure, but you need fast random access (option of parallel "real" access would be best) on a flash drive -- PCIe SSDs for instance.
Basically in order to balance read and write performance only a fraction of each database page (with changed records) needs to be written sequentially in batches to the end of a file.
Each revision is indexed under a RevisionRootPage and these are indexed under the UberPage with keyed tries.
This opens up a lot of opportunities for analysing data and its history.
The reduction in memory has allowed us to do parallel processing, for a pretty significant speed-up!
[0](github.com/LaurentRDC/npstreams)
If you sign up by month it will be much cheaper too.
[1] https://www explorablelabs.com
Disappointingly, I ended up writing my data to a text file, sorting with unix sort (w/ some tuning on --parallel and --buffer-size) and reading it back.
One of my biggest beefs with Unix is that so much powerful functionality is locked up in command line utilities and not accessible as libraries.
Whole companies have been built around the problem. e.g. syncsort.com
I hope some interesting techniques don't get locked up and lost in proprietry mainframe software.
Uhm... Isn't it called a database?
Jokes aside, one of the core use cases of things like databases is random access to a part of data that's too large to fit in RAM.
In-memory databases are a thing, but that's a very specialized use case.
I realize that relational databases are not the right box to fit certain kinds of data into, but you have to put your data somewhere that allows it to be efficiently manipulated. What is that if not a “data base”?
Next thing you know, grandma is typing:
SELECT * FROM GOOGLE.COM WHERE TITLE LIKE "%CAT PICTURES%";
That's just an analogy, of course. But perhaps you can imagine your own list of reasons why MySQL didn't replace Apache.
The answer to your question is that we could not possibly use any off the shelf software to provide the interfaces we needed (we were writing our own). By the time we had bytes that we needed to store somewhere, our data was about a GB of complicated structure that we accessed as memory-mapped files. (Look at https://capnproto.org/ for a rough analogy to the kind of access we needed.)
And what I didn't say was that the DBA was literally recommending MySQL BLOBs, storing a bit more than 0.5 MB (512 * 512 * 2 bytes of image data) in each row, having a thousand rows or more per CT scan. The performance of that would have been absolute crap. It made literally zero sense.
However, if you are simply "pulling blobs up" of compressed (or uncompressed) volumetric data... well yeah. Don't put that in a sql database.
Also, another hint: process all inserts in the bulk in a single transaction. I was amazed at the speed that SQLite ingests data! Speed demon! (Don’t forget PRAGMA synchronous OFF too)
And yes you will need to rewrite a small amount of your code but you're doing so here as well.
At least you will be able to scale out to much larger volumes in a consistent way.
And 90% of all tasks are batch orientated.
yeah, that's what I thought.
(assuming the data is organized in files to begin with)