Running Awk in parallel to process 256M records
ketancmaheshwari.github.io
ketancmaheshwari.github.io
But, I started being somewhat confused by something:
> Fortunately, I had access to a large-memory (24 T) SGI system with 512-core Intel Xeon (2.5GHz) CPUs. All the IO is memory (/dev/shm) bound ie. the data is read from and written to /dev/shm.
> The total data size is 329GB.
At first glance, that's an awful lot of hardware for a ... decently sized but not awfully large dataset. We're dealing with datasets that size at 32G or 64G of RAM, just a wee bit less.
The article presents a lot more AWK knowledge than I have. I'm impressed by that. I acknowledge that.
But I'd probably put all of that into a postgres instance, compute indexes and rely on automated query optimization and parallelization from there. Maybe tinker with pgstorm to offload huge index operations to a GPU. A lot of the shown scripting would be done by postgres, the parallelization is done automatically based on indexes, while eliminating the string serializations.
I do agree with the underlying sentiment of "We don't need hadoop". I'm impressed that AWK goes so far. I'd still recommend postgres in this case as a first solution. Maybe I just work with too many silly people at the moment.
There'd be much less setup overhead
Ok. Maybe those are enterprise concerns: Sqlite doesn't scale regarding multiple users reading and writing. Of course it's a read only dataset, but do you know the bouquet of views and derived tables data scientists create around a read-only dataset? Hah. Oh and of course these are not critical, but if they get lost, shit hits the fan because it takes multiple weeks to rebuild them.
I've been in that swamp enough times to just install postgres and stop caring. Takes me 2 more hours now, but avoids weeks of discussions in the future.
apt-get install postgres
create user
create databaseThose are massive, massive machines.
I know you said "automated parallelization" but... in newer versions of Postgres, what does it take to trigger some of the automatic query parallelization?
Table scans, index scans (b-tree & bitmap only I believe), joins, and aggregations. There may be some limitations with joins, such as hash joins duplicating hashes across processes or merges requiring separate sorts.
Although I agree that for this problem Postgres looks like a better option, I wouldn't discard awk or other Unix tools. It is way faster than people think it is, easy to use and solves pretty well quite a lot of use cases, such as pre-processing records before database ingestion, aggregations, some simple queries on temporary data...
Just wanted to say that the resources were not needed. I simply had access to them and were not being used much at the time so thought why not try them.
I don't think it's really displayed here because everything is run through swift, but it can also compose much better with other tools, a database frequently won't have cli tools that integrate well (not sure about postgres). There are 322 source files being globbed for instance, move that to a make file and you can automatically re-run without wasting time on data already processed. If it was in a database you'd have to track source files and manage changes somehow.
> At first glance, that's an awful lot of hardware for a ... decently sized but not awfully large dataset. We're dealing with datasets that size at 32G or 64G of RAM, just a wee bit less.
Note that it doesn't require that much memory, it's just using it to boost performance because it's available, this could have been developed on a dual core 4GB laptop outputting to spinning rust and all that would change is the directory it's output to (and running time), but the same machine may choke on database queries over that much data.
@ketanmaheshwari, which supercomputer is this?
This is an SGI UV300 system similar to this one: https://www.hpcwire.com/2016/05/11/tgac-installs-largest-sgi...
I believe databases like postgres etc. would probably run on this system. I did not choose that route because I wanted to see how far I can go with Awk.
First, I would need to figure the right schema and populate the database.
Second, it would need some creative SQL acrobatics that I would probably not be comfortable with.
Third, it would probably be hard to perform quick tests at a small scale that I can perform easily with text files.
Fourth, the solution would probably be hard to port elsewhere where postgres is not available. Most Unix systems have awk available.
Fifth, programming the postgres db in a higher level language would require connector api libs which would be additional effort.
Note that this is not a production work -- it started simply as a hack to see how far I can go without getting into serious rabbitholes and giving up. Surprisingly, with Awk I went really far and never fell into any rabbithole so to speak.
[1] https://github.com/lanl/MPI-Bash [2] http://hpc.github.io/libcircle/
The results were indeed embarrassingly parallel.
I am a fan of some languages that seem better equipped to utilise our modern many-core machines and I'd still write a longer-living system with better guarantees in those languages -- but many people ignore shell goodness at their own peril.
I've been developing a simulation software that does Monte Carlo over realizations of a simulated universe, and xargs and (later) parallel were super useful for parts of the workload. All the parallel job instances run the same simulation code, but for a different simulated universe, all controlled by a random number seed, so you can generate an ensemble of simulations with basically:
head -NumberOfSimsWanted seeds.txt | xargs simulation.pyI quite disliked the journey. The UNIX shell legacy shenanigans are real and can be a pain. But you can still bear a lot of fruit without going to the more arcane corners of it.
Switching to Opus saved a great deal of space on my phone's SD card.
Thanks for mentioning AWS Batch, I didn't know about it and will look it up.
I'm sure I could have done the parallelization itself in Rust given enough time, but honestly I found the `parallel` command to be so easy and resulting in so little overhead that it didn't really even seem worth it to spend the time on a language-native solution.
Obviously there are cases where parallelisation is very superfluous. F.ex. if you have 1000 records of something, splitting the work of processing those (and we're talking really light processing like JSON encoding) among 4-8 threads is actually detrimental to performance.
https://www.gnu.org/software/parallel/
*SAPANS: as simple as possible, and no simpler
Do note that awk and GNU Parallel are tools that use/contain C-like syntax, which is abhorrent to some factions. GNU Parallel is written in Perl.
Sometimes it improves performance, sometimes not. Sometimes it improves productivity, sometimes not. I am not sure of the global balance.
* The OS provides process isolation and management, as well as various IPCs to deal with their interconnection? Let's use a single process and create threads, and then fibres. More performance? yes, a bit; more headaches? yes, many.
* The display protocol/server provides lots of primitives for managing graphical interfaces, together with many benefits? Let's just use it to get a canvas and reimplement everything ourselves inside that canvas (and thereby turn a snappy interface into something sluggish when the display is not local).
* A GUI toolkit provides easy to use, multiplatform interface elements? Let's take a single program which only use a single canvas from the toolkit, and reimplement a whole GUI toolkit within (in one of the worst languages, over a program not made at all for the job, otherwise it wouldn't be fun).
* The OS provides management for multiples users and groups, together with 2 systems to manage resources access rights? Nah, let's run everything as one user, and then try to implement some separation/protection within. Advanced level: run everything not only as a single user, but inside a single program, and then try to reimplement some isolation within it.
> The OS provides process isolation and management, as well as various IPCs to deal with their interconnection? Let's use a single process and create threads, and then fibres. More performance? yes, a bit; more headaches? yes, many.
Many one-off CLI script programs don't have the time to deal with N process forks (where N = CPU threads). This overhead is proven. It's objectively much quicker to spawn N threads inside the same process especially if your separate tasks need to converge and/or exchange messages here and there.
And runtimes like Erlang's BEAM VM have unquestionable benefits (like preemptive green threads / fibers that exchange immutable messages), although they aren't well known for big raw processing muscle.
> The display protocol/server provides lots of primitives for managing graphical interfaces, together with many benefits? Let's just use it to get a canvas and reimplement everything ourselves inside that canvas (and thereby turn a snappy interface into something sluggish when the display is not local).
Which one is that? Linux's X11 is flawed in many regards, with the ability of each program to freely capture keystrokes being one of the most egregious. I can't blame people for trying to escape such hell and implement other platforms for the same thing -- without the drawbacks.
> A GUI toolkit provides easy to use, multiplatform interface elements? Let's take a single program which only use a single canvas from the toolkit, and reimplement a whole GUI toolkit within (in one of the worst languages, over a program not made at all for the job, otherwise it wouldn't be fun).
People wanted multi-platform desktop programs. F.ex. in my rather small country the only reason a software for managing shop stock and customer purchases, backorders, inventory etc. succeeded was because it was written in Java Swing; many shops used old Windows machines but no small amount of them also used Macs and some even used Ubuntu netbooks. That software would never sell back in 2007 if it wasn't OS-neutral.
Also, Qt has commercial licensing that requires paying to use. Many, me included, don't want to deal with that. The programming world is complex as it is and if I suddenly find myself having to pay royalties, personally, to Qt, years after I stopped working for the employer I developed a desktop app for, that would be catastrophic. So many just dodge such potential bullets.
> The OS provides management for multiples users and groups, together with 2 systems to manage resources access rights? Nah, let's run everything as one user, and then try to implement some separation/protection within. Advanced level: run everything not only as a single user, but inside a single program, and then try to reimplement some isolation within it.
That's quite fair, however people want much more isolation than the OS usually offers -- hence stuff like FreeBSD jails, LXC containers / Docker, Windows' Sandboxie program etc. It's not enough to just run, say, Chrome, as a different user if you suspect Google is scanning your machine.
---
Again, I agree with your premise but the examples aren't really accurate. A lot of the modern OS-es and and their technologies aren't as top-notch as we'd like (the Linux kernel and network infrastructure probably being some very rare exceptions).
Only if you intend on modifying Qt and distributing it. Otherwise, you're free to build commercial software using Qt under the LGPL license.
Unfortunately I no longer have a working example but as I recall it wasn't difficult to keep n number of scripts running and popping entries until the queue emptied.
But what you mention is very interesting and I'll have it in mind for the future. Thank you.
The xargs/parallel mode can also scale up to many more nodes, if you are willing to introduce ssh into your shell scripts & do a minimum of stdio piping.
I wrote a python mapreduce program named "batchman" a few years ago, which was aimed at running awk on hundreds of machines (& then being able to fire more things based on the results).
zgrep on gzipped logs + awk + a central scheduler is extremely fast, to the point of beating Splunk to lookup exact things like ("does this zid have another session which they kept logging in while the game was throwing out-of-sequence errors?").
I thought of utilising my 4 machines at home sometimes for such things and would hugely benefit from implementing this using standard tools.
My code is poorly documented and badly written, but it works at least the ssh_copy_id one I use enough to keep updated[1] (that ssh logins with passwords to place the ssh keys).
That's my script, the brains of it is basically an ssh connector.
Paramiko is as simple as
s = SSHClient()
s.connect(hostname=h[1], username=h[0], password=pwd)
(_in, _out, _err) = s.exec_command( cmd, bufsize=4096)
Originally, I had awk scripts to process ^A separated, I'd read the out stream back and formulate my next script over the logs to drill-down in a loop.[1] - https://github.com/t3rmin4t0r/batchman/blob/master/ssh_copy_...
Edit: Done.
I love these kind of “counterculture” uses of UNIX tools to solve hard problems. The boundary of “where to stop using xargs, awk, and grep and start using Python” is pretty blurry for a surprising number of tasks, especially if you’re willing to invest in mastering those tools!
!($1 in a) && FILENAME ~ /aminer/ { print }
This uses a regular expression. As regex is not actually needed in this case, you might be able to get better performance with something like this: !($1 in a) && index(FILENAME, "aminer") != 0 { print }I suspect if it was open source, it would probably be the most popular big-data storage and computing platform.
250M records isn't big for a lot of data platforms. I've seen plenty of solutions in traditional RDBMSs that would easily scale an order of magnitude larger than this, if not several.
Moving into anything analytically focused, like an OLAP engine, or just a columnstore in an RDBMS easily gives you another order of magnitude or few.
Querying things like these on such a small dataset should take seconds, not minutes.
The solution that took 9 minutes involved processing the abstract in each record. The abstracts are quite sizeable on some of these publications. Processing millions of them took time.
The solution that took 48 minutes involved a nested loop effectively reaching the iteration count for 216 years times 256M records which comes to about 55B.
Hope this clarifies things a bit but I am not claiming this to be the most optimized solution. I am sure there is scope for refinements -- this was my take on it.
Is the proper answer 'just learn it'? Are these tools one of these things (like musical instruments or painting) where the initial learning phase is tedious and frustrating, but the potential is basically limitless?
My tools of choice were awk, R, and Go (in that order). Sometimes I could calculate something within a few seconds with awk. But for various calculations, R proved to be a lot faster. At some point, I reached a problem where the simple R implementation I borrowed from Stack Overflow (which was supposed to be much faster than the other posted solutions) did not satisfy my expectations and I spend 4 hours writing an implementation in Go which was a magnitude faster (I think it was about 20 minutes vs. 20 seconds).
So my advice is to broaden your toolset. When you reach the point where a single execution of your awk program takes 48 minutes, it might be worth considering using another tool. However, that doesn't mean awk isn't a good tool, I still use it for simple things, as writing 2 lines in awk is much faster than writing 30 in Go for the same task.
The solution would be much less error prone and most likely much quicker as well.
How do you know this, i.e., can you share a benchmark -- the test data, not the results -- that you ran?
Data.table was the clear winner for my use case.
"After "exploring" many options, kdb+ was the clear winner for my use case. I process an average of 1.6TB per day."
See how silly that sounds without any details or credibility? Does "clear winner" mean it was faster? How does a person know if the recommender even tried both solutions? What if someone else's use case is dfferent from the recommender's? How could someone verify that one solution was faster than the other?
The easiest way is to provide sample data and a processing task and let people try the two solutions for themselves.
The naming conflict makes googling the differences fairly challenging.