Puzzleshot #003 - Multiple large files, similar questions
Once again, we’re back with trying to validate aspects of a data delivery without some of the more modern tools. This time, the delivery is over 500 million rows spread across 500 gzipped files (1M rows each) and you need to validate the aggregate per region counts of column 4 across the whole dataset.
In Puzzleshot-002, we learned how to use awk’s associative array functionality to cycle through the data once and do a count, avoiding traditional methods such as sort and uniq due to their memory requirements. Unfortunately serially processing the dataset takes too long even on a modern system with 32 cores and 512GB RAM.
How would you speed up this operation?
Reveal solution
For Puzzleshot-003 we suddenly needed to process over 500 million rows on a reasonably modern system, but were hit with the serial nature of our last process. Given that our system has 32 cores and 512GB RAM, running a single awk instance to process everything makes the process much slower than it should be.
Enter parallel - This utility allows us to pass a list of data along with instructions of what to execute. It then handles all the input / output based on the number of instances we wish to run.
Let’s look at the first part - parallelizing the task. parallel allows access to the data passed in via the {} substitution. Our previous solution (single-threaded) used the following:
zcat datafile.csv.gz | awk -F',' '{counts[$4]++} END {for (region in counts) print region, counts[region]}'
Now we can do the following in single-threaded form for our new data like this (note the wildcard * for filename):
zcat datafile*.csv.gz | awk -F',' '{counts[$4]++} END {for (region in counts) print region, counts[region]}'
This works exactly as before but across all files in a single-threaded fashion. To run this in parallel, we can do the following:
ls datafiles*.csv.gz | parallel -j 24 "zcat {} | awk -F',' '{count[\$4]++}END{for (i in count) print i\":\"count[i]}'"
This will get the list of files matching the wildcard and then prints out the counts from each region in each file as before. In this case, parallel will run up to 24 tasks at a time - this can be tweaked based on your use case / number of cores / server utilization, etc. A couple things to note - because we’re doing this as a one-liner, we need to escape the dollar-sign on the $4, otherwise parallel or the shell may substitute it for us. We also need to escape the double-quotes around the : in the awk’s print statement to keep the line together.
The problem with this solution is the data comes back as each region’s counts per file (ie, 2000 rows of output = 4 region counts * 500 files). We need to add a little more to print out the calculation. As always, there are multiple ways to do this, but we can actually just pipe all of our output into another awk instance like follows:
ls datafiles*.csv.gz | parallel -j 24 "zcat {} | awk -F',' '{count[\$4]++}END{for (i in count) print i\":\"count[i]}'" | awk -F':' '{totals[$1]+=$2} END{for (i in totals) print i": "totals[i]}'
In this case, all the output from the parallel tasks are then parsed and split on the : character, giving us the region and counts. Note we don’t have to escape the entries on the 2nd awk command as it’s a single-command as before. We total it up like we did in the single-threaded version and reduce all of them to a final count. This version will be drastically faster (wall clock time) on bigger hardware as far more of the system is being utilized.
There’s a lot of power that parallel and normal unix tools give you - this is basically a simple map-reduce system in the vein of Hadoop or Spark without anywhere near the overhead of setup / coding / etc.