The new MacBook Pros apparently have north of 2GB/s of disk bandwidth for their internal SSD.
https://9to5mac.com/2016/11/01/the-late-2016-entry-level-13-...
The new MacBook Pros apparently have north of 2GB/s of disk bandwidth for their internal SSD.
https://9to5mac.com/2016/11/01/the-late-2016-entry-level-13-...
It's worth keeping in mind, because it means that a whole lot of workloads that people think they need a cluster for end up being far faster and cheaper to put on a single server these days, given the cost and overheads of sufficiently high speed network interconnects to move the data around at a fast enough to compete.
Of course there are plenty of workloads where you still need a cluster, but people often don't realise just how much you can put into a single commodity server today at a reasonable cost, nor what kind of overheads going from a single server to a cluster tends to add.
Or does it depend on the scenario?
If its disk, maybe a one box approach is suitable.
If CPU, where many current (F)OSS RDBMS are limited to one-core-per-query, a cluster can help.
Additionally a cluster can manage a larger working data set in RAM than a single box can do.
If you're going to partition the data because of limitations like that, it may still be faster to run a "cluster" on a single server, depending on data volume.
> Additionally a cluster can manage a larger working data set in RAM than a single box can do.
That's true in theory. In practice I've never seen anyone max out the amount of RAM possible in off the shelf x86 servers. I'm sure it happens, but it's not common, and it's fairly unusual for people to even approach the point where adding ore RAM starts getting more expensive than adding more boxes.
You can currently fit at least 8TB of RAM in the biggest x86 boxes; possibly more by now, and boxes that can fit 1TB-2TB are still relatively cheap.
Unless you plan to exceed at least 1TB, it's likely going to be cheaper to stuff more memory in a box than add enough extra boxes to compensate for the communications overhead.
Of course there are plenty of datasets where you would exceed 1TB for your dataset, but most people don't get anywhere close to that.
I would say that before a hardware cluster (as opposed to dirty "hacks" like sharding on a single box) makes sense if you are buying new hardware (the maths of course looks different if one option is to make do with existing servers), you will need to have requirements that exceeds one or more of:
- 48-64 x86 cores.
- 1TB-2TB of RAM.
- 2GB/sec aggregate disk bandwidth.
For all of these, it's worth noting that you can't compare like for like, as the moment you go to a cluster you have the according overheads in both CPU, RAM and disk IO of having to spread the work out, write data to disk more places, duplicate data in memory etc..
Note also that you can do better than that on plain x86 hardware, but the above is the point where the incremental cost of increases starts to really hurt and eat into any savings vs. a cluster.
I'm currently experimenting with a 4 node cluster running Apache Drill (on top of MapR's Hadoop distro), with 500GB RAM, 96 cores, and 28TB of SSD. For some analytic queries that span multiple date partitions, its actually easy to see 70% RAM capacity being utilized, and 80% of the CPU. The size of the data set is currently about 7TB, but that replicates for both data locality and redundancy.
We are a Postgres shop, and have tested with some of the options that run a "single node cluster" but one of the limiting factors is that RAM is divided between those DB processes, and performance is impacted. Plus we then need to implement some form of replication to obtain redundancy, fault tolerance, or HA.
Quite simply we are not using a cluster because our data won't fit, but because we can get better performance and reliability on 4 nodes. We likely won't have the need for 400.
There are cases where for non-interactive data processing we definitely don't need a cluster. At its simplest, multiple copies of Awk scripts can be run using GNU Parallel, Python scripts can be parallelized with the Multiprocessing module, and Go has been useful for concurrency (especially to us Python programmers).
aws s3 cp --recursive bunch-o-data s3://some-bucket/
spark-ec2 --region eu-west-1 --identity-file s.pem --key-pair=spark --instance-type m3.2xlarge --slaves 40 -v 1.5.2 launch my-cluster
Is significantly easier than making PG work easily at the 100GB scale in my experience. spark-ec2 is a script that ships with Spark to make it easy to set up a cluster in AWS.[0] http://spark.apache.org/releases/spark-release-2-0-0.html#re...
Our HBase cluster had 6 * 2Tb disks: about 8.5Tb of usable storage (the other 3.5Tb accounts for data replication/duplication) per host in the cluster. However, you need about 200 bytes in memory per kb on disk and should assign only 32Gb of heap to HBase. That's 2.5Tb wasted, per host. Couldn't just plug those disks out and use them somewhere else: you need all the disks in parallel to overcome the IO/bandwith bottleneck.