Using PostgreSQL as a Data Warehouse
narrator.ai
narrator.ai
- Instead of updating tables build their replacements under a different name then rename them. This makes updating heavy-to-compute table instant. Works even for schemas: rebuild a schema as schemaname_next rename the current to schemaname_old then rename schemaname_next to schemaname.
- Keep all the source data raw and disable WAL, you don't need it for ETL.
- Set memory limitis high.
And lots of other good tips for doing ETL/DW in postgres. It's here: https://www.youtube.com/watch?v=whwNi21jAm4
I really appreciate having data in postgres. It's often easy to think that a specialised DW tool will solve all your problems, but that often fails to consider things like:
- Developer experience. Postgres runs very easily on a local machine, more specialized solutions often don't or are tricky to setup.
- Learning another tool costs time. A developer can learn postgres really well in the time it takes them to figure out how to use several more specialised tools. And many devs already know postgres because it's pretty much the default DB nowadays.
- Analytics queries often don't need to run at warp speed. Bigquery might give you the answer in a second but if postgres does it in a minute and it's a weekly report, who cares?
- Postgres is boring and has been around for many years now, it will probably still be here in 10 years so time spent learning it is time well spent. More niche systems will probably be superseded by fancier, faster replacements.
I would go so far as to say don't necessarily need to split out your DW from your prod DB in every case. As soon as you start splitting out a DW to a separate server you need some way to keep it in sync, so you'll probably end up duplicating some business logic for a report, maintaining some ingestion app, shuffling data around S3 or whatever. Keeping your analytics in your prod DB (or just a snapshot of yesterdays DB) is often good enough and means you will be more likely to avoid gnarly business-rules going out of sync between your app and your DW.
https://github.com/mara/mara-pipelines
And additional packages can be found at the Mara org:
Data warehouse workloads tend to be very IO intensive and could be highly disruptive to a production db. ETL is hard but it's a price to pay to isolate these two very different workloads.
DW lets you gain insight from your available data, whether there's a lot of data, or not.
DW with postgres is probably the right choice if your data isn't petabytes level.
DW has always described the storage solution not the workload.
More than once I've pointed at such teams as an explainer when pitching a service-team approach to standing up a new application. (Caveat, lest it backfire if you try the same, check who you're talking to first; they're not always fans of "the BI gang")
Not to mention that they have hardly any monitoring/alerting, which means that the data in the DW is not reliable, which can be really frustrating for users.
Overall I get the sense that these are unskilled folks who oversold what they could do but ultimately have very little clue what’s really going on.
Indeed, you can even guarantee backwards compatibility by versioning views you expose as APIs to other applications, and having each depending application create a view in its own schema that references the view--in Postgres, this means that only backwards-compatible migrations to the original view will go through! So when you find you can't make a change you need backwards-compatibly, you know it's time to make a new version, and you can keep the two versions of the API in sync by recomputing data in the new view to the old one until the depending applications can upgrade :)
Well, don’t do that then. I think the parent comment suggested that it was an anti-pattern.
Martin Fowler seems to agree: https://martinfowler.com/bliki/IntegrationDatabase.html
The scenario presented on the website (a separate organization controlling the database schema and application developers having to negotiate with them to get what they want in) is not what I'm proposing here anyway, and seems to be deliberately conflating the use of the database as an API and some sort of organizational pattern that is mostly orthogonal to that. There is no reason applications can't have their own private schemas that are managed by the application developers, who then choose what data to expose as public (and what as private), while still sharing a database with other developers. And for applications where database operations are particularly performance sensitive or safety critical, there's also no reason why not exposing your schema to other applications implies that application developers are allowed to have full control over it. I've been in organizations at various points in time that do both of these things.
Of course, there may be other restrictions that make this a bad idea (such as security restrictions that for whatever reason can't be enforced within the database), but for consumption for BI processes this is not really applicable anyway.
IMO it's not fruitful to throw blame on teams or call them unskilled. Usually, everyone is trying their best to do their job, and things that make this not happen should be seen as "how can we improve/fix this" - better communication is almost always the solution.
Database level integration is tempting at first since it’s zero effort up front, but it usually ends badly. It’s well worth it to make the BI team go through an API, even spend the effort on a well-designed one to give them a nice experience doing so.
In an ideal world yes. But sometimes, the engineering team wants nothing to do with the data itself, and so the BI team is often left to do what they can with what they have in order to do their job at all.
Being a data engineer, both sides of the coin often fall on me, and so I understand the difficulties on both side. Almost always (anecdotally, so take with grain of salt), the problem lies with upper management because they aren't willing to invest the time to actually build a proper infrastructure in place; the focus is on building features to satisfy clients, without much thought on how this can affect the overall internal system, and so both sides have to rush to do things without thinking how it affects the other side. However, it likely won't change because at the end of the day, the client is happy - whether or not this is a "unfortunate reality" of business/tech, is a discussion for another day.
It's probably more a symptom of upper management valuing neither the core team nor the BI team, thinking software writes itself, and business insights come for free with the MBA. We're really in the same boat here. All the more reason for core team to make nice hooks for the BI team. They won't be appreciated by management, they probably won't be appreciated by the customers, but at least they can be appreciated by each other.
But yeah, we could complain on and on about all the problems of shortsighted management... it's a tale as old as time.
A data warehouse is not a capability, it's a data system. Company A with a warehouse and company B without do not necessarily differ in any capabilities.
I've used a slight variant of this in the past: I'll have a table (e.g. my_schema_new.my_table) that gets updated by an ETL pipeline. I'll then also have a matview (e.g. my_schema.my_table) that's just SELECT * FROM my_schema_new.my_table. As long as I can give this matview a UNIQUE index, I can then REFRESH MATERIALIZED VIEW my_schema.my_table CONCURRENTLY for zero-downtime updates.
(Of course, if you're using this in a non-interactive data warehouse context, you might not care so much about zero-downtime updates; this is more for application-facing views. The REFRESH ... CONCURRENTLY pattern can be faster for incremental updates, but often struggles with larger changes to the underlying table, as it's essentially applying a diff between the two versions. Also, it only works in cases where your users can tolerate data that's as stale as the scheduled time between REFRESHes.)
This works really well, although a few minor warnings.
If you run backups, they will block DDL changes like creation and rename of a table. Now this might block all reads on your table until backup completes. Make sure you do this on a schema that isn’t being backed up or do other workarounds.
Other than that, this trick is really nice for instant replace and can also reduce disk space churn / fragmentation.
I like schema swaps too for the tables rather than renames but the pattern is the same.
isn’t this what materialized views are for?
Refreshing a materialized view blocks selects from it though, rendering it inaccessible while the refresh is taking place. There is the CONCURRENTLY option but that may be more costly and requires a unique index on the view.
There's also the issue of rebuilding several dependent views. Say you have 3 dependent views that are all being shown in some dashboard: 1, 2 and 3. You refresh view 1 OK but view 2 fails for some reason. Now view 1 and view 3 are out of sync and view 2 is broken.
Encapsulating the 3 dependent views in a schema and only doing the name swap once all steps are OK prevents you from ending up in this state. If view 2 breaks you wont do the name swap so you'll be serving stale data, which is probably better than no data for 2 and out-of-sync data for 1 and 3, until you can fix it.
It is a hybrid row column store with excellent compression and performance.
It would be interesting to see how it compares if narrator would try it out. Benchmarks would be cool.
One very neat feature I am enamored by is “continuous aggregates”. These are materialized views that auto-update as you change the fact table.
Continuous aggregates are a great idea. InfluxDB had “continuous queries” (but the implementation of influx generally is not so neat), and firebolt has “aggregate indexes” which are much the same thing.
I think all olap dbs will eventually have them as staple, and that they will trickle down into oltp too.
My uninformed assumption is if I do a group by over all rows in a table that they may not perform better.
I'll look into their continuous aggregates -- that could be one way to get around the cost of aggregating everything if it's done incrementally.
Fetching data becomes much, much slower once your index is too large to fit in main memory. TimescaleDB segments the indexes into chunks (the “hypertables”) and makes sure these chunks are all “small enough”.
This alleviated further by having the data sequential by time; inserting new data does not need to alter older index chunks, which is what makes inserts fast.
I can imagine that if that’s not the case and your inserts are altering “older” chunks so data needs to move between lots of chunks could make the database prohibitively slow.
I do not have specific experience with TimescaleDB, but I have some experience scaling PostgreSQL directly and with Citus (which is similar, but not the same). But depending on the nuances of your use case, I can envision a number of scaling strategies in vanilla Postgres to handle your use case. A lot of what Timescale and Citus does is abstract some of those strategies and extend them. Which is just a vague way of me saying: I think I could probably come up with a scheme in vanilla Postgres to support your use case, and since Timescale/Citus makes those strategies even easier/better I am fairly confident they would also handle that use case.
As an example I currently have a table in my current Citus schema that is sharded by column hash (e.g. "type" enumerator) and further partitioned by time. The first part (hash based sharding) seems possibly sufficient for your use case.
Beyond the most simple applications in that domain though, there are more exotic options available to both Timescale and Citus that could come into play. For example, I know Citus recently incorporated some of their work on the cstore_fdw into Citus "natively" to allow columnar storage tables directly:
https://www.citusdata.com/blog/2021/03/05/citus-10-release-o...
And being HTAP, timescale ought do better than a classic pure column store on upserts and non-appending inserts too.
Of course if your table is big and the keys are unordered you still might get excellent performance if your access pattern is read heavy.
Of course you can still mix in classic Postgres row-based tables etc. Timescale just gives you a column store choice for each table.
Computing aggregates against contiguous values in memory or disk for a single column will always be faster than reading records with differing value offsets/alignments. These operations can benefit from SIMD and other hardware optimizations you don't get with row-aligned data.
Upgrading to 12 (or 13) seems like the better option here, whenever you're able to do so. The improvements are very much worth it.
Fortunately there's a workaround - using the `with .. as not materialized ()` hint sped up my query 100x.
In my personal experience: the organisation doesn't invest in a storage solution that offers snapshots/COW (e.g. ZFS or SAN or whatever). Then they wait to upgrade until their disks reach a usage that the upgrade has to be done in-place. Then they become like rabbits in the headlights and never upgrade.
That said, given the performance implications, if someone wants to use PG as a warehouse upgrading to 12 is a no-brainer.
- Citus data, not Citrus data.
- In tables, column order matters. Order const-width non-null columns before all other columns (best so that there's no unnecessary alignment padding). Then maybe some low null fraction fixed-width columns, then ordered by query access frequency. A lot of time can be spent extracting tuple values, and you can save a lot of time using the offset caches by correctly tetris-ing your columns. Note that dropped columns are set to null, so a table rewrite (by redefining your table) may be in order for 100% performance.
Adding or removing a column in this system does not require postgres to rewrite the whole table, which means that old data can stay in the table effectively forever, as long as there are no other rewrite-required DDLs performed. Additionally, columns can be updated to SET NOT NULL / DROP NOT NULL, further complicating the whole system you're trying to optimize.
Eventually postgresql might support some form of table column reordering that permanently optimizes out deleted columns, but I think it's unlikely to happen anytime soon. There are lower hanging fruits on the tree; altering existing table definitions in a backwards-incompatible manner is a lot of effort for likely very little gain.
As for column packing: Maybe this can be implemented for CREATE TABLE with an option, but in the current transactional DDL framework this cannot be implemented for ADD COLUMN, because we can't reorder columns.
Where are my DuckDB people? (Think SQLite for OLAP workloads.)
If rows are around 1kB, then full-dataset queries over 100m rows will cost less than $0.5 each on BQ -- less if only a few columns have to be scanned. Storage costs will be negligible and, unlike pg, setup time will be pretty much nil.
It feels expensive, but then again running and maintaining a big data platform is inherently very expensive. If you consider the fully loaded cost of a single big data engineer to maintain such a platform to be ~$250,000 (likely an underestimate if you want similar performance characteristics to BigQuery [disclaimer: never used it myself, but I assume its performance is near-unbeatable]), that'd be ~500,000, 100 GB queries. Which makes a $0.50 query feel reasonable relatively speaking. GCP also sells dedicated "slots" (as they call it) which as I understand it is an abstraction over a CPU. If you buy said slots, the marginal cost of queries is $0, but you may be subject to queueing. No idea what a "slot" actually represents however.
It's not a datastore to power a crud app, or anything requiring frequent queries, but it's a great place to stash gobs of logs that you may need to query at some point. Or it's great for serverless batch workloads and is often cheaper in both time and money than firing up spark clusters or something similar to do the work.
Quite frankly, it's awesome. But sure, they do use it as a tool for lock-in, and for some cases it would be prohibitively expensive.
I find 50c to read 100GB from disk, do useful work on it (including running javascript code or ML models if you are so inclined) and returning a result in seconds... pretty damn incredible.
There's a lot more to value than the price.
No server fees and no fees when autoscaling up for heavy computations.
Yugabyte appears to an application as essentially Postgres 11.2 with all psql features (even the row / column level security mechanisms) but, apparently, handles replication and sharding automagically (is it DHT, similar to Cassandra)?
Are there any serious technical comparisons with some conclusions, without marketing bs?
Yugabyte uses actual Postgres code for the query parsing top layer and then translates that into to operates on it's key/value store which handles replication and distributed. CockroachDB is similar but has built everything from scratch in Go. There are similar examples like TiDB and Vitesse for MySQL as well.
Cockroach is a no go for me because of the licensing restrictions. For example, password authentication only in the BSL licensed code. Backups in enterprise.
https://www.cockroachlabs.com/blog/distributed-backup-restor... https://www.cockroachlabs.com/docs/v20.2/backup.html
Well, that’s not really correct is it. ClickHouse for one definitely has them as Snowflake the last time I used it.
This is a lot of work to go through to avoid using the right tool for the job. Just use something like ClickHouse or even DuckDB and reap the benefits of better performance with less caveats.
But, as things grow, you tend to run into problems with analytical loads adversely impacting, and even knocking over, the production system by locking resources and causing timeouts. Especially if you're allowing analysts to run ad-hoc analytical queries.
Long story short, yes, resiliency is expensive, but it's not always more expensive than not having resiliency.
Use a single database for transational and analytical workloads.
Replicate your transational database as is for analytical workloads.
Remodel your data and replicate to the exact same technology stack.
Remodel your data and replicate into specialized tools for analytics.
I have never seen anybody that actually needs the last one. But the largest environments I've looked are government databases with a few thousands of people working on (there are bigger envs out there).
Synchronising the data can be done in stages as well. Daily loads are pretty straightforward and usually small/straightforward enough to get going quickly and maintain. As your needs increase you up the frequency or start employing more sophisticated methods.
Snowflake does not have indexes, and ClickHouse indexes are what they call "data skipping indexes". BigQuery, Redshift, Netezza, and Vertica also do not have support for indexes.
Primary key is range index - to quickly locate records. Secondary indexes are data skipping indexes - to quickly skip blocks of data.
That’s still an index though isn’t it? Might work slightly differently, but the purpose is still the same.
What I think you meant to say was that Snowflake does not offer user-manageable indexes, not that they don't use indexes at all.
On top of that BRIN (block range) indexes are usually used to capture the value of sorting by pruning I/O. I don't see these mentioned here -- they seem like a good open-source version of this idea.
To my mind the big differences for data warehouses are the following:
(a) Table scans are relatively cheap thanks to columnar structure and compression. It's much more important to tune compression than indexes. If you can reduce stored data size by 10^3, you don't need an index. That's exactly the opposite of row stores like MySQL.
(b) Data warehouses don't use indexes to maintain referential integrity, because it's not something they really care about in the first place.
My DW experience is with ClickHouse, but I think it illustrates a lot of the principles.
The fact is that for large datasets scans on denormalized fact tables parallelize well, which means you can (a) offer stable performance and (b) scale more efficiently. This is important for use cases like web analytics, where users play around with different dimensions and measures but still expect consistent response. Note also, dimensions for things like Year, Month, Week, and the like compress absurdly well. It is often way faster to scan these values than to join them.
Also, under "Reasons not to use indexes", #3 says: "Indexes add additional cost on every insert / update". Yes, but then data warehouses usually aren't continually updated during the workday. It's a one-time hit sometime during the night, during your ETL run. (Not a Pg DBA, but presumably you can shut of index updating during data load, and then run it separately afterwards for higher performance?)
This isn't really what you're saying, but citus [1] ships a distributed Postgres. A lot of the things they improve would help massively with analytical workloads actually.
I'd assume yes, but I haven't personally used you guys. I'm just aware of you and broadly how you scale Postgres.
I ran a pg data warehouse in the 8.x and 9.x days with about 20TB of data and it performed great.
What we personally do in practice is put everything into a single time-series table with 11 columns. [1]
* partitioning your fact/aggregate tables (which was mentioned)
* rolling up old data reducing granularity as data ages can help with record count and overall db size
* PG triggers and stored procedures can be used to manage slowly changing dimensions
* the hstore and json column types are super useful for implementing quazi-nosql storage along side traditional relational storage
* window functions and CTEs (in 12) are great for writing analytical style queries
* Implementing incremental loads with a staging/buffer table and possibly with plpgsql can really make connecting it all easier
Not saying I didn't enjoy the article. It always makes me happy to see people realizing how suited PG can be to analytical workflows especially for small to medium workloads which represents most of what people want to do.
Typically only the more recent data was analyzed.
We retained all source data sou that reprocessing could happen when adding facts or dimensions.
Analysis, scheduled or ad hoc, only ever happened against the data cubes. Their partial aggregate nature along with partitioning was part of how the performance was so good.
Ugh, this is nightmare. I wish they would come up with a better system than forcing this on users.
For data warehouses inserts can happen in bulk on a regular cadence. In that case it can help to vacuum right after. I'm not sure if it has a huge impact in practice.
Also, do you have any numbers on how PG performs once it's configured?
Foreign data wrappers are another thing that might be compelling — dunno if Snowflake has an equivalent.
I don't have any numbers but PG has served me fine for basic pseudo-warehousing. Relative to real solutions, it's pretty bad at truly columnar workloads: scans across a small number of columns in wide tables. The "Use Fewer Columns" advice FTA is solid. This hasn't been a deal breaker though. Analytical query time has been fine if not great up to low tens of GB table size, beyond that it gets rough.
In my own testing PG performed very similarly to a 'real' warehouse. It's hard to measure because I didn't have the same datasets across several warehouses. Maybe in the future I'll try running something against a few to see.
And often having a homogenous database stack is a plus. If your production systems are all MySQL, then trying to get away with using MySQL for analytics too is a smart move etc.
I’ve seen so many tiddly data warehouses. Most companies don’t need web scale, and they overbuild and over complicate when they could be running on a simpler homogenous stack etc.
Cost is the most common reason I’ve seen. An RDS instance is about 1/2 the cost of Redshift per CPU and then Snowflake is slightly more expensive than Redshift (often worth the extra $$).
Also, if you’re dealing with less than 10GB of data the difference in performance will be barely noticeable, so at modest scale the most cost effective solution is RDS.
Because you don't want to learn and maintain a new kind of software in a small organisation with limiter resources, for example. Each new language, framework, database and other kind of tool adds a lot of cognitive load for everyone involved.
Customers who demand / are legally obliged to ensure their data does not leave their territory, is one big reason.
Snowflake give you some options, but if you get a big contract with a customer in another region your entire analytics platform is unavailable to that customer.
The biggest thing holding me back from using pg as a data warehouse is RDS not having support for instances with ephemeral drives / cost for PIOPS.
I need 100s of thousands of iops not thousands.
Thanks Cedric for sharing your experience with using PG for data-warehousing <3
https://www.toolbox.com/tech/data-management/blogs/2-petabyt...
Since 2008 improvements in parallel query execution (and numerous other improvements) in the core project plus the availability of forks/extensions which abstract and/or modify various bits for improving scalability (see Citus and Timescale) it's never been easier to scale Postgres to some truly staggering heights.
While I wouldn't want to speak in absolutes, there are very few applications where I think Postgres wouldn't be a viable choice as a data warehouse.
Emphasis on warehouse as I wouldn't want to suggest Postgres as an ideal candidate to be a data lake. The difference between them for me being whether or not the data is structured/processed. Similar in definition to this article:
https://medium.com/@distillerytech/data-warehouse-vs-data-la...
Personally, I have experience scaling core PostgreSQL (9.4) to handle ingestion of monitoring data for web servers to the tune of 2-3 terabytes a day. Not the grandest of scales, but enough to have seen a few bumps along the way...and, for what it's worth, I think it is surprisingly easy to scale.
I wouldn't want to sign up to scale Postgres to handle exabyte data loads, but single digit petabytes? Sure.
https://techcommunity.microsoft.com/t5/azure-database-for-po...
And at petabyte-scale, I personally think it qualifies as a data warehouse.
We use it all the time and their aggreations can get really advanced and perform well to the level where we run most analytics on demand. Sure we're not pushing to “big data” levels, max a few 100k records, but in reality I believe that's the average size of the majority of business datasets (Just an estimate, I have nothing to back this up)
Been building with it for 5 years now and it's been a breeze. Especially with Atlas. I think our team has not spent more that 3 days in total on DB dev ops. And with Atlas Lucene text search and data lakes, querying data from S3. What's not to love.
I haven't tried Atlas myself (and why would I? - I try and avoid lock in) but since Postgres supports json column types which has been my go to instead of Mongo & it has been an absolute breeze. Especially since it can be indexed and scanned with postgres sql.
Excel on a laptop is also a viable option at that scale
However, once you get to a 600 M row dataset, you will likely encounter significant performance/cost issues with a noSQL database.