Bigger data; same laptop
frankmcsherry.org
frankmcsherry.org
The amount of personal computing power every person with a modern computer has at their fingertips is staggering and almost nobody has any idea how to use it.
I'd like to think that's starting to change.
BTW if 90% time is spent in garbage collection, something is definitely wrong in the choice of tools or methods.
Funny ideas like your program needing to complete before the universe grows cold or shouldn't address more memory than there are particles in the universe don't get lots of attention at the academic level.
I think at the engineering level, where the theoretical CS gets turned into real systems, real resource constraints start to become serious problems. Things like the laws of physics and time really become important.
So for a 300GB graph, imagine you need to run some computation on data that size every day -- because of business reasons. So now you have to fit it into a 24 hour run-cycle instead of 40. So now you have to figure out how to shave off 40% off the run-time. Boom, you're now having to deal with performance optimizations.
For fun, let's say each time I get the result of that data, I can make a decision that earns $10,000, but at 40 hour run-times I can only make that decision 2 or 3 times a week instead of 5. So not only is it 40% too slow, but I'm only making 40% as much money.
Also, now you're dedicating a piece of hardware to that one task. And the data better never grow.
So you decide to rewrite the entire process in C, spend 2 months doing it, and now it runs in 5 hours instead of 40. I sunk $30k into the rewrite, but I'll make that back inside a week and it has room to grow as the data grows.
My other alternative is to figure out some distributed method that will get it down to <24hr run-times, but now I have to manage a compute cluster, with associated facility, electricity and management costs. I probably have to hire, if even part-time, a human or two to manage the cluster, so there goes additional overhead and administrative costs. Maybe I just go with Amazon, but there's downsides with that as well.
The solution for this problem is being worked on mostly at the engineering level, not the research one. And that's simply because the business forces are driving it to happen.
This is the real-life set of alternative choices computing groups are starting to contend with. It made more sense to distribute when 300GB of graph data was a lot. But now you can get a single machine that can fit that into RAM let alone disk. Single machines are outrageously powerful right now, but nobody has bothered to realize it.
It's an interesting space.
[1] http://www.frankmcsherry.org/graph/scalability/cost/2015/01/...
If you can write it so that it runs in 1/100th the time, do you even have to worry about it running for 24/7?
if you wanted to run this same algorithm on a 'real' single computer with ecc and high-performance cores, nothing would stop you, and it would get even faster, and still be radically less expensive in every dimension than the "big data" clusters, for this problem set.
Luckily networking is getting so good that remote desktops feel almost like local most of the time.
So basically big systems only make sense if the computation time on a single system exceeds the transfer time off of that system to the remote nodes + the time for that data to work it's way up the local node memory hierarchy.
This can happen when the data is big and the computation has unavoidable serialism in each part of the computation. Otherwise you'r likely to be able to coax a single system to just do it faster and with far less complexity than working on a compute cluster.
(I'm ignoring big systems with fast shared memory, like classic super-computers, where there effectively is zero transfer time since each node uses the same memory pool)
According to Google quick search, the bandwidth of DDR3 is ~12.8 GB/s, or ~104 gbps. As I understand it second-hand (read: cheap) infiniband givies ~40 gpbs, and newer, high-end gives ~100gbps[1]. Now, if you're using regular gigabit ethernet (or even 10gps with cheap-ish switches) -- you'd probably see around 5-8 gbps -- or if we're being very generous, on the order of 1/10th of cpu-ram bandwidth.
Sata2 is rated at ~6 gbps for comparison. In other words, 10gpbs ethernet/infiniband should be on-par with local SSDs, but "proper" infiniband should push the node-node ram-ram performance up compared to that. It appears fiberchannel is also moving towards the ~100 gps mark[2], but isn't there quite yet. And you'd have to have something on the other end -- with high speed networking it is quite natural to assume that both source and destination is RAM.
It's actually pretty crazy with ~40gps infiniband -- I'm considering getting some second-hand gear just to play with at home. Seems both cheaper and more reliable to work with than try to get 10gps ethernet to work (on a home/hobby budget!).
Eg:
http://www.ebay.com/itm/Flextronics-F-X430066-8-Port-SDR-Inf...
Or:
http://www.ebay.com/itm/Mellanox-MHGA28-XTC-InfiniHost-III-E...
All that said, I think the article makes a lot of great points -- and it does seem a bit funny that people seemingly are thinking of 1 gbps as "fast" when so much has happened on the cpu/ram side since that was remotely true.
[1] http://www.hpcwire.com/2012/06/18/mellanox_cracks_100_gbps_w...
[2] http://www.pcworld.com/article/2096980/fibre-channel-will-co...
Though, you'd be surprised at how often people get stuck on distributed problems where they just can't seem to get their cluster to process more than 80-100MB/sec and it turns out it's because they ran 1Gb ethernet between their nodes.
What these posts rail against is a tendency (post-MapReduce/Hadoop) to motivate and evaluate distributed systems solely in terms of their batch performance. This is a particularly common problem in graph processing research, where 2008-era general-purpose systems (i.e. Hadoop) have terrible performance, and so provide a straw-man baseline that a reasonably competent grad student should have had no difficulty in outperforming.
"Does my new system outperform a single thread?" seems like a proper hypothesis to test for a distributed system that claims only a performance motivation. Instead our grad student finds that there is a long list of peer-reviewed systems to compare against, and doing that should be enough, right? It would be fine if those systems had been tested against the first hypothesis, but in most cases they have not, so the new system can be much worse than a single thread (for a small set of benchmarks) and this will go undetected. The COST work is a remedy for this issue.
To return to your original question, I do think that distributed systems are good for some things. First, the OP concerns running graph algorithms on genuinely big data: 128 billion edges, and 3.5 billion vertices. A distributed system should be able to outperform Frank's laptop on that graph, and hopefully this post will encourage developers of new systems to aim higher and evaluate their systems on this scale of data.
The second is that distributed systems offer far more than pure batch performance. For example, geo-distribution is critical to providing a good user experience in web properties; and serving interactive queries with low latency and high availability is a particular interest of mine. We touched on the first part in the Naiad paper [1], and I expect this will become especially important for inference in online learning systems [2]. Hopefully the COST work will force the authors of these systems to try harder in their implementation, and we will all benefit as a result.
[1] http://dl.acm.org/citation.cfm?id=2522738
[2] http://www.cidrdb.org/cidr2015/Papers/CIDR15_Paper19u.pdf
I fantasize about this work inspiring a competition like the sort benchmark [2], but for single-threaded graph algorithms. Not only would this encourage people to develop usable single-threaded implementations, but it have the side effect of increasing the COST for all distributed implementations....
As for providing scalability break points, I think this work goes some of the way, but doesn't take into account the fact that the input data for a particular real-world computation might not reside in a conveniently-compressed format on a local SSD. Extracting the Common Crawl from an HDFS cluster just to replicate the blog post's experiments could take much longer than the computation itself; extracting the social graph from Facebook's Tao [3] to do offline processing might be similarly non-trivial. Academic experiments don't typically take this cost into account, but recent work like Musketeer [4] has started to consider it, and even offer ways to automate the system choice problem.
[1] https://www.usenix.org/conference/osdi12/technical-sessions/...
[2] http://www.sortbenchmark.org/
[3] https://www.usenix.org/conference/atc13/technical-sessions/p...
[4] http://www.cl.cam.ac.uk/research/srg/netos/camsas/musketeer/
My concern was that I saw the blog did the run on real life big(actually big) data and outperform everything on the market. This got me thinking, are we still not there on the problem size/computing power curve so as to justify use of distributed computing?
Of course distributed makes sense in a lot of places specially web and where problem is embarrassingly parallel but running graph algos is a bug chunk of the market and it seems for most firms a laptop will do.
Again, thanks for sharing.
I'd also wager that those Spark, GraphLab, or Naiad programs would be far from optimal, by at least an order of magnitude. My ideal would be to see someone attack the distributed graph processing problem with the same eye for systems performance issues that the TritonSort folks brought to the distributed sort problem [3] and MapReduce [4].
[1] http://law.di.unimi.it/webdata/twitter-2010/
[2] http://law.di.unimi.it/webdata/uk-2007-05/
Then you need to deal with distribution between a large number of systems over network... This will likely force you to choose a different algorithm, but at least you'll have something the benchmark it against from local best case scenario.
When it comes to graphs, I think good partitioning is very important in order to minimize average case latency. So it'd greatly help if you can somehow partition the data to minimize the number of nodes referring to a graph node on a remote system. When walking the graph, instead of immediately sending a query to the target system, it's probably better to batch the queries. Or not if latency counts more. Either way, latency will hurt badly. On a local system it's 100 ns. When the data is on a remote system, you're probably talking about 1000000 ns.
To combat latency, maybe the best way is not to send the request in the first place? Maybe you can ensure nodes referring to a remote system are duplicated to some extent on multiple remote systems. It could help performance -- or hurt it more, like when you need to modify the nodes.
Maybe you should have counters on nodes with connections to nodes on remote systems. If some remote system is referring to it more than local system, migrate the node there instead. It's of course not hard to see how this strategy could blow up. :-)
Please don't take above ideas seriously if you're developing a graph processing system. I wrote those to illustrate the biggest basic problem of distributed systems, latency. Distributed systems are very hard.
No it doesn't.
Reading doesn't cause (much) wear on an SSD - writing does.
(To be technical - it does, but it only requires 1 write / 100,000 reads or so. - https://en.wikipedia.org/wiki/Flash_memory#Read_disturb )
I've added the timings from the original post to the tables in this post, which also makes clearer the intended point: the existing systems don't yet have measurements for data at this scale (except Facebook's modified Giraph, which I want to hear more about :))