Why Databricks Is Winning
cloudnativeenterprise.substack.com
cloudnativeenterprise.substack.com
We run Hadoop & Spark internally, but the team is underfunded and stuck in a constant cycle of fighting fires. And the result (and part of a larger push of the company due to the same cycle of under-funding and culture issues) is that we're moving our petabytes of data to cloud providers into their systems. Not only is the cost of doing this dwarfing that it would take to actually fix our issues, but we're going to lose the people who know how to design and manage petabyte scale hadoop clusters.
We wind up in a situation where we locked up data fundamental to our company and our position in the market with a 3rd party, and losing the talent that would allow us to maintain full control over the data. If the service increases prices, changes it's offering, or we get to a point where the offering doesn't meet our needs- we're fucked.
It's nice that Databricks has a nice "offramp" that you can take to go somewhere else, but the general idea is the same.
If you want to see de-skilling in action, hard core, go look into the service providers, wireless and fixed. They are running on fumes and attrition-victorious teams that last had a new technology in the late 70s/early 80s because everyone good at networking went to the FANG predecessors and then FANG proper.
Anyone actually talented just leaves after 2-4 years, so you get zero cohesion or snowballing effect of talented tech leaders who other good people actually want to work with.
These places become total career death. Eventually it becomes a place where only people with kids or other significant life obligation go, specifically because they value no-ambition bureaucracy that trades off compensation and risk taking in exchange for work life balance.
Those places are so soul crushing. Avoid at all costs.
1) The team is competent but picks up migration work to arbitrary technologies, and approaches with no clear ROI. These migrations block feature development and never seem to end e.g. teams ceaselessly migrating from GCP to AWS, to kubernetes, to podman from mysql to postgresql etc.
2) The team is operationally heavy and generates arbitrary requirements for everyone else to follow. The toolchain seems to get worse over time, and the number of hoops to jump through to get anything done endlessly grows e.g. Wait 2 weeks and get three business approvals for a server which you aren't allowed to have root access to.
3) The team has big ideas, but the business constantly under invests. The team is called to fight every fire but unable to stop the fires through any meaningful project. The company ends up on a platform thats constantly on fire.
When weighing these execution risks building an internal platform for just about anything looks incredibly expensive. I've only been at 1 company out of 6 which nailed the internal platform tooling requirements. The only reason I can attribute to their success was through quarterly NPS surveys on the developer experience for every major piece of the companies toolchain + hard to meet SLAs for uptime.
The underlying symptom is bad project management and more work than the team can handle but convincing the business that the ops team actually needs to be twice the size and have dedicated managers falls on deaf ears.
Every company outsources something fundamental. Does Google mine its own metal? Generate its own electricity?
Even if you did it on prem, that doesn’t save you from license renewal costs or upgrades. You can write the software yourself, but that’s not cheap either.
I don't think anyone is arguing one should maintain their own silica, atoms, independent universe, etc .
Solipsystem, Inc. begs to differ. Having it your way is a booming business nowadays.
(j/k, borrowed an SF title https://en.wikipedia.org/wiki/Solip:System )
Why is that different from "lose the people who know how to design and write accounting systems from scratch?" That was what happened when packaged accounting systems showed up. I'm not sure why you would want to preserve knowledge of Hadoop if other technologies are more efficient.
It's not like stem cells worry about losing the ability to excrete bile when they become neurons.
* AWS has a managed Spark offering called EMR
* EMR pricing (https://aws.amazon.com/emr/pricing/) is lower than Databricks pricing (https://databricks.com/product/aws-pricing)
* Databricks notebook development experience is better than EMR (but still really basic compared to IntelliJ / PyCharm text editing)
* Both Databricks & EMR have proprietary Spark runtimes
* Databricks is building a Spark runtime in C++ that might be faster (Delta Engine)
* Spark lets you process massive datasets easily, with small teams. 2-3 person teams can build data ingestion pipelines to clean & process terabytes of data a day. It's an incredible technology.
* The difference between the PySpark & Scala APIs confuse the hell out of people
* Whether or not ppl can run Python machine learning models on Spark clusters confuses people
* Overreliance on notebooks causes big issues (no version control, tests, deployment process, dependency management)
The big data ecosystem is constantly evolving and you need to study constantly to keep up.
The last three of those things might be valid issues, but since when are notebooks not just as subject as any other source code format to version control (I get that the difference between the UI and the on-site structure may make typical diff tools less-than-ideal, but VC itself is unaffected.)
I don't mean to sound harsh about other Data Scientists from a non software engineering background but the standard workflow is to fiddle around with a notebook until you can get a result. That's as far as it goes, no real robustness to it.
That's a pretty big generalisation but in organisations where they "home grow" their Data Science capability many of the online courses don't cover production level Data Science.
All the notebooks are in one place. Some are for important production jobs, other are for data exploration.
It's easy to make a little edit in a notebook and accidentally break production jobs.
Comparatively harder to make an edit in a git repo and do a deploy that'll break production jobs (e.g. if the JAR doesn't compile or the CI errors out cause the tests don't pass).
Notebook based production jobs get even more dangerous when NotebookA depends on NotebookB and so on.
You're right that all those issues can be sidestepped if you build projects in version controlled Git repos, test the code, and deploy JAR / Wheel files.
Speaking of testing, can you let me know if this PySpark testing fix worked for you ;) https://github.com/MrPowers/chispa/issues/6
My point is that you can do that even without jars/wheels - you can do VC and tests of notebooks. For example, https://github.com/alexott/databricks-nutter-projects-demo
Problem is that it’s limited. You can’t, for example, commit multiple files in one go (that improves a little bit in https://docs.databricks.com/projects.html), or merge changes a colleague made with your changes. You also have to use the UI Databricks provides. You can’t use a git CLI or whatever GUI you prefer. (all AFAIK, but I’m fairly certain about it)
There is also my rinky-dink open source project, Flintrock [0], that will launch open source Spark clusters on AWS for you.
It's probably not the right tool for production use (and you would be right to wonder why Flintrock exists when we have EMR [1]), but I know of several companies that have used Flintrock at one point or other in production at large scale (like, 400+ node clusters).
[0]: https://github.com/nchammas/flintrock
[1]: https://github.com/nchammas/flintrock#why-build-flintrock-wh...
I think I should start a company around minimalistic data tooling or the like - the amount of waste seems large across the industry.
I saw the de-skilling a few year ago, where a guy stiched together a compete application from a couple of SaaS APIs. Cool, but it somehow does not impress me.
Where your model breaks down is when you have folks throughout your engineering org who have data needs but don't have folks to spin up bespoke pipelines like this and spend time optimizing them. You have data, you want to write something SQL-ish or have some nice APIs, and you want to get results dumped into a predictable place. When you have mixed workloads, mixed data sources, and those jobs are being tweaked and changed frequently, the actual underlying compute cost is hardly the issue. Getting the data, chewing on it [fast enough] without having to spend much time optimizing, getting it to its destination, and making that happen regularly and reliably is where the value is.
I personally had this experience when first wanting to learn spark, and it really turned me off to the whole spark ecosystem. Curious if you have any suggested resources that do a good job on this?
Scala & PySpark are both great options. Lots of devs are terrified of Scala, so PySpark is more popular now. I'd say Scala has a slight technical advantage, see here for more details: https://mungingdata.com/apache-spark/python-pyspark-scala-wh.... Both are great overall.
I wrote a book that's a practical introduction to Spark: https://leanpub.com/beautiful-spark/
Most of the training materials are theoretical, which makes Spark seem really intimidating. You can learn some basic Spark principles and get up-and-running with production workflows quickly.
[1] https://spark.apache.org/docs/latest/running-on-kubernetes.h...
The DataBricks notebooks have a lot of value (for now) but ever since running Spark on Kube... I have had literally 0 cluster issues. It's absolutely shocking, coming from a YARN-based environment, where I was constantly plagued with ResourceManager issues, autoscaling issues, preemption issues, network disconnect issues...
Have you seen any good Spark vs. Dask benchmarks?
In my experience, Dask really shines when you implement custom numpy computations that could only be done in spark UDFs. We saw a decent performance difference there, but for common built-in computations I’d imagine that spark has better performance.
Edit: after some googling I found this paper with benchmarks.
https://arxiv.org/pdf/1907.13030.pdf
> Results show that despite slight differences between Spark and Dask, both engines perform comparably. However, Dask pipelines risk being limited by Python’s GIL depending on task type and cluster configuration. In all cases, the major limiting factor was data transfer.
https://towardsdatascience.com/supercharging-hyperparameter-...
https://towardsdatascience.com/random-forest-on-gpus-2000x-f...
Disclaimer: We produced those benchmark. I'm a founder of Saturn (https://www.saturncloud.io/) and we focus on providing Databricks-like capabilities with Jupyter + Dask + Prefect, so I definitely have strong feelings in this area.
That said, Databricks is a much broader platform, with all of the collaboration environments and is generally much more programmable than Snowflake.
It's interesting how the modern data lake is developing in this way, recreating many patterns from the traditional database for distributed systems and massive scale: SQL and query optimization, transactions and time travel, schema evolution and data constraints...
Having started out as a database developer / DBA many years ago, working with data lakes today reminds me in many ways of that early part of my career.
I wrote a post tracing a common interface from the typical relational database to the modern data lake.
This said, cost didn't really spiral or become an issue for us, especially when you take into account the cost avoidance of administering the Spark cluster and all of the tools. However, the pricing model did feel a little opaque.
We found that notebook-based development is actually an antipattern for software engineering. It was ostensibly helpful for the narrower "data science" use case, but we have a much more robust ETL and research platform we built on our own using Pandas, Dask, Prefect and AWS.
And personally I hated writing code in notebooks. If you're attached to that, you can basically get the same thing by using PyCharm in scientific mode with cell execution.
I can build a business like that: sell Hadoop consulting services at a loss, below the market price. Win lots of business by being really cheap; grow as fast as your cashflow allows, which will require annual investment rounds of ever increasing size.
When I look at Google Trends, Hadoop's long-term trajectory looks like "dead by 2024". Do we expect that Databricks will be able to justify a $8B capital raising if Hadoop doesn't really exist? If not, what's the path to profitability?
(Not trying to say anything bad about Databricks or their management, I'm just genuinely wondering what's going on.)
1. Your metrics are wrong about revenue/employee. We’ve heard they’re growing well by all standard SaaS company metrics 2. Hadoop consulting business is just a completely incorrect description of the platform. Hadoop is dying because of companies like Databricks, and all that needs to exist to justify their valuation is people bringing them large data workloads.
Their biggest problem is release maturity. Between Azure being down and them being down, stuff is down too often
Not a traditional file system, or even a HDFS clone.
I don't think there's a good one size fits all UI that can be applied the different types of work that take place with data and consumption by users. This is evidenced in Databrick's own feature set which includes an integration with R Studio and Tableau to Databricks as a data source rather than a work environment.
Snowflake predicate pushdown filtering seems quite promising: https://www.snowflake.com/blog/snowflake-spark-part-2-pushin...
Think both these companies can win.
The idea of storing your data in Snowflake and querying it in Databricks is pretty silly. Why would you want to do that? Why not just use Snowflake's compute? Sure you could argue Spark has some transformations that are hard to express in SQL, but that is why Snowflake introduced Snowpark.
And they’re Apples and Oranges. Snowflake is more comparable to Athena and BigQuery.
- https://managingml.substack.com/p/ml-model-training-is-an-et...
As an ML manager one of the things that I dislike about the Databricks model is the reliance on notebooks and bespoke interactive cluster workflows.
ML workloads are stuck in the stone age because there’s no common pattern to plug them into existing ETL frameworks, but really that is what they are.
Model training is an ETL that just happens to need unique domain tools to visualize how it’s succeeding / track observability metrics. But apart from that, you really need to automate model training and experimentation.
It is a total antipattern to ever use notebooks, or even non-notebook code developed interactively. Rather you should be submitting tasks to a task queue that maps your ML exploratory workload or training workload to a scheduled compute resource, runs it with observability baked in, and treats outputs like schema’d outputs of more traditional ETLs.
Instead of that, Databricks is much more of an early 2000s MATLAB model. Make it addictive and easy for unpragmatic researchers detached from real production use cases, then figure employers will have to adopt it since all the expensive-to-hire researchers can only be productive with it.
Long term I think it’s a very bad gamble. Just consider how much open source Python tools have eaten MATLAB’s lunch in the past 10 years.
We do have real-time collaborative notebooks. However, we also have long-running notebook scheduling where notebooks get picked by a runners and executed against certain environments. One cool feature is that you can watch the notebook's output as it runs, even if you had closed the browser or logged in from another computer. Anothr cool feature is that we do automatic experiment tracking, so we detect parameters, models, metrics, and the notebook automatically, removing the cognitive load from the use to "remember" doing that. People rarely remember clerical things, but even when they do, we didn't want to pollute the notebooks with tracking code of say MLflow. This happens automatically.
These long-running notebooks are simply jobs. There are other types of jobs: building a Docker image for the model so that you can push it to a registry is also a job, so is deploying a model to get a "REST" endpoint and a live dashboard to monitor model performance.
One of our guidelines is: "Anything that is possible with point and click must be possible with an API call", so we make sure to make everything programmable. For example, scheduling a long-running notebook can be done sending in a request. We also are working to add in webhooks and use cloudevents spec so that people can take information from our system, such as a job having finished, and trigger things in other systems.. Another is workload portability.
Again, it's still early and we're prioritizing the issues that have frustrated us the most for the past years in our projects. There's a long way to go, especially that there are many such as yourself who may have worked on problems we've never seen, or a scale/data we haven't handled before.
- [0]: https://iko.ai
I feel like instead of using the tool a lot of folks have specific expectations and if the tool doesn't fit exactly into how they work then it's just written off.
I’ll give a concrete example. In my org we use Dataproc on GCP as a model training task execution paradigm. You define your base environment via some Docker container, put it in GCR, and then define Dataproc jobs in terms of the base environment, the backing compute resources, any GCS bucket connections, and any ML-specific config like hyperparameters.
A human being never under any circumstances triggers these jobs. Instead a human user deploys the config as a cronjob or regular job in Kubernetes, and then a scheduler picks them up and runs them. For experimental workloads only, developers can manually trigger a Kubernetes job.
Each job consults config, spins up the appropriate Dataproc cluster, runs the job (with visualization tools exposed on ports at the cluster node IPs), and saves artifacts to GCS when done.
All of this is controlled via clean and easy internal CLI tools and wrappers to make it simple for any developer.
The number one thing this ensures is that no work ever exists in notebook format, beyond tiny scratch work a developer might do strictly to debug code or try a small data proof of concept.
The number two thing this ensures is complete reproducibility. Since every possible training task must go through code review, commit all config to version control, get impounded into a container, and execute via a deployed Kubernetes job, it is by definition impossible for someone to execute an ad hoc task that other engineers can’t rerun or have to follow weird setup steps to recreate (it’s all impounded in the container).
The third thing it ensures is that all accuracy, monitoring and results artifacts are explicitly tied to the Kubernetes job that controlled the process. It is not possible for some accuracy result to float around untethered from a job ID that uniquely and conclusively ties it to all relevant code for the job. This can be facilitated through MLflow or whatever else.
Getting data scientists to “wear corrective shoes” and learn to reorient their way of working to align it with this process has universally paid dividends, both for letting the data scientists experiment faster and more reliably, and for ensuring model training adheres to SRE-related compliance and best practices, so it is pluggable into various tools and constraints those non-ML support teams need in order to do their jobs and offer support to ML teams without getting hit with unstructured notebook spaghetti and bespoke execution paradigms.
I think this is wrong. This is what ML researchers tend to think when they are ignorant of why DevOps best practices exist and what other teams need in order to provide underlying infrastructure support. You’re only thinking of “what you need” from the point of view of the developer experience of the ML engineer, which frankly is usually the least important part by a wide margin.