COST in the land of databases
github.com
github.com
In fairness, graph workload performance, compared to other types of workloads, is unusually sensitive to a diverse set of details that are rarely captured adequately by researchers. As often as not, you are actually measuring things unrelated to the graph algorithm itself. Measuring graph databases properly would require a level of specification for how to measure such things that hasn't happened yet. And many people like it that way, it makes it relatively easy to contrive details that make your particular algorithm or implementation look good despite the reality that serious weaknesses are pervasive across graph database platforms. A generally, objectively strong graph database platform doesn't exist AFAIK -- they all kind of suck in unique ways.
Just as in e.g. relational domain you can propose well defined operations (such as join two relations based on a common foreign key) and devise benchmarks for these operations, it should be possible to do the same precise thing for graph databases.
For a graph database, you would want (a) ingress, (b) comps, and (c) recovery benchmarks. The (b) comps per above can be based on a set of common patterns of queries applied to a parametric graph spec.
> In fairness, graph workload performance, compared to other types of workloads, is unusually sensitive to a diverse set of details that are rarely captured adequately by researchers.
That reads: we don't understand graphs.
Figuring out how to scale-out non-trivial problems is difficult. Many of these scale-out solutions are going to be worse than local processing, and many won't add meaningful value at any scale. However, there's real value in approaching difficult problems, and there's real value in "the fastest way to do X in a cluster" even if it's not "the fastest way to do X". That's because scale-out is a real requirement to do a lot of interesting things, and because if we can find better ways to scale out we can start to compete with single-machine solutions (and "custom cluster" and "supercomputer" solutions) then we have made real progress towards something that's faster at all scales. It's a valuable research goal.
Now, it would be great if the database community (and the systems research community at large) was more explicit about that. They should set the standard that authors need to clearly differentiate between "improved on result X under constraints Y" and "moved the state-of-the-art forward under all constraints".
It's also good to read systems research skeptically, and think about whether you really need to pay the cost of scale-out. For companies who's output isn't primarily "research", it makes a lot of sense to constantly look for "the cheapest way to do X", not "the cheapest way to do X in a cluster".
https://github.com/frankmcsherry/blog/blob/master/posts/2017...
I think that answering the "fastest way to do X in a cluster" question is only particularly interesting if X is something that you need to do in a cluster. But clearly they are testing X that is capable of being done on a mid-range laptop. There's a real question of whether the X that they are testing will actually scale to larger datasets than the laptop (or a decently beefy single node) could handle; it's entirely possible that the cluster solution will blow up at that point. And of course, if the single-threaded implementation is faster, then you could presumably just do "X in a cluster" by doing "X on a single node" and using the rest of the nodes to mine your CPU bound cryptocurrency of choice.
As the author hi,self points out, amusingly.
These numbers are for sure more interesting. SEED does better than any number I can actually achieve on my laptop. That demonstrates non-triviality, in that their system isn't just a laptop computing the answer and 156 cores mining bitcoin.
I really enjoyed the snark in this, and the series of blog posts he links in the beginning. I tend to be nice to a fault myself unless someone has shown themselves to not be acting in kind (and even then I confirm), but I admit I take quite a bit of voyeuristic pleasure in people being assholes in smart and interesting ways. Vicariously living and all that, etc...
That article is https://www.usenix.org/conference/hotos15/workshop-program/p...
*for a given definition of decent.
>None of them present any evidence that they are any better than a single-threaded implementation on any problem at any scale. So, if nothing else the published papers are not yet right. But I suspect they are mostly wrong.
It's (a) always there, and (b) always hilarious. From operating systems comparisons that need "we didn't actually slow down the processor physically" comparisons to look good to these "we look great, compared to that guy over there licking the floor."
"You can have a second computer once you've shown you know how to use the first one."
In general, I think people would be more convinced if there were real cases in industry where engineers switched from distributing graph algos to single-machine algos because they found performance benefits were better. In all my years, I haven't heard of a case like this (it's always been the other way around).
A database or hadoop cluster is not just an execution engine, but rather a managable, reproducible execution engine. Sure you can code up a smarter algorithm on your personal i9 workstation. However having something available as part of DB or Hadoop ecosystem enables reproducible execution while taking care of compliance (e.g. privacy), security, access control etc.
Eventually we will have a containerized extensions to DB query execution flow that will allow us to have best of both worlds but for 99% practicing data scientists downloading a Gigabyte sized subset to their own workstation is not a viable option.
The real opportunity lies in ability to add arbitrary compute/memory capacity to DB query execution flows.
Also, downloading 1GB is no big deal, downloading 1TB is a big deal, so I think your example has more merit in the > 1TB range. i.e. things that don't fit on a laptop easily.
Also most commodity compute nodes have allocation around 16 Gb. The big difference comes in 16~256 Gb range where having a single powerful node can make a huge difference.
10% cost for the benefits that this stuff can give you, that is pretty easy to swallow. 100%-500% perf. cost is just a very, very, very hard to swallow pill.
also rules out e.g. the fairly common practice of using R or pandas for some ad-hoc processing.
This is essentially the reason why all organizations are adopting Spark, since it allows you to write imperative code on dataframes and build ML models.Maybe you've worked at a job or two where nobody can comprehend not using distributed computing, as you describe, but it's nonsense to claim that "all organizations" work that way.
No, "all organizations" are not adopting Spark.
All organizations which already have a Hadoop cluster. You do not need Spark to use dataframes.
Never claimed this. To clarify Spark allows you to directly port Pandas code while leveraging existing Hadoop cluster infrastructure. And distributed computing is terrible for machine learning.
Distributed computing (Both traditional hadoop/spark and latest TF/PyTorch with parameter server) are essential for scaling ML beyond a certain point. Maybe you've worked at a job or two where nobody can comprehend not using distributed computing, as you describe, but it's nonsense to claim that "all organizations" work that way.
If you have experience routinely training models on Terabytes of data intended for production deployment. I am happy to hear. There is a vast difference between training a model on your machine for research and building a reliable ML system that scales across large datasets and teams while taking infrastructure costs into account.Here's how I would leverage Hadoop infrastructure to use Pandas: delete Hadoop so I've got more disk space to run Pandas.
I don't get to train ML on terabytes of data very often. I do NLP, so "terabytes" means training a background model on the entire Common Crawl. Usually I'm doing something more specific and interesting than learning about random web pages. But when I do deal with the Common Crawl, I deal with it on one computer. Terabytes are not scary.
How does distributed computing even help? ML models need memory locality, sometimes to the extreme of being localized within a GPU's memory. And the limiting factor is the ability to iterate over the data. Sending the data over a network during training would be the worst thing you can do there.
delete Hadoop so I've got more disk space to run Pandas.
Except when you have ~2000 node cluster that runs 10,000 ETL tasks daily all of which are IO bound, And you cannot "just" uninstall. However during certain periods the same cluster has significant underutilization this opens up possibilities of doing lots of cool stuff for almost zero cost.I can understand your confusion. Training embedding model on Common Crawl is a toy problem. I recommend you thinl from perspective of a Tech company ideally in a production ML setting to understand the cost tradeoffs that go into making these decisions. Regarding memory locality if the problem is small enough its possible to tune Spark to use fewer or even just one worker with enough amount of memory allocated.
Sending the data over a network during training would be the worst thing you can do there.
When your data itself is in 100s of Terabytes and sharded across multiple racks, often a properly tuned spark pipeline is more reliable and performs well. Again there is a difference between using Common Crawl subset that you manually download filter train etc and in ensuring that models are updated/trained automatically daily across 100s of TB data.You are posturing.
The reason Spark ecosystem has been so popular is because it enables these types of computations without breaking the model.
Today’s OSS big data query execution environments are BnL neither secure nor compliant.
This is not a knock. They were simply not designed for those constraints. They were built for academic or single tenant use cases without separation of duties / control.
This also applies to commercial products such as Splunk.
There is tremendous investment going on the last couple years in raising the security and compliance bar. We’ve worked on this with the usual suspects.
But barring a handful of proprietary stacks (the big three CSPs, and a couple enterprise on prem bare metal distros) getting close, we are not there yet.
Trying to land with your security or risk teams or regulators that the CSP pulled it off will likely take you longer than provisioning a compliant laptop you’d keep locked up with two keys.
Unless you have several tens of man years invested in in-depth security wrapping these environments, or can choose a big three CSP w/o answering to anyone, today I’d still recommend the trusted laptop build approach for truly sensitive algorithms and computations.
today I’d still recommend the trusted laptop build approach for truly sensitive algorithms and computations.
You are utterly wrong. Algorithms and computations (especially ML kind) are never sensitive, its the data which is always sensitive. And that ALREADY exists on the cluster.If an Organization already has a Hadoop cluster containing data. You are suggesting that somehow having it downloaded to a secured laptop is better? Than say Spark running on top of the cluster? I think you are deeply mistaken. The cluster instances are already protected (if not you have a bigger problems). Also while an organization might not have Hadoop, they surely have an RDBMS, in which case the algorithms are even more useful.
scaling for the fun of scaling.
I think your arguments are misguided, for every compute bound task where Hadoop/Spark undeperform, 1000 other ETL type tasks where hadoop is indispensable. As a result any organization running a large hadoop cluster will already have underused compute capacity for free and a maintainance staff which already taking care of the cluster. Thus from an organizations perspective the time difference is not material especially for batch jobs, this is the reason why Presto and Spark have been so successful. They enabled underutilized hadoop clusters to be used for ML and data science while delivering reasonable performance for Zero cost.“The” cluster, singular ...
This is the catch. If you need security and compliance, you can’t today derive the benefits of re-use by other teams.
Given today’s distros, and assuming your threat model needs to account for insider threat, you need a different cluster for each data ownership grouping and data sensitivity level. For four teams with three levels of data, you’d need twelve completely independent clusters.
Unless, as noted above, you’ve done a ton of in-house multi-tenancy work to provide full stack security and compliance assurances and audits.
His "smarter algorithms" aren't.
You are conflating a couple of things here.
Of course for real analyses, you want to run them on standardized environments with appropriate access controls.
But nothing says that "standardized environment" has to use a distributed database that, from all accounts, seems to harm more than help.
You could just have a RDBMS or even filesystem, which is appropriately backed up or replicated, on a single beefy VM or physical machine.
> The real opportunity lies in ability to add arbitrary compute/memory capacity to DB query execution flows.
Only if the tools used to provide that ability actually allow you to take advantage of it. If scaling up to 160 nodes doesn't allow you to be any faster than 1 single laptop, you probably could have spent a lot less money on just adding more RAM to a single beefy server than a 160 node cluster.
Athena is built on Presto, which has full support for transactions. Each connector provides varying support (isolation level, multi-statement writes, etc.) for transactions, limited by the underlying data source.
Edit: regarding transaction though, I wonder if Athena could even be used as a "meta sharding" layer, when all of the underlying data sources support transaction, but data is too large to fit in a single instance. The advantage would be to not implement sharding logic in code. Not sure about the performance though. Just a thought.
Edit2: not sure if I read it correctly, but looks like right now Athena does not support transactions yet (http://docs.aws.amazon.com/athena/latest/ug/creating-tables....) although using Presto there could be a future that they do support.
CSV is a very unwieldy storage format for Athena + S3, since Athena charges you by data scan sizes.
ORC is pretty much optimized for S3-like use-cases (of being read over HTTP), so you'll find that a columnar structure like that would be a much nicer way to store and query often.
There's a pretty good run-down of the options and formats here - http://tech.marksblogg.com/billion-nyc-taxi-rides-aws-athena...
"Our first paper is Scalable Distributed Subgraph Enumeration or "SEED". We will have a future post about rules for picking acronyms."
I suspect the author would agree.
He is limited by not having an "expensive macbook with 4 real cores" ?? He gave up on a test because it "paged out the 15GB of data".
Of course at these scales a cluster is x5 slower.
Can someone please explain - am i missing something here ?
BUT
With only 2TB of local storage, your 4TB dataset (or output set) has to pass over the net for every execution. At 10Gbs - this alone can take 20 minutes to 1hr.
If you want to write it to the local SSD - multiply that by x3 or x10 ?! ....
This 4TB of memory is still bounded by 4 x 30MB of L3 cache. Which means your single thread implementation will be slowed down by x3 or x5 due to memory latency.
Your multi core implementation will probably suffer even more.
Distributed system are VERY hard, but dealing with them is inevitable for certain workloads.
That's what I mean. You have to get really big to have Big Data. Fine, you can't do PageRank on a petabyte of web crawls using one machine. But the datasets that people use for benchmarking, at least the benchmarks that are made public, you definitely can. You can go far larger than a typical benchmark dataset and still do it on one machine.
Note that he already demonstrated that for many of the tests, his older MacBook was faster than these 160 core clusters. There were a couple where the older MacBook was not faster than a 160 core cluster, but he suspects that a newer one might be.
The one that paged out happened to be one where he artificially broke it up into 160 pieces to simulate what it would be like if he had 160 similar cores, but that caused enough issues with memory locality that swapping dominated the results. This artificial partitioning was to give a sense for what results on a 160 core cluster "should" be, but the swapping distorts the results enough that for that problem it wouldn't even be a very good rough estimate.
> Is this guy for real?
This isn't supposed to be a publication-quality paper. This is a blog post, using the resources he has immediately at hand, to show why if you follow cutting edge distributed database research you might be led down a wrong, and very costly, path, which you could avoid with a much simpler implementation on much cheaper hardware.
Note that he references an earlier paper that he did publish on this topic; this blog post is mostly just doing a similar check on some more recent published work, to see if much has improved in the field:
https://www.usenix.org/system/files/conference/hotos15/hotos...
So this complaint should be levelled at the papers and the academic communities- why is it acceptable to only provide empirical evidence for these distributed algorithms at data scales where they cannot out perform simpler implementations? Why not require empirical performance measurement on truly large scale examples for these kind of algorithms?
Strictly speaking, I think that term attributes more organization to the academic CS establishment than my experience supports.
The point here is that, if they're trying to establish that Smart Technique 2 is better than Smart Technique 1, they need to first establish that Stupid Technique Q doesn't work at all. Which they haven't.
If they don't have access to "truly big data" to do so, well, most of the people who want to use the Smart Techniques probably don't either.