Say you've got a giant table for of emails, and there's a bunch of indexed fields (people maybe want to filter their emails by tags, read/unread, etc.). If you keep them in one giant table, you need to keep those giant indexes fully in memory to get good performance. BUT, you know 99% of requests are for recent emails, people only occasionally look at old ones. So you partition your table by timestamp, and instead of running queries like `select * from email where status = 'unread' order by timestamp desc limit 50`, you run queries like `select * from email where status = 'unread' AND TIMESTAMP >= <pretty-recent> order by timestamp desc limit 50`. This only hits the first partition, and almost always returns the 50 emails you need - if it doesn't, you run subsequent queries to hit older partitions.
But now, say you partition by month, and keep 2 years of data, almost all queries are being serviced by a partition containing 1/24th of the total data, so keeping indexes for that 1 partition in memory is WAY easier than keeping indexes in memory for the entire giant table. When you have to query old partitions, it's slower and there's still some swapping happening, but that's rare anyways. So you can get away with having way less memory for your db, without performance degradation on most queries.
We use it to store 150M records a day, and be performant enough to have a customer facing analytics dashboard querying against it for real-time results. Pretty amazing approach.
For those who don't know, TimescaleDB is a extension to Postgres that makes it easy to use and scale for time-series data. It works with both PG9.6 and PG10.
Some folks might be interested in two recent blog posts we wrote about PG10 partitioning:
1. Technical write-up about PG10 partitioning: https://blog.timescale.com/scaling-partitioning-data-postgre...
2. Why TimescaleDB's partitioning is easier-to-use and more performant than PG10's native partitioning for time-series data: https://blog.timescale.com/time-series-data-postgresql-10-vs...
For bulk operations like (auto-)VACUUM, CLUSTER or dumping/restoring data it can make a big difference, though.
There are realistic scenarios where each row can end up costing you an 8kb page fetch, reducing query speed by ~250x for narrow tables as you're bound by memory bandwidth, or worse, disk bandwidth.
There's also an additional performance benefit in being able to skip the index scan when querying by month, and just sequentially scanning the entire partition which is usually 2-3x faster to access the same amount of data.
That makes sense.
> There are realistic scenarios where each row can end up costing you an 8kb page fetch, reducing query speed by ~250x for narrow tables as you're bound by memory bandwidth, or worse, disk bandwidth.
I've actually had a similar scenario. CLUSTERing by the index solved that. But of course running CLUSTER on a huge table is very slow. If you can partition the data so you only have to cluster one of the partitions, that's a huge win.
> There's also an additional performance benefit in being able to skip the index scan when querying by month, and just sequentially scanning the entire partition which is usually 2-3x faster to access the same amount of data.
Didn't think of that.
Thanks for the clarifications.
The declarative partitioning (added in PostgreSQL 10) is also transparent for the optimizer, i.e. it can understand how the data is routed to partitions, and can leverage it while planning/executing queries. For example when a table is partitioned on "a" and your query does "GROUP BY a" then in some cases the database can do the aggregation per partition - which should be more efficient in general (smaller hash tables, less data to sort, ...). Or when joining tables partitioned in the same way, it may be possible to do by joining the matching partitions (again, more efficient).
Obviously, PostgreSQL 10 only introduced the "core" declarative partitioning, and many such goodies are currently being worked on - either for PostgreSQL 11 or following version(s).
https://www.postgresql.org/docs/10/static/rules-materialized...
Once your indices are larger than than your available memory, write and read performance plummets. This obviously only applies to fairly large tables, but in our case it will be a godsend.
Partitioning looks to be going the other way. Create a big table, partition for convenience, drop partitions easily when you are ready to dispose of them. I can definitely see some cases where this would be super handy.
Jokes aside when you need to operate reporting on a multi-billion row scale which only needs to look at a very specific subset of your data I could see the use of partitioning. It's the transparent cousin to "copy it into it's own table".
The biggest win I see though is fast deletes.
Partitioning is not like copying to another table because the partition is already its own table. It saves the entire expensive operation of copying that data and you get the additional benefit that it stays current.
Copy data out of tables, put them into formats conducive to analysis (document stores, relational, column oriented stores, etc). You can then tie everything together in your language or from something like Presto.
With well designed partitioning scheme you can simply drop a partition, and you're done. Much faster/cheaper.
But this is just one one benefit of partitioning - I've mentioned the possible benefits for planning elsewhere in this thread.