Partition Large Tables

Partitioning divides a large table into smaller ones, so that a query scans only the data relevant to its conditions. This topic covers how to decide whether a table is worth partitioning, how to create and maintain partitioned tables, and which constraints apply to them. For the complete syntax, see CREATE TABLE and ALTER TABLE.

When to partition a table

The benefit of partitioning comes from scanning less data. When a query filters on the partition key, SynxDB reads only the relevant partitions and skips the rest, a process called partition pruning. Removing or archiving historical data by partition is also far faster than deleting rows.

Partition a table that meets these conditions:

  • The table is large. Its data should exceed the physical memory of a single segment.

  • Queries consistently filter on a particular column. If queries rarely reference the partition key, partitioning adds overhead without any benefit.

  • The data has a clear lifecycle. For example, you retain it by month and archive it a month at a time.

  • The data divides evenly. If most rows land in one partition, partitioning achieves nothing.

Important

Partitioning and distribution are two different things, and you must set both. Distribution determines which segment a row lands on. You specify it with DISTRIBUTED BY, and its goal is to balance data across segments. Partitioning determines which child table within that segment a row lands in. You specify it with PARTITION BY, and its goal is to reduce the data a query scans. The partition key and the distribution key are usually different columns.

Choose a partitioning method

SynxDB supports three partitioning methods:

Method

Use case

Partition key

RANGE

Continuous values such as dates or numbers. The most common choice.

One or more columns

LIST

Discrete values such as regions or categories.

A single column or expression only

HASH

No natural dividing line. Use it only to spread data evenly.

One or more columns

You also need to choose the granularity. Partitions that are too coarse prune little data; partitions that are too fine cost more in metadata and query planning than they save. For a log table retained for three years, 36 monthly partitions are usually a better choice than more than 1,000 daily ones.

Use multi-level partitioning, that is, partitioning within a partition, only when a single level cannot deliver the pruning you need. Multi-level partitioning multiplies the number of partitions, and the optimizer handles it differently from single-level partitioning. See Verify that partitioning takes effect.

Create a partitioned table

Creating a partitioned table takes two steps. Declare the root table with PARTITION BY, then create each partition with PARTITION OF. The root table itself holds no data.

Partition by range

Range boundaries are inclusive at the start and exclusive at the end.

CREATE TABLE sales (id int, sale_date date, amt numeric)
    DISTRIBUTED BY (id)
    PARTITION BY RANGE (sale_date);

CREATE TABLE sales_2025 PARTITION OF sales
    FOR VALUES FROM ('2025-01-01') TO ('2026-01-01');

CREATE TABLE sales_2026 PARTITION OF sales
    FOR VALUES FROM ('2026-01-01') TO ('2027-01-01');

Numeric ranges use the same syntax with numeric boundaries.

Partition by list

CREATE TABLE users (id int, region text)
    DISTRIBUTED BY (id)
    PARTITION BY LIST (region);

CREATE TABLE users_cn PARTITION OF users FOR VALUES IN ('cn');
CREATE TABLE users_us PARTITION OF users FOR VALUES IN ('us');

Partition by hash

Hash partitioning divides data by modulus and remainder. Set the modulus to the total number of partitions and number the remainders from 0.

CREATE TABLE events (id int, payload text)
    DISTRIBUTED BY (id)
    PARTITION BY HASH (id);

CREATE TABLE events_0 PARTITION OF events FOR VALUES WITH (MODULUS 2, REMAINDER 0);
CREATE TABLE events_1 PARTITION OF events FOR VALUES WITH (MODULUS 2, REMAINDER 1);

Create a multi-level partitioned table

To add a level, declare a partition as a partitioned table itself. The following example partitions by date, then by region.

CREATE TABLE logs (id int, log_date date, region text)
    DISTRIBUTED BY (id)
    PARTITION BY RANGE (log_date);

CREATE TABLE logs_2026 PARTITION OF logs
    FOR VALUES FROM ('2026-01-01') TO ('2027-01-01')
    PARTITION BY LIST (region);

CREATE TABLE logs_2026_cn PARTITION OF logs_2026 FOR VALUES IN ('cn');

Only the lowest level holds data. Intermediate levels and the root table are empty.

Partition an existing table

You cannot convert an existing ordinary table into a partitioned table in place. Create a partitioned table, then either attach the original table as one of its partitions or insert the data into the new table. For the attach approach, see Attach an existing table as a partition.

Choose a storage format for each partition

Partitions of the same table can use different storage formats. Partitions created afterwards inherit the format declared on the root table, and you can override it when creating an individual partition:

CREATE TABLE sales_2025 PARTITION OF sales
    FOR VALUES FROM ('2025-01-01') TO ('2026-01-01')
    WITH (appendonly=true, orientation=column);

This suits data that cools over time. Keep recent partitions in the heap format while they still receive updates, and switch historical partitions to a column-oriented AO or PAX format for better compression and scan performance. A table attached with ATTACH PARTITION can also use a different format. For the differences between formats, see Choose a Table Storage Model.

A partition can also be a foreign table, which brings external data in as a partition. When the external data source is unreachable, a query against the root table fails as a whole rather than skipping that partition.

Load data into a partitioned table

Insert into the root table. SynxDB routes each row to the matching partition based on the partition key value. If no partition can hold the row, the insert fails.

To avoid such failures, define a default partition that absorbs rows matching no other partition:

CREATE TABLE sales_default PARTITION OF sales DEFAULT;

Check the default partition regularly. Data accumulating there usually means a partition is missing, or that unexpected values have entered the partition key.

Verify that partitioning takes effect

Use EXPLAIN to confirm that a query scans only the partitions it needs. Single-level and multi-level partitioned tables produce different plan shapes, so read each accordingly.

