We're also doing this at Thumbtack. We run all of our Spark jobs in job-scoped Cloud Dataproc clusters. We wrote a custom Airflow operator which launches a cluster, schedules a job on that cluster, and shuts down the cluster upon job completion. Since Google can bring up Spark clusters in < 90s and bills minutely, this works really well for us, simplifying our infrastructure and eliminating resource contention issues.
Awesome stuff, glad to see folks leveraging the possibilities! Perhaps as a follow-up you could write a guest blog on how this works for you! Feel free to ping me offline.
We built a coordinator that would spin up specific categories of machines for each stage (some stages were MR jobs, some were hadoop streaming jobs) -- for example when doing in-memory work it was useful to have fewer nodes with more RAM, etc.
We're also looking at making this a bit easier with some API changes (additions, really) and work with other OSS projects, like Apache Airflow. Streaming is a very interesting use case; a focus for us as well. Keeping clusters running for months/years presents some interesting challenges.
Same here. Built a pooling algorithm so we can have a maximum number of pools for each job type and distribute more jobs onto them when the maximum is reached. Clusters shut down when idle. It's served us well and saved loads of time and money.
interesting idea. I can see it being worthwhile with per-minute billing, which is not something supported by AWS. Also EMR charges a per-instance premium, does anyone know if Google Cloud does similar?
I'm curious how shuffle data is handled. Does the cluster intelligently scale down and move the shuffle data, or will the entire thing keep running while waiting for a single skewed reducer to finish? Or does the entire thing run on a single instance??
(Disclaimer: work on Dataproc) Though Dataproc doesn't natively have an autoscaler, Spotify wrote one as part of Spydra: https://github.com/spotify/spydra
The autoscaler here doesn't solve the downscaling issues mentioned. We've seen it work fine for scaling up, quickly. For now we're not scaling down, as the job will finish quicker than individual tasks will time out when unlucky.
Good question! You probably are familiar with the bandwidth and throughput power of the underlying storage system of GCS, Colossus, through use of BigQuery. BigQuery Storage and GCS storage both leverage Colossus. It's silly fast :)
Others can chime in more intelligently wrt Spark/Hadoop specifically, but I'll point out that read latency from GCS would definitely be higher than local-disk HDFS (esp Local SSD). Throughput, depending on your configuration, could be much better with GCS. Spark/Hadoop don't take the same care to optimize the storage-to-compute route as BigQuery, as evident by some bits of Hive performing serial FS operations.
So my answer is, it depends on the configuration of the job, the cluster, how data is written, choice of disk, and et cetera.
That said, when talking about price-performance, flexibility, scalability, and ease of operations, I suspect the "job-scoped clusters" setup would have a far superior TCO. We should try and do the math one day :)
For most use cases, GCS is going to give you better performance than using PD. GCS removes some headaches, like replication. In so doing, when you read from GCS you can often read from a large number of places in parallel; you're also less likely to run into bottlenecks because the VMs have pretty freaking amazing network bandwidth.
To be fair, there are some exceptions:
1: Small files (several KB to maybe 1MB)
2: Many reads
With GCS, you pay the tax for network overhead + SSL, which means reads are slower and scanning tons of small files can also be less performant.
For Cloud Dataproc, we also do provision HDFS on PD for intermediate output, because it's usually more cost effective than trying to do everything in GCS (which you can do, but you're going to pay for class A ops, and we want to save people money.)
To go into another level of detail, GCS also used to lack immediate list after write consistency. Now that GCS has been changing to support it, you also get a good story on clusters hitting the same bucket without hacks (in our case we used an NFS cache on the cluster.)
Dennis Huo, the tech lead manager for Dataproc also gave a pretty good talk covering some of this at Google Cloud Next [1]. I also gave a talk w/ Michael Yu on the GCS team [2] talking about a change we're in the process of making - swapping out the Java SSL provider with one from Conscrypt which nets about a 2x improvement in reads, making a great thing even better. :)
Can anyone tell me how, as a sole developer, it's possible to gain real-world experience with distributed Hadoop and Spark given the massive computing resources required? It just seems like a closed shop to me.
Define "massive". You can learn all the most important aspects of the APIs and programming model on a single machine, either by running in "local" (non-distributed) mode or by running a few VMs to simulate a real cluster. And you can spin up a cluster of 100 machines on either AWS or GCE for about a dollar per hour.
GCP offers a $300 credit that expires after 1 year. It's a good way to get your feet wet with Hadoop and Spark via Cloud Dataproc. (Google Cloud employee speaking.)