Clickhouse Local
clickhouse.tech
clickhouse.tech
This is not a rhetorical question, I would really like to know why it gets so much attention here.
The basic idea behind Presto is that it federates other databases, and supports doing joins across them. From what I understand, the problem that it solved at Facebook is bridging the gap between different teams; if a team has MySQL and another has files stored on HDFS, it doesn't really matter because all you do is query Presto and it'll query both under the covers. The alternative is setting up data pipelines, and dealing with the ongoing issues of maintaining those data pipelines.
BigQuery is much slower and is much more expensive for both storage and query.
Databricks (Spark) is even slower than that (both io and compute), although you can write custom code/use libs.
You seem to underestimate how heavily ClickHouse is optimized (e.g. compressed storage).
Is it any more compressed than Apache Hive's ORC format (https://orc.apache.org)? Because that's increasingly accepted as a storage format in a lot of these analytical systems.
https://engineering.fb.com/core-data/even-faster-data-at-the...
https://www.altinity.com/blog/2019/7/new-encodings-to-improv...
Clickhouse manages the whole distributed storage, ram caching, etc. thing for you.
In my experience, a unified single purpose vertically integrated solution will be faster than a bunch of kitchen sink solutions bolted together.
DataBricks is essentially Spark, and I shouldn’t need a whole spark cluster just to get database functionality. It also costs money.
Unless I’m mistaken, Presto is just a distributed query tool over the top of a separate storage layer, so that’s 2 things you have to setup.
I have no experience with BigQiery but I’ve heard good things about it and Redshift, however but if the rest of your infra isn’t on GCP/AWS then that will probably be a blocker.
Clickhouse is open source, comes with convenient clients in a bunch of languages as well as a HTTP API. It’s outrageously fast and has some cool features and makes the right trade-offs for its use-case, large range of supported input/output formats, built-in Kafka support and the replication and sharding is reasonably straightforward to setup.
There's large complexity and cost overheads to Hadoop solutions, and not everyone has actual big data problems. ClickHouse hugely outperforms on query patterns that would devolve into table scans in a row store, while working at row store volumes of data without a bunch of big nodes.
[1] I say "perhaps" because I've only just started using it having migrated from MonetDB, but have no experience of alternatives like Presto.
[1] https://github.com/MonetDB/MonetDBLite-R/issues/38#issuecomm...
The closest open source thing that matches its feature set is Presto, but that one is quite different.
Apache Druid is supposed to be very mature, but also very difficult to set up and manage. I've not used it myself.
There's also Vespa, but I don't know how well it performs with large numbers of columns.
A lot of people use Elasticsearch for analytics. Being based on Lucene, it's kind of columnar, and it can perform very well indeed on aggregations.
InfluxDB may be good, but it's not fully open source.
On the opposite side of the spectrum you have other open source projects like questDB that have full focus on core performance: constantly optimise to get as much as possible from a single processor core. You can't scale out (at least yet), but given how fast it is on single core, it will be pretty powerful if they chose to go this route.
Greenplum: I've not used it, but it does support columnar tables, so maybe it's comparable.
I wouldn't be surprised if this inverted at large scale (say 30+ machines). Druid data servers are rebalanced automatically; if you're on AWS and decide to scale up by adding a new data server, it will automatically load its assigned subset of data from S3. If AWS kills one of your data servers, then other data servers will automatically load from S3 some of the data that server used to carry, in order to reach the desired replication factor again.
Last time I checked ClickHouse had no automatic rebalancing at all, which sounds horrendous unless you're running at very small scale or willing to have people babysit it. I haven't operated ClickHouse at large scale though, so if I'm wrong I'd be happy to hear how people manage ordinary tasks like scaling up and down, replacing dead instances, changing instance types to adjust CPU/mem/disk, etc... with let's say a 100 TB compressed dataset.
Another difference is that Druid can index all dimensions. So if you plan to run queries with filters that only match a small fraction of rows, then Druid can be faster than ClickHouse. Conversely, if your queries have filters that match many rows, then ClickHouse will be faster because it has higher raw scan speed. (At least that was the case about a year ago. Since then, Druid has added vectorized aggregation, which I haven't benchmarked, but I'd bet that ClickHouse is still faster at doing full table scans.)
IMO these 2 things are the main elements to think about when choosing between ClickHouse and Druid.
I do think rebalancing is a weak point for ClickHouse, although for our use case that would not be so much of an issue, and it feels like it is on the roadmap for ClickHouse this year, but we will see. And if you are on Kubernetes, some of that headache may be handled for you with the ClickHouse Kubernetes operator.
I will say, that Druid indexing comes at a heavy cost in hardware for ingestion.
We find ClickHouse can easily ingest at least 3x the rate of Druid on the same hardware, and since Druid is asymmetric in design, you then have to get even more hardware to handle the queries.
Even with the vectorized aggregation, ClickHouse is beating Druid for full table scans at least, especially high cardinality data. But the vectorized aggregation has some restrictions to get on the fast paths, so that may improve. as those are removed.
Overall, I find ClickHouse much easier to work with and manage compared to Druid. ymmv
My company migrated our time-series data from InfluxDB to ClickHouse last year (I personally led this, in fact), and the performance difference is night and day.
While I liked a lot of what Influx could do, it was also nonstandard in bizarre ways (Clickhouse behaves more like a subset of SQL), sometimes shockingly immature, and despite appearing fast when we first started using it, so slow that it was a considerable bottleneck.
But if you have a lot of time-series metrics (only numbers), you might be better with specialized time-series databases like Prometheus + VictoriaMetrics with Grafana for visualizing it.
However, influx, being schemaless, can be outstanding for rapidly prototyping ephemeral metrics, as adding a new measurement is zero-cost (just start writing it). It also plays great with grafana for building dashboards. Finally, I much prefer the ergonomics enhancements of the influx query language (v1, the v2 "flux" language looks terrible to me personally), particularly how duration strings are a first-class datatype ("group by time(1h)").
Interested to hear more about clickhouse performance, haven't had a chance to use it for anything where performance would matter significantly, although am aware it can be very fast.
The notion that you will get approximately the same query performance with all column stores is false. There can easily be an order of magnitude difference depending on the implementation. Take GROUP BY as a paradigmatic example of what OLAP stores do. Of course the way to implement GROUP BY is with a hash table but little tricks make all the difference and a lot of love went into the clickhouse implementation. Just to give you a taste: there is a custom hash table with specializations for different key types (e.g. it will store a precomputed hash for strings but not for integers and use it to speed up equality test). Variable length data is stored in arenas to reduce allocator pressure. Data will of course be aggregated in several hash tables in different threads and then merged together, but if there is a lot of keys each table will additionally be sharded so that the merge step can be performed in parallel too.
Of course you shouldn't trust random claims on the internet that clickhouse is fast and should do a small case study yourself. Then you'll appreciate how easy is to setup a clickhouse instance or a small cluster. It can easily slurp up most common formats. It is just a single binary with minimal dependencies that will run as-is on any modern linux. There is just a single node type (compare this to druid madness).
You are right that there is a lot of limitations and, how should I put it, quirks. This is resoundingly not a general-purpose database and someone used to the comforts of e.g. postgres will encounter some nasty surprises. Bugs are unfortunately common, especially in the newer functionality. But performance is its main feature and it makes many users of clickhouse put up with its limitations.
Actually, for low cardinality columns, _not_ using a hash table will speed up things.
For example, a dictionary-encoded column for states might have a few values in its dictionary (1 => CA, 2 => FL, 3 => NY) and the data looks like an array of numbers (eg. [1, 1, 1, 2, 1, 3, 3, 3, 1, 3, 2]). The fastest way to aggregate is to actually use an array, as the dictionary index conveniently maps to an array index.
Then, when it comes to merging several of those arrays, they're turned into hash tables.
Combine enough of those optimizations and proper data layouts, and you end up with several orders of magnitude of performance differences between engines.
As the parent says, try it yourself on your own data. That's all that counts.
Disclaimer: I wrote the blog article and we sell support for ClickHouse.
As a column store engine that supports (hybrid) SQL, manages its own storage and clustering, and has only one external dependency (Zookeeper), I believe its main competitor is Vertica which can be very expensive. I assume Oracle, IBM, and MS have column stores as well and for a cost.
Greenplum is a Postgres fork, and same scale requires much more hardware. Citus is row based, and therefore will lag in scan time for many OLAP query patterns. Presto, Hive, Spark, all of the "post-hadoop" options may scale larger, but will also lag in scan time, and have significant external dependencies - mainly storage.
Clickhouse is easy to install, configure a cluster, load and query. It does have limitations, but currently all horizontally scaled database platforms do.
The low latency to query execution is really nice.
If a conversion happens, you need to generate the views/cubes again in reporting dashboard, clickhouse makes it cheap and easy to run such operations over commodity hardware.
If you don't have this, you'll be using big query and it might not be as fast.
- the multiple table engines that are heavily optimized for specific data access patterns: ReplacingMergeTrees which essentially make records mutable, SummingMergeTrees that allow us to progressively build pre-aggregated data, AggregateMergeTrees that allow storing the intermediate aggregation state of most aggregation function and compose them at query time over multiple groups (example, store a p95 aggregation state hourly and query the daily p95 by composing them), and more.
- column data types is extensive and includes nested columns
- the architecture is relatively simple making it easy for developers and on prem users to deploy a single node local clickhouse very easily
- it is very efficient in inserting big batches of data which works really well for our use case were we ingest massive amount of errors.
- data skipping indexes, bloom filter indexes
(yes, as vlad@sentry.io mentioned below we are hiring for the team that manages storage and thus clickhouse)
https://blog.sentry.io/2019/05/16/introducing-snuba-sentrys-...
This basically replaces most of my usages of SQLite.
When its SQL "dialect" matures, Clickhouse will eat MySQL lunch, then PostgreSQL.
Therefore I think the lack of OLTP will not matter much and that clickhouse will be widely used, but also misused when it becomes too fashionable.
If you need deletes and transactions, look elsewhere, but Clickhouse seems to be great for what it's been designed for.
For example, aside from the lack of transactions, Clickhouse is designed for insertion. There's an INSERT statement, but no UPDATE or DELETE statements. You can rewrite tables (there's ALTER TABLE ... UPDATE and ALTER TABLE ... DELETE), but they're intended for large batch operations, and the operations potentially asynchronous, meaning that they complete right away, but you only see results later.
ClickHouse has many other limitations. For example, there's no enforcement of uniqueness: You can insert the same primary key multiple times. You can dedupe the data, but only specific table engines support this.
There's absolutely no way anyone will want to use ClickHouse as a general-purpose database.
So I insist: everyone will WANT to use clickhouse as a general purpose database, and will create ways to make it so (ex: copy table with the columns you don't want filtered out, drop the original, rename)
It is just too fast and too good for many other things, so it will expand from these strongholds to the rest.
A personal example: I am migrating my cold storage to clickhouse, because I can just copy the files in place and be up and running.
I know about insert and the likes, I have a great existing system - but this lets me simplify the design, and deprecate many things. Fewer moving parts is in general better.
After that is done, there is a database where I would benefit from things like alter tables or advanced joins, but keeping PostgreSQL and ClickHouse side by side, just for this? No. PostgreSQL will go. Dirty tricks will be deployed. Data will be duplicated if necessary.
* https://github.com/ClickHouse/ClickHouse/pulls?q=is%3Apr+mer... -- Recent work to enable merge joins
* https://github.com/ClickHouse/ClickHouse/pulls?q=is%3Apr+s3 -- Same thing for managing data on S3 compatible object storage
There's been a lot of community interest in both topics. Merge join work is largely driven by the ClickHouse team at Yandex. Object storage contributions are from a wider range of teams.
That said I don't see ClickHouse replacing OLTP databases any time soon. It's an analytic store and many of the design choices favor fast, resource efficient scanning and aggregation over large datasets. ClickHouse is not the right choice for high levels of concurrent users working on mutable point data. For this Redis, PostgreSQL, or MySQL are your friends.
> As i'm sure things like pivot tables and rolling windows are a PITA in SQL
I can't speak for clickhouse, but group-by and window functions are a very standard part of any SQL analysts toolbelt.
Because I have more data than what fits locally, there’s a data pipeline that pushes more in, and I only need to work on a subset.
Storing everything in flat csv/parquet etc is useless when there’s more than fits on your local/single machine memory or if you want to search/subset etc some of the data or do anything that’s larger than memory without having to write spill-to-disk stuff in Python/pandas.
My bigger point to the OP was why bother using a DBMS that is specifically tuned for running fast analytical queries if they only intended to use it as a storage layer and pushing all of the analysis into Python. Use a solution focused on storage, not analysis, if that's the use case.
Additionally, if I have my data in an actual database, I can attach tools like Tableau and Superset directly, as opposed to having to take it from storage and then put it in a database anyway before being able to visualise/use it.
[1] https://github.com/TileDB-Inc/TileDB
[2] https://docs.tiledb.com/developer/api-usage/embedded-sql
Disclosure: I am a member of the TileDB, Inc. team
There are some examples in the article cited by dang: https://news.ycombinator.com/item?id=20163017