PostgreSQL 11 Partitioning Improvements
pgdash.io
pgdash.io
This set of changes in 11 is pretty great and makes it much more usable out of the box. For those not reading the article the short list of improvements:
- Updates move records across partitions
- A default/catch all partition
- Automatic index creation
- Foreign keys for partitions
- Unique indexes
- Hash partitioning (where as 10 was just time/range based)https://github.com/postgres/postgres/commit/3de241dba86f3dd0...
The relevant limitation is described by CREATE TABLE docs that say
Partitioned tables do not support <literal>EXCLUDE</literal> constraints
however, you can define these constraints on individual partitions.
Also, while it's possible to define <literal>PRIMARY KEY</literal>
constraints on partitioned tables, creating foreign keys that
reference a partitioned table is not yet supported.
(the wording is slightly broken in the commit, but got fixed later).* Dynamic repartitioning to balance query load/space/create new partitions when tables get over a certain size. Without this, when you get more than a few hundred TB of data, you'll end up with terrible hotspots where single partitions are much more loaded than others and theres nothing you can do about it.
* Query replanning: Query plans can sometimes have a lot of error. By replanning, one can start on another query plan if the existing plan is grossly wrong in its estimates.
* Lock Eliding instead of *Exclusive locks. Some major operations, such as schema changes, take locks which effectively stop the server responding to any other requests till the operation is done. Instead, these operations should work on a copy-on-write image of the database, repeatedly re-syncing with recently changed data until the final change necessary is small, and then atomically commit the change.
You can say that it should aim for ~100GB per partition, partitioned by time, and it'll aromatically add or remove partitions.
You can go even further with the time scale extension for postgres.
(Which I personally agree with, given that you almost never get these extension—or support for arbitrary extensions—on Postgres DBaaS hosts. Upstreaming a feature = making it available on these hosts.)
Providing first-class partitioning is something that extensions can't possibly provide, due to how many places this has to integrate with.
On the other hand, manipulating partitions (which is effectively automatic table creation) seems like the sort of magic I expect to only exist in extensions.
As mentioned in a different comment, it does sound like people are looking for round-robin load distribution rather than partitioning.
I'm sure this is one of the things people on pgsql-hackers will be talking about, but I wouldn't hold my breath for it to get into PG12. I do expect this to emerge in an extension first, and then getting the most useful bits into core (either in the form of built-in features, or hooks allowing extensions to be smarter).
[1]: https://static.googleusercontent.com/media/research.google.c..., Section 5, "and also splits tablets that have grown too large"
[2]: https://storage.googleapis.com/pub-tools-public-publication-..., Section 1, "it automatically migrates data across machines (even across datacenters) to balance load"
For other types of partitions, partition hotspots would be the fault of the DBA partitioning on a bad column. The DB could theoretically half partition ranges if a range partition is very large, but there are other things that are more important.
#2 is a nice performance feature which has been talked about before, but do any DB's actually have this yet? I suspect that it is quite tricky, and won't be universally beneficial.
Is #3 really a problem? I certainly wouldn't expect schema changes to perform well, and increasing complexity just to help in this area would seem silly.
Not the parent, but obviously they refer to the online schema migration problem. Without that online schema migrations are basically impossible. Or you get handed the burden of implementing or finding a tool that does it for you.
And yes, pt-online-schema-change[1] is such a tool, but for MySQL. MariaDB has this feature built-in with their `ALTER ONLINE TABLE` [2].
So something similar for Postgres would be godsent.
[1] https://www.percona.com/doc/percona-toolkit/LATEST/pt-online...
I have been using it exclusively now instead of pt-online-schema-change.
[1] https://dev.mysql.com/doc/refman/5.7/en/innodb-online-ddl.ht...
In the nice world of all data having an even load, sure...
But in the real world, you can easily get a handful of users in your "Users" table sending millions of requests per second, and in that case, you would really like to re-partition so that they don't happen to all end up on the same partition.
Implementation can be as simple as allowing splittable partitions, so that whenever load/size gets too high, a partition can be split in half, and half the records moved to a new host. The partition map is only a few kilobytes, so is globally shared/updated. There are no concurrency issues, because during the splitting process, either the old or new partition is responsible for each record, and both the old and new partition hosts can respond to a read or write query, either themselves or by forwarding it to the other host.
Whenever two neighbouring partitions both see low load/space usage, join them with the same method in reverse. By joining only neighbouring partitions, you can't suffer terrible fragmentation and blowing up the size of the global partition map.
A more immediate use case is partition-aware algorithms. For example joining equi-partitioned tables can be done at the partition level, and the smaller the partitions the faster the join (likely). Or partition-aware aggregation, when the partition key is included in the GROUP BY keys. And so on. In that case it makes sense to split the largest partitions to make them smaller.
While a can imagine a way to make hashed partitioning sub-divide buckets, it won't really be useful, as you can still have a single user blowing their bucket out of proportions with others being almost empty, which partitioning can never solve.
What you want is round robin distribution and parallel query execution, not partitioning. This of course only makes sense if the tables are on different disk systems/hosts or if all tables fit in memory.
Partitioning primarily deals with isolating related data into small pools so that queries do not need to touch full tables. The purpose of the buckets is to group related data so that a single partition may serve a query (instead of having to touch the entire table), thereby lowering the load from each individual query, not distributing all queries.
Load balancing deals with distributing all queries to ensure even node load, not increasing efficiency of queries (which it does not do at all). Data distribution leading to round-robin accesses across nodes/"partitions" is best. This is the exact opposite of partitioning, which tries to group all that data into a single partition.
Hash partitioning on serial or uuid might give you round-robin like behavior, but it's not really what partitioning is meant for.
We don't have (1) yet - at least not in core PostgreSQL, but this part should be doable in an extension until it gets into the core. We do have (2) to a large extent, although I'm sure there are gaps and room for improvement.
And when combined with ability to place partitions to other hots, that will be another step.
No one is claiming partitioning alone magically solves scalability, but IMHO it's an important building block in doing that (for a large class of applications).
SQL Server is getting it in 2017 https://docs.microsoft.com/en-us/sql/relational-databases/pe...
(I didn't downvote you, btw)
> The columns in the index definition should be a superset of the partition key columns. This means that the uniqueness is enforced locally to each partition, and you can’t use this to enforce uniqueness of alternate-primary-key columns.
Are there any workarounds for global indexes across partitions outside of the partition key? Say if you have a a table partitioned by (some_date) with a unique primary key (id), how would one speed up looks up that are purely by id? Having per-partition local indexes would require an index scan per partition.
I supposed you could have a separate/smaller (id, some_date) non-partitioned table (or maybe partitioned on HASH(id)) and then do the indirection yourself wherever you're doing "SELECT ... WHERE id = ..." but seems like a kludge. Any other workarounds?
That would solve nearly all the major remaining limitations.
That's probably the easier bit, I'd guess that adjusting a lot of the relevant in-memory structures to share relations, uniqueness checks across relations, adjusting the locking-model to deal with multiple underlying relations etc. is going to be more work.
I'm not sure why the bloom filter would require good spatial locality?
In addition to this, in 10.4 it was mentioned spider is working to have internal joins. This means it would be possible to have 1 table defined with the full data set, then a second table that just has id, datefield partition by id. When you select WHERE id = X, it will read from the small table partitioned by id, find the dates, then join against the larger table to pull in remaining columns/data. This isn't mentioned in the presentation linked, but it was brought up in discussion at the mariadb developers unconference.
[1] page 6: https://mariadb.org/wp-content/uploads/2018/03/Merging_patch... (note this is about pages for 10.4 but the VP functionality is in 10.3)
If so, you could potentially create 1. a table that maps ID ranges to partitions, and then 2. create a stored procedure that takes a set of IDs and returns the set of partitions you must scan to find them.
For applicable workloads, it means you get nearly perfect scaling up to the number of partitions you have.
PG11 will be huge.