A comparative analysis of SQL-on-Hadoop systems for interactive analytics
arxiv.org
arxiv.org
Is there any advantage in using Hadoop for such a use case?
Even something like PG with the Citus extension to compare parallelism would have been relevant.
EDIT: Although another comment suggests that the data for the paper could have been gathered as long ago as 2016, which could mean Citus was not yet open source.
Don't jump on the SQL-on-Hadoop bandwagon without understanding the trade-offs. Here are some:
- If you use Parquet as your storage format, you lose indexes (ORC does support indexes though). If you use JSON, you lose a bunch more stuff including predicate pushdown, which are critical for WHERE clauses.
- With schema-on-read, you lose a bunch of potential performance optimizations that come with schema-on-write.
- Joins on a distributed database are much more difficult and usually much less performant than on a single node.
- Depending on the SQL-on-Hadoop system you choose, you stand to lose a ton of ability to perform complex analytic queries with pivots, windowing functions, etc. which are commonplace in most traditional SQL databases.
- You also lose the ability to rewrite data easily -- UPDATES and upserts are much more difficult (HDFS is an append-only file system, and is best suited for storing completely immutable data; its optimal mode is write-once read-many).
- HDFS also doesn't like small files, so you either have to buffer your writes before writing nice and large 128MB blocks, or have a nightly job to do small file compaction. For enterprise data, small rewrites are often needed to incorporate corrections, backfills etc. due to data arriving late, or due to errors that need to be fixed, and Hadoop does not handle this case well at all.
So as you can see, traditional SQL databases, if your data fits in them, present you with a tightly-coupled optimized environment that distributed databases are often not able to. They support also easy replication (real-time replication can be achieved via CDC).
However, if the scale of your data is beyond what traditional SQL databases are able to handle (and this differs from database to database; some MPP SQL databases can handle large datasets just fine), and if you need redundancy, SQL-on-Hadoop solutions may make sense.
One advantage of SQL-on-Hadoop (via Spark) offers is the ability to resume jobs in a long data transformation pipeline. If large SQL queries fail, they typically return nothing. With Spark, if you cache the intermediate results, you can carry on where you left off.
Longer coffee breaks.
I'm curious how Hive LLAP does, compared to Impala and Drill.
But having been part of the academic machine, I kind of understand. These things happen.
EDIT: I notice that a short paper was submitted to IEEE in 2017, so the study was probably done in 2016, which may explain the omission.
Hortonworks has been pushing their latest Hive LLAP benchmarks pretty hard. They're doing a kind of a Hive3 "roadshow" at the moment as well.
I was also really surprised that neither Hive nor Presto were included. Clearly the author is biased against FB originated projects \s.