Ask HN: Sorting massive text files?
time cat FILE | wc -l 2608560847
real 11m18.148s user 1m35.667s sys 1m33.820s [root@server src]#
Any suggestions on how I can go about getting unique records from this type of file?
time cat FILE | wc -l 2608560847
real 11m18.148s user 1m35.667s sys 1m33.820s [root@server src]#
Any suggestions on how I can go about getting unique records from this type of file?
http://vkundeti.blogspot.com/2008/03/tech-algorithmic-detail...
Thus, sort -u <filename> is your go-to for simple jobs. (Note that you'll need to have enough extra disk space to hold all of the temporary files.)
If you need to do something more sophisticated (e.g. joining lines from a web server log into sessions, then sorting sessions), you can still use divide-and-conquer, but you have to be smarter. Divide the file into N parts based on some logic (i.e. divide into files based on session ID), then sort the lines in each individually, then merge the results back together.
This is what map/reduce frameworks are made to do, of course, but something like Hadoop may be overkill unless you plan to do this type of thing often.
Here's a blog post I wrote about it:
- a target to split the input into chunks
- targets to sort each of the split chunks
- a target to merge the sorted chunks using "sort -m"
Then, if you are using GNU make, it's just a matter of using the -j flag to control how many chunks are sorted in parallel.1) Do sort -u --buffer-size=60GB FILE. You'll sort it all in memory, and that'll be a great speedup.
It's easier to scale up than scale out if you have the money, and your dataset fits in memory, so don't bother with Hadoop for something as simple as that. What do you want to do after you get the uniques?
You could use the same approach to simply take the unique tokens of the output of the final "count".
May be overkill if you you can read the file in less than an hour, but this approach (divide and conquer) may be a good inspiration.
sort -u <filename>
would be faster than all that pipingIt may also help to use compression with the temporary files (something like "--compress-program lzop"), but I've never tried that.
Did you really think sort(1) just bails if you pipe it more data than can be contained in memory?
lzop -9v $file
lzop -dc $file.lzo | sort ...
You shouldn't use cat(1) if you don't need to. You're removing the ability of wc(1) to do faster things than read(2) from a pipe, e.g. mmap(2) the file, which is probably what cat is doing.Similarly don't use uniq(1) if sort's -u option will suffice as you're forcing 35GB of data to be written down a pipe to uniq(1) when it's quite possible that `sort -u' would produce 100MB or so.
split -d -C 1G FILE split.FILE # separate into 1G files on line boundaries
# then, on separate cores
sort -u split.FILE.00 -o sort.FILE.00
sort -u split.FILE.01 -o sort.FILE.01
...
sort -u split.FILE.N -o sort.FILE.N
# then, on one core
sort -m -u sort.FILE.*
You should be aware that you may not get what you expect from sort unless you set LC_ALL="C". You should also pass "-S" to sort to set the main memory size, probably to a value like "1G". Of course, at some point when you're doing a lot of distributed processing you'll just want to use Hadoop, as other posters have noted.http://en.wikipedia.org/wiki/External_sorting
The reason to split is not to replicate mergesort, but rather to get around single threading in GNU sort. Of course, there are better options if this is a regular task, but GNU sort just happens to be really common and easy if it's a one-off.
I do wish I had a better split (since split seems to read the whole file) and a better sort (that was parallel and reasonably common) though.
Both can be rewritten as "sort FILE | uniq" and "time wc -l FILE".
And given the data size, both non-cat uses would likely be faster as well.
Of course, if you want to learn hadoop then go for it. But it's probably more practical just to let it run :-)
How much memory do you have available?
The Bloom filter is a good way to go. Here is one possibility:
If you have 2 billion unique lines here's how much memory you need for a Bloom filter:
Error rate Memory
------------ --------
1 / thousand 3.35 GB
1 / million 6.70 GB
1 / billion 10.05 GB
1 / 10 billion 11.16 GB
1 / 100 billion 12.28 GB
If you have 12 GB available you can reduce the chance of missing a unique string down to almost zero.
http://en.wikipedia.org/wiki/Cuckoo_hashing http://www.ru.is/faculty/ulfar/CuckooHash.pdf
You can get your hash table utilization to almost 100%. If you're willing to accept a 64-bit hash function and you have 2 billion unique strings you'll need just over 16GB to do it all in one pass.
A 64-bit hash is is just barely enough for 2 billion items. You might get a handful of collisions. Add more bits to reduce the probability. 128-bits would double the required size to 32GB but you're unlikely to ever see a collision in your life time.
If 16GB or 32GB is too much you could do it in multiple passes. Make N passes through the text file and adjust the range of hash values you accept to test only 1/N each time.
8 passes through the original file with a generous 128-bit hash would take only 4GB of RAM.
I have the feeling this is more work than you would want to put into the problem, though.
Your 'cat' is not needed by the way, sort takes a filename as argument. And sort can '-u'!
It would be interesting if sort would fail on this (why?) or how long it would take.
edit: Why on earth are you doing this as root?
Also, what do you want to do with the unique records? That might effect what initial processing method is best for your goal.
Maybe a kind of divide and conquer could work? Split into several files, do the sort | uniq on each of them. Then merge them, checking for duplicates on the way. I think merging should be almost as fast as line counting, at least linear in the size of the two files.
Edit: I guess it would be slower than counting, because presumably it would write to a new file (the merged file). But still, it should be linear.
The only thing is that you have to split it up and merge after sorting (for which unix sort was ok enough).
Not sure why I got that result, but even with increased buffer size for unix sort it didnt much differ. I also didn't run the splitted sorts in parallel, which would of course have been a good idea.
You may have had a data set that tickled something that plays to Timsort's advantage; Timsort was basically designed to encounter that case as often as possible on real data.
then if you have multiple cores on your machine, you can run multiple instantiations of the same script on those smaller files and aggregate the results later.
if you're more ambitious, you could look into using a lightweight MapReduce framework
first split the file into chunks that will comfortably fit in ram, depending on the encoding of the file and the language's default encoding, allow maybe ~2x blow up in memory size, though experimenting is more accurate than guessing.
Then sort each of these files in place.
then have a program then opens all these smaller sorted files, and does an incremental line by line merge that compares them over all these files, with the case of two lines being equal to drop one, and write the lines that are unique to the result file, and then tada!
i actually first dealt with this problem on a programming interview, and I quite liked how instructive an example of out of core programming it was.
If that's what you're after, you don't need to sort at all, you can pick out your uniques with a 3 line script in your favourite language, split up the file first, if necessary.
I'd guess it's a gigantic spam list.
...and this is a problem most databases have addressed.
That sounds at least halfway feasible to me - you assume 1/3 duplicates, and 20 bytes per line, you'd get ~8GB worth of hashtable entries (i.e., if you want the hashtable 70% filled to limit the amount of collisions, you'd need 12GB of virtual memory to back the hashtable, but only rarely access it since you're using the Bloom filter).
(To the person who downvoted it: can you say why you don't like the idea?)
first thing I'd try is probably pseudo-this:
while((entry=readDataPoint())!=null){
sha1e = sha1(entry);
//database has log N index on sha1e
if(database.contains(sha1e)) continue;
//database object figures out insert syntax
database.insert({sha1:sha1e,data:entry});
}The above can be written: awk -F(dataseperator) "{print $(number of seperation to print)}" FILE | sort -g | uniq | sort -g
...| sort -g | uniq | sort -g
The second sort is doing nothing. And sort can do uniq, saving time and I/O. $ (seq 10; seq 5) | shuf | sort -gu | fmt
1 2 3 4 5 6 7 8 9 10
$I know for a fact that they do do it on smaller batches, though. It takes a lot of room, as you can imagine!