Transforming Postgres into a Fast OLAP Database
blog.paradedb.com
blog.paradedb.com
1. Based on Clickbench results, pg_analytics is still far from top-tier performance. If you're looking for a high-performance OLAP database, you should consider the top-ranked products.
2. The queries tested by Clickbench are really simple, far from real data analysis scenarios. You must be aware of the TPC-DS and TPC-H benchmarks, because ClickHouse simply cannot run these test suites, so their Clickbench does not include these datasets.
Lastly, I want to say, if your enterprise is of a certain size, separating OLTP and OLAP into two databases is the right choice, because you will have two different teams responsible for two different tasks. By the way, I use StarRocks for OLAP work.
1. On Clickbench, make sure you're doing an apples-to-apples comparison by comparing scores from the same instance. We used the most commonly-used c6a.4xlarge instance. While a few databases like DuckDB rank higher, the performance of Datafusion (our underlying query engine) is constantly improving, and pg_analytics inherits those improvements.
Then again, people only care about performance and benchmarks up to a certain threshold. The goal of pg_analytics is not to displace something like StarRocks, but to enable analytical workloads that require both row and column-oriented data or Postgres transactions.
2. We're working on TPC-H benchmarks. They're good for demonstrating JOIN performance and we'll have them published early next week.
Most of the time, all that matter in terms of performance is user's tolerance. Once that is reached, operational complexity becomes a lot more important. We use raw Postgres for analytics, knowing that projects like these and cloud offerings like AlloyDB will make our lives easier (in terms of performance) as time goes.
pg_bm25 looks awesome too! Next up, take fdw to the level of Trino/Drill, and we dont need anything else other than postgres and its extensions!
You touch on this in your next sentence, but really, how many people need that kind of performance?
https://benchmark.clickhouse.com/#eyJzeXN0ZW0iOnsiQWxsb3lEQi...
With regards to ParadeDB, we rely on the Datafusion SQL parser, which can transform the Postgres SQL dialect into a Datafusion logical plan that can be executed by Datafusion. We actually have an open PR that adds support for user-defined functions...it will likely get merged within a few days.
If you decide to try it, I believe they now have a way to load in arbitrary Postgres extensions!
1) How does this deal with backups? Presumably the deltalake tables can't be backed up by Postgres itself, so I guess there's some special way to d backups?
2) Similarly for replication (physical or logical). Presumably that's not supported, right? I guess logical replication is more useful for OLTP databases from which the data flow to datalakes, so that's fine. But what's the HA story without physical replication?
3) Presumably all the benefits are from compression at the storage level? Or are there some tweaks to the executor to do columnar stuff? I looked at the hooks in pg_analytics, but I see only stuff to handle DML. But I don't speeak rust, so maybe I missed something.
3) We store data in Parquet files, which is a heavily compressed file format for columnar data. This ensures good compression and compatibility with Arrow for in-memory columnar processing. We also hook at the executor level and route queries on deltalake tables to DataFusion query engine, which processes the data in a vectorized fashion for much faster execution
How large part of the plan you route to the datalake tables? Just scans or some more complex part? Can you point me to the part of the code doing that? I'm intrigued.
1) You can find the code here: https://github.com/paradedb/paradedb/blob/996f018e3258d3989f...
For deltalake tables, we send the full plan to DataFusion.
2) Re: WALs and replication -- We are currently adding support for WALs
Or Debezium to Kafka for a more hand coded solution?
Debezium -> kafka -> kafka engine -> MVs -> Replacing merge trees works like a charm for me.
- Greenplum;
- Citus with cstore_fdw;
- IMCS (Konstantin Knizhnik);
- Hydra;
- AlloyDB;
For example, Greenplum is one of the earliest, fairly mature, but abandoned.
We see pg_analytics as the next-generation Citus columnar, with much better performance and integration into the wider data ecosystem via Delta Lake, and eventually Iceberg
why start with Delta Lake instead of Iceberg?
I imagine it might have had something to do with timing and the state of Iceberg at the time you started this effort?
Personally I found it very hard to reason about and thought that Clickhouse's strategy for managing columnar data much more reasonable.
The existing version of pg_analytics uses delta-rs to manage Parquet files stored within Postgres. In the future, we plan on integrating external object stores. This means that you'll be able to query any Delta Lake directly from Postgres.
Iceberg support will come later, once the Rust implementation of Iceberg matures.
Case in point regarding OLAP in particular, I am currently trying to solve a problem where I have a high number of categorical dimensions, and I want to perform a “count distinct” over any combination of dimensions, grouped by any other combination of dimensions, filtered by specific values in each dimension. E.g., count(distinct a.col1, b.col2), count(distinct a.col1), count(distinct b.col3) from table a join table b using (id) group by a.col4, b.col7.
Sounds obscure when I word it that way, but this is actually a pretty “generic” problem that appears whenever you want to filter and count the number of distinct property combinations that occur within a fact dataset of transactions or events that has been joined with other dimensional datasets. A naive implementation is exorbitantly expensive (and impractical) if you have to join many large tables before grouping and performing count distinct.
However, this specific problem manifests in various equivalent forms mathematically: model counting of boolean expressions, low rank factorization of sparse high dimensional boolean tensors (each row in your transaction dataset corresponds to a value of “true” in a sparse tensor with dimensions indexed by the values of your columns), minimal hypergraph covering set, etc.
Is there a database already out there that’s optimized for this fairly common business problem? Maybe...? I searched for a while but couldn’t easily separate the startup database hype from the actual capabilities of a particular offering. Plus, even if the ideal “hypergraph counting database” exists, it’s not like my company is just going to replace its standard cloud SQL platform that serves as the backbone of our entire product with a niche and fragile experimental database with questionable long-term support. It’s much easier to just translate one of the latest tensor factoring research papers into a Python script, plop that into the data processing pipeline, and output the simple factored form of the transactions dataset into a new table that can be easily queried in the standard way.
The custom types/indexes introduced by PostGIS won't work with deltalake tables. Even if it were possible, the benefits of using deltalake tables to execute geospatical queries are unclear, since PostGIS indexes are already optimized for this task.
But, what ParadeDB enables is for geospatial tables and deltalake tables to exist within the same database.
It's a lot more mature than people think :)
We also ship our individual extensions (pg_analytics and pg_bm25) pre-compiled as .deb under our GitHub Releases, and as a fully complete product, ParadeDB, as a Dockerfile and Helm chart
This would allow easier reporting within a very active db without too much bother.
Collations at the column/operation level are not yet supported but we're working on it.
In philosophy, we believe in playing into the ecosystem. We use DataFusion to avoid needing to write a vectorized query engine, Arrow to avoid needing to build an in-memory representation, and Parquet to avoid needing to build columnar disk storage. Citus columnar/Hydra appear to be working from first principles directly within Postgres, storing data in Postgres blocks, and writing vectorized execution operator by operator