FastSpark: A New Fast Native Implementation of Spark from Scratch
medium.com
medium.com
Spark is still in-between yarn and kubernetes.
That's apples to oranges - because dask does not expose a SQL syntax that needs a query optimiser.
Also pyspark has the additional issue of serialisation between python and jvm. Turns out that just getting rid of that is a huge performance boost.
Most operations on dataframe-like objects can be described in SQL operations. Spark supports these operations and Catalyst can optimize query plans for these.
You are correct in that Dask does not optimize for this because Dask operations are more primitive hence it does not have the correct level of abstraction to do query optimization, only task graph optimization. Which reinforces my point that if you have a SQL-like workload on Dask.dataframes, chances are Dask may not outperform Spark.
Every Spark production engineer i know, translates the SQL written by data scientists back into high performance RDD code.
That's the advantage of Dask - there is no SQL abstraction needed. Pandas Dataframes are already the lingua franca of data scientists.. in fact, orders of magnitude more than SQL ever will be.
TLDR - Dask doesnt need SQL because the people who will push Dask to production already are far more comfortable in Dataframes than they ever will be in SQL.
You may still argue that spark RDD is faster than dask (and you may indeed by right)...but not having a SQL engine is not a problem for Dask.
Not really. The Spark API has equivalent calls in both Scala and Python, with Scala being the superset. Spark's SQL is a high-level abstraction that internally mapped to these operations.
> Every Spark production engineer i know, translates the SQL written by data scientists back into high performance RDD code.
This would be very unusual and rarely advisable with Spark > 2.0. Spark Dataframes are generally more memory-efficient, type-safe and performant than RDDs in most situations, so most data engineers work directly in Spark Dataframes -- dropping to RDDs only in specific situations requiring more control.
If you know data engineers who are somehow translating SQL into RDDs (except in rare circumstances) you might want to advise them to move to Spark > 2.0 and change their paradigms. They might be working with older Spark paradigms [1] and might have missed the shift that happened around 2.0 and missing out on all the work that has been done since.
> That's the advantage of Dask - there is no SQL abstraction needed.
SQL is only a language to access the dataframe abstraction (Spark Dataframes, Pandas dataframes, etc.) -- the fact it is higher-level means it is amenable to certain types of optimization.
If you take the full set of dataframe operations and restrict it to the set that SQL supports (group by's, where's, pivots, joins, window functions, etc.) you can apply query optimization.
Dask does not, and hence allows more powerful lower-level manipulations on data, but it therefore also cannot perform SQL-level query optimization, only task-graph level optimizations.
This Spark vs Dask comparison on the Dask website provides more details [2].
> Pandas Dataframes are already the lingua franca of data scientists.. in fact, orders of magnitude more than SQL ever will be.
I wonder if this is where our misunderstanding lies -- I sense that you might be thinking of SQL strictly as the syntax, whereas I use SQL as a shorthand for a set of mathematical operations on tabular structures -- which is equivalent in the subset.
[1] https://databricks.com/blog/2016/07/14/a-tale-of-three-apach...
Native dependencies usually mean you’ll need docker. Spark pre-dates docker, and just relatively recently added the Kubernetes runner, which makes dockerized jobs easy. But historically it hasn’t been easy to run a job in a containerized environment with the native dependencies you need. You can ship native deps with your job, but that’s not easy, especially if you need a rebuild with each job.
The main advantage of Spark is flexibility and interoperability. You save time by not having to write something optimized on day 1 (for something you might throw away). And you get SQL support, something Beam / Hadoop don’t have (certainly not for Python). There are lots of benchmarks where Spark SQL is not a winner, but the point is Spark will help you save development time.
A Spark inspired framework written in modern C++.
There’s nothing automatic about it, you or someone else will need to put a lot of work into leading the community, merging pull requests, debugging, etc.
(Sad to say, promotion too, in a lot of cases.)
The best rule of thumb I'm aware of is: unless you can't fit your computation on a single machine or your jobs are likely to fail before completing from the size and length involved, you are generally better off without Spark or similar systems. And if sampling can get you back onto a single machine, then you're really better off.
Perhaps I'm alone here but I'd prefer the title say Apache Spark explicitly.
[0] https://en.wikipedia.org/wiki/SPARK_(programming_language)
:( I guess I can read Medium posts only during first couple of days in a month
Stop hosting your content on a platform that holds it hostage so that it can make money off it without giving anything back to you.
The examples seem to be implemented in pure Rust. No one is going to port their Spark jobs to Rust in the shot term. Have you evaluated perf with Python etc?
If you're still seeing significant speedups, you might want to bottle this up and seek VC because a managed service along the lines of 'databricks but 10x faster' would certainly get traction.
This is an economy where content competes for clicks, not clicks competing for content. The author of that content wants me to see it, Medium doesn't want me to see it. I don't care enough to try to circumvent their arrangement.
Given the number of votes on my root comment, it seems neither do most people.
Learning this just made my day!
You can use the following options individually or in combination.
Option 1 : Pipeline DB extension (PostgreSQL)
Option 2 : Service broker in commercial SQL databases or building PUSH/PULL queue if not supported. There are many libraries in each programming language which tries to do that. Also see option 4.
Option 3 : Using CDC or Replication for synchronous or asynchronous streamed computation on single or multi node cluster
Option 4 : Transducers. For example, you can compose many sql functions or procedures to act on a single chunk of data instead of always doing async streamed computation after each stage of transformation.
For data cleaning, processing, analytics, ML on decently large datasets? Spark wins out
Dask+Perfect is a much better experience all round including perf, with virtually none of the cluster management hell involved.
Unlike Airflow, this lends itself to microbatching and streaming. Plus a bunch of housekeeping items ticked off that Airflow never got around to. With a bit of devops engineering time, you can have perfect manage the size of your worker cluster on k8s and scale it up/down with ingest demand, etc.
I'll say one thing though. The Perfect website used to be a lot more technical and explicit about what it is and isn't. Now it's mostly sales gobbledegook. Maybe not a good sign. I've seen this happen before with dremio.
Do you run dask on k8s ? I have been concerned that dask does not leverage kubernetes HPA for autoscaling...but instead chooses to run an external scheduler.
How has your experience been ?
Not the most SEO-friendly choice of name. Great product though.
Data cleaning -> PL/SQL and various inbuilt functions for the transformation of data (or new UDF if required at all)
Processing -> PostgreSQL Parallel processing on the local node and Citus DB extension for distributed computing and sharding
Analytics -> Many options here. Materialized views OR Triggers OR Streaming computation with PipelineDB extension OR Using Logical replication for stream computation
ML -> PG support variety of statistics functions. It also supports PL/R and PL/Python extension to interface with ML libraries.
Also, there are various kinds of Foreign Data Wrappers supported by PG - https://wiki.postgresql.org/wiki/Foreign_data_wrappers
PG is great but it's not suitable to be a feature store and sure as hell not suitable to fan out ML workloads. In a modern ML stack, PG might play the role of the slow but reliable master store that the rest of the ML pipeline feeds off.
depends on the scale? Not everyone processes petabytes of data.
> PG might play the role of the slow
You have any benchmark in your hand to support this? I believe highly optimized C code in PG can be significantly faster than Scala inside Spark.
There's no question about this. If you can express your task in terms of PG on a single instance, then you probably should.
When you get to more complex tasks, like running input through GloVe and pushing ngrams to a temporal store, PG offers very little - which is fine, it's not at all what PG is designed for. Inter-node IO eclipses single node perf, which is why Spark is used despite being a terribly inefficient thing (although in the case of Spark, it's so inefficient that for interim sized workloads you'd actually be better off vertically scaling a single node and using something else). PG won't help at all with these tasks.
Also, that smorgasbord of extensions GP listed isn't offered by any cloud vendor as a managed service afaik, meaning you must roll and manage your own. Depending on your needs, that might be a show stopper.
why exactly you think PG will not do this well?
Then you join first and second table.
Also, I think typical scenario is to resolve embeddings in your model code or data input pipeline.
Correct. PG has no place in this workload other than being the final store for the model output. And even then, you'd be using a column store like Redshift or Clickhouse. PG not even suitable for the ngram counters because its ingest rates are way too slow to keep up with a fanned out model spitting out millions of ngrams per second in addition to everything else going on in the pipeline.
You -could- probably do it all in PG. But that'd be a silly esoteric challenge exercise and not something anyone would try on a project. I am sure you recognise that.
But you have a year worth of historical data that you want to work with. If you're able to process 1m ngrams per second, it'll take a couple of days to get through that. You probably want to get closer to 10m/s if you're tweaking your model and want to iterate reasonably quickly. Of course there's ways to optimise all that and batch it and whatnot, but basically any big data tasks with the need to work on historical data and iterate on their models, quickly end up with kafka clusters piping millions of messages per second to keep those iteration times productive.
Ultimately this post is about Spark, and the comment that started this was someone listing PG 'replacements' for traditional ML pipeline components. If you need Spark, you're at scales where PG has no place.
Also in typical ML pipeline as I mentioned you can generate ngrams in input function of your model (Dataset API in TF), you don't need to store it somewhere.
It's also graph analysis and ML models.
ML models - I already mentioned how to uplift R and Python functions to SQL function. even if you are not using PostgreSQL many other databases help you with uplifting and interfacing with existing ML libraries through FFI