GPORCA handles single-level partitioned tables. The plan contains Dynamic Seq Scan and states the number of partitions scanned:

 Gather Motion 2:1  (slice1; segments: 2)
   ->  Dynamic Seq Scan on sales
         Number of partitions to scan: 1 (out of 3)
         Filter: (sale_date = '2025-06-01'::date)
 Optimizer: GPORCA

The PostgreSQL optimizer handles multi-level partitioned tables. The plan gives no partition count and instead names the partition actually scanned:

 Gather Motion 2:1  (slice1; segments: 2)
   ->  Seq Scan on logs_2026_cn logs
         Filter: ((log_date = '2026-06-01'::date) AND (region = 'cn'::text))
 Optimizer: Postgres query optimizer

The second shape is expected and does not mean that partitioning failed. A specific partition name in the plan shows that pruning happened.

If the plan scans every partition, check whether the query filters on the partition key. Pruning cannot happen when the filter column differs from the partition key, or when a function wraps the partition key. For tuning queries against partitioned tables, see Optimize Query Performance with GPORCA and Execute Queries in Parallel.

View the partition design

To view the partition key of a root table:

SELECT pg_get_partkeydef('sales'::regclass);

To view the partitions and their boundaries:

SELECT c.relname, pg_get_expr(c.relpartbound, c.oid) AS bound
    FROM pg_class c
    JOIN pg_inherits i ON c.oid = i.inhrelid
    WHERE i.inhparent = 'sales'::regclass;

Maintain partitions

Partition maintenance is valuable because an operation on a whole partition works on that partition as a unit rather than row by row. To remove a month of history, drop the corresponding partition instead of running DELETE.

Add a partition

CREATE TABLE sales_2027 PARTITION OF sales
    FOR VALUES FROM ('2027-01-01') TO ('2028-01-01');

If a default partition exists, SynxDB checks it for rows that belong to the new partition, which makes this operation slow when the default partition holds a lot of data.

Attach an existing table as a partition

ALTER TABLE sales ATTACH PARTITION sales_2028
    FOR VALUES FROM ('2028-01-01') TO ('2029-01-01');

Rows already in the attached table become part of the partitioned table. Its column definitions and distribution key must match the root table. This is a common way to load and validate data in a standalone table before bringing it into the partitioned table.

Detach a partition

ALTER TABLE sales DETACH PARTITION sales_2025;

The detached partition becomes an ordinary table. Its data remains, but it no longer belongs to the partitioned table, and queries against the root table no longer reach it. This is a common way to archive historical data.

Truncate and drop a partition

To empty a partition while keeping the partition itself:

TRUNCATE sales_2025;

To remove the partition along with its data:

DROP TABLE sales_2025;

Rename a partition

ALTER TABLE sales_2025 RENAME TO sales_2025_archived;

Index a partitioned table

Create the index on the root table. SynxDB creates a matching index on every partition, and on any partition added later:

CREATE INDEX ON sales (sale_date);

Index choices for a partitioned table follow the same rules as for an ordinary table. See Create and Manage Indexes.

Use the classic partitioning syntax

Besides the syntax above, SynxDB supports the classic Greenplum partitioning syntax, which declares a batch of partitions at once through START, END, and EVERY:

CREATE TABLE sales_classic (id int, sale_date date)
    DISTRIBUTED BY (id)
    PARTITION BY RANGE (sale_date)
    (START ('2026-01-01'::date) END ('2027-01-01'::date) EVERY ('1 month'::interval));

This syntax mainly serves tables migrated from Greenplum. For new tables, use the syntax described earlier. The gp_max_partition_level parameter limits only the number of levels the classic syntax can create and has no effect on the newer syntax. For the full classic syntax and further examples, see CREATE TABLE and ALTER TABLE.

To create and drop partitions automatically on a time schedule, use the pg_partman extension. See Manage Partitioned Tables with pg_partman.

Best practices

  • Partition only genuinely large tables. On a small table, the planning overhead outweighs what pruning saves.

  • Choose the column that queries filter on most often as the partition key, usually a time column.

  • Do not use the same column as both the distribution key and the partition key, or data concentrates on a few segments.

  • Keep the total number of partitions in check. More partitions mean more planning overhead, and creating or altering the table takes noticeably longer. At a few thousand partitions, a single CREATE TABLE can take tens of seconds, and the cost grows faster than the partition count.

  • Avoid unnecessary levels. Do not add a second level when one level prunes well enough.

  • Define a default partition and check regularly whether data is accumulating in it.

  • Verify pruning with EXPLAIN once the design is in place. Do not assume it works.

Limitations

Partitioned tables have the following limitations:

  • You cannot partition a table that uses the DISTRIBUTED REPLICATED policy.

  • To create a unique or primary key constraint on a partitioned table, the partition key must not include expressions or function calls, and the constraint columns must include every partition key column as well as every distribution key column. This is because the individual indexes behind the constraint can enforce uniqueness only within their own partition; the partition structure itself has to guarantee that no duplicates exist across partitions.

  • An exclusion constraint cannot span a whole partitioned table. Create it on each partition individually.

  • You cannot mix temporary and permanent tables in the same partition hierarchy.

  • A partition cannot inherit from another table, and you cannot change its inheritance while it belongs to a partitioned table.

  • GPORCA does not handle multi-level partitioned tables. Such queries run through the PostgreSQL optimizer instead.

  • If any partition is a foreign table, TRUNCATE cannot run against the root table, and that foreign table partition cannot be emptied on its own either. Truncate the remaining partitions one by one instead.

  • gpbackup backs up the definition of a foreign table partition but not its data. The backup log reports this as Skipped data backup of N external/foreign table(s). The external data source maintains that data itself.