Skip to content

Stabilize metadata sync - part1 - #8878

Open
Onur Tirtir (onurctirtir) wants to merge 14 commits into
mainfrom
ms-scoped
Open

Onur Tirtir (onurctirtir) wants to merge 14 commits into
mainfrom
ms-scoped

Conversation

@onurctirtir

@onurctirtir Onur Tirtir (onurctirtir) commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

Closes #7264.

DESCRIPTION: Reduces memory usage during metadata sync by avoiding unnecessary PostgreSQL and Citus cache growth (controlled by citus.metadata_sync_cache_flush_interval, set to 1000 by default).
DESCRIPTION: Reduces the number of network roundtrips required to sync Citus metadata tables during metadata sync (controlled by citus.metadata_sync_set_batch_size, set to 1000 by default).
DESCRIPTION: Ensures proper cleanup of shell tables left over from previous failed metadata sync runs.
DESCRIPTION: Logs metadata sync progress to the logs.
DESCRIPTION: Prevents syncing wrong distribution column types in colocation metadata to workers when the type is a custom type or is in a schema whose name needs quoting.
DESCRIPTION: Fixes a duplicate key error on pg_dist_colocation while syncing metadata to workers when a user type has the same name as the built-in type of the distribution column.

For easier reviews, each commit can be reviewed separately.
Each commit's message also provides more details about it.

There are also couple of big improvements scoped out from this phase, see issue #8879.
This PR proposes the safest changes we can merge, backport to older releases, and enable by default.

Also, seems we're not doing the best we can do to reduce the peak memory usage in "collect" mod, but this has already been the case and collect mode is there only for testing purposes.

@codecov

codecov Bot commented Sep 25, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 94.50867% with 19 lines in your changes missing coverage. Please review.
✅ Project coverage is 88.81%. Comparing base (bfb8600) to head (d936926).

Additional details and impacted files
@@            Coverage Diff             @@
##             main    #8878      +/-   ##
==========================================
+ Coverage   88.77%   88.81%   +0.04%     
==========================================
  Files         290      290              
  Lines       65149    65408     +259     
  Branches     8225     8249      +24     
==========================================
+ Hits        57835    58094     +259     
+ Misses       4938     4933       -5     
- Partials     2376     2381       +5     
🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

FlushDistTableCache() could dereference not-found (negative) cache entries
while flushing the distributed-table cache during metadata sync. Guard
against negative entries so the periodic cache flush is safe.
On clusters with millions of distributed tables the coordinator buffered the
whole snapshot and never released cached distributed-table entries, driving it
toward OOM. Periodically flush the distributed-table cache during metadata
sync to keep coordinator memory bounded, controlled by the new
citus.metadata_sync_cache_flush_interval GUC (default 1000).
pg_dist_object records were synced one row per round-trip. Accumulate them into
multi-row batches flushed at citus.metadata_sync_set_batch_size (new GUC),
collapsing N round-trips to N / batch_size on clusters with many distributed
objects.
Follow the batching optimization done for other Citus metadata tables
while inserting into colocation / tenant-schema metadata too.
…sync

Follow the batching optimization done for pg_dist_object while
inserting into pg_dist_partition/shard/placement metadata too.
…the table during metadata sync

Previously, as we were creating shell tables in a different remote
xact than the one we insert the pg_dist_partition row for the
Citus tables, having shell tables on a worker that don't have
relevant pg_dist_partition records, if the metadata syncing fails
during the process.

And as we rely on worker-local pg_dist_partition entries while
cleaning-up worker shell tables in the early stages of the metadata
sync, this was causing leaving dangling shell tables on the worker.
As a result, later in the metadata sync, we were experincing
"table already exists" kind of errors while trying to (re)create a
shell table.
Add metadata_sync_batching, which syncs 110 distributed schemas with
2 tables and a foreign key each to a node, with
citus.metadata_sync_cache_flush_interval and
citus.metadata_sync_set_batch_size set to 50, in both transactional
and nontransactional sync modes.

Add metadata_sync_batching_edge_cases, which covers flush interval 1
with a negative Citus table cache entry, flush interval 0, batch size
1, and a mix of object kinds in the same batch.
When sending a colocation group to a node, the coordinator sent the
distribution column type's schema name through quote_identifier(). The
node then compared that string with pg_namespace.nspname. For a schema
like "Q S", the node compared '"Q S"' with 'Q S', so the lookup failed
and the node stored 0 as the distribution column type.

This happened in both places that send colocation groups:
ColocationGroupCreateCommand (used when a table is distributed) and
SendColocationMetadataCommands (used during metadata sync). Send the
plain schema name instead, like we already do for collation schemas.
GetRemoteTypeNamespace returned NULL for types in pg_catalog, and the
worker-side join used "typeschema IS NULL OR ..." for them. So a
pg_catalog type such as int4 or text matched every type with that name
in any schema. When a user type with the same name existed (for example
a domain named int4), the join returned two rows and the insert into
pg_dist_colocation failed with a duplicate key error.

Always send the schema name, including pg_catalog, and always match on
it. Add two shadowing domains to metadata_sync_batching_edge_cases to
cover this.
The worker looks up the distribution column type and collation of each
colocation group by name, and stores 0 when it cannot find them.
Metadata sync used to send the colocation groups before it created the
types and collations on the node. So after citus_add_node(), a colocation
group whose distribution column uses a custom type, an extension type or
a custom collation got type 0 or collation 0 on the new node.

Send the colocation groups right after the dependencies are created.
No worker-side code reads pg_dist_colocation while the dependencies and
their pg_dist_partition records are created, so this is safe.

This relies on the previous commit. Without it, a user type with the
same name as a pg_catalog type would now exist on the node when the
colocation groups are sent, and the insert would fail.

Remove the workaround in metadata_sync_batching_edge_cases that created
the types and the collation on worker_2 before the sync.
…reates the table during metadata sync

Pin the colocation id sequence in seclabel so that the colocation ids
logged in the metadata sync commands don't depend on earlier tests.

Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
"bounds how many objects' rows are packed into each such "
"statement. Larger values emit fewer, larger statements "
"(peak coordinator memory stays bounded by the batch, which "
"is reset after every flush); 1 restores one statement per "

@colm-mchugh Colm (colm-mchugh) Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Wondering if the GUC description could/should cover all metadata phases ? Given that the batch sizing controls colocation, tenant-schema and partition metadata, the set of metadata layers impacted is bigger than mentioned here, which may give the impression that the scope of this GUC is narrower than what it controls. Suggest changing the first sentence to something like:

Eligible per-object metadata layers (including pg_dist_colocation, pg_dist_schema, pg_dist_partition, 
pg_dist_shard, pg_dist_placement and pg_dist_object) are synced over the serial metadata connection by 
rendering rows into multi-row VALUES lists fed to set-based metadata statements.

*/
LogMetadataSyncPhaseBoundary("starting", "inter-table relationship");
SendInterTableRelationshipCommands(context);
LogMetadataSyncPhaseBoundary("finished", "inter-table relationship");

@colm-mchugh Colm (colm-mchugh) Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It looks like there's no regress tests that include this new logging of phase boundaries so how about a simple test whose expected output should include the starting and finished messages for metadata-sync phases, e.g.:

SET client_min_messages TO DEBUG1;
SELECT 1 FROM citus_add_node('localhost', :worker_2_port);
RESET client_min_messages;

(localhost,57637,t,"ALTER EXTENSION")
(localhost,57638,t,"ALTER EXTENSION")
(2 rows)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we add some tests for edge cases for GUC values - e.g. outside of the allowed values:

SET citus.metadata_sync_cache_flush_interval TO -1; -- expect rejection
SET citus.metadata_sync_set_batch_size TO 0;        -- expect rejection
SET citus.metadata_sync_set_batch_size TO 10001;    -- expect rejection

SET ROLE metadata_sync_non_superuser;
SET citus.metadata_sync_cache_flush_interval TO 1;  -- expect success
SET citus.metadata_sync_set_batch_size TO 1;        -- expect permission error
RESET ROLE;

@colm-mchugh Colm (colm-mchugh) left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Lgtm, few non-blocking comments on testing

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Metadatasync Progress Report

2 participants