Skip to content

Stabilize metadata sync for large clusters #8879

Description

On top of #8878, there are also couple of big improvements we can do:

  1. Sending dependency commands in parallel.
    We could use the adaptive executor or invent a new pool-based executor to send a stream of DDL commands using parallel connections to the new worker.
    This requires tracking dependencies between the objects being created during the sync.

    A few gotchas:

    • We should not send long (very close to 64 characters) CREATE TABLE commands in parallel; otherwise, it is possible to get PK failures when Postgres inserts the short name it generates for the relation type record.
    • There might be other conditions that could result in deadlocks or similar PK failures under concurrency; this needs further exploration.
  2. Sending DROP TABLE commands to clean up the existing shell tables on the worker in parallel.
    We should first detach all partition relationships between existing shell tables on the new worker, and do the same for any foreign key relationships.
    Then we can send DROP TABLE commands using parallel connections as described above.

    One gotcha:

    • Not sure what happens if two tables we are dropping in parallel cascade to same object (directly or transiently). Would we get deadlocks?
  3. One other stabilization issue we see during metadata sync is that we sometimes get errors because of growth in the lock table on the coordinator.
    This happens because, while generating DDLs for various purposes, the helpers we use acquire AccessShareLock on some types of objects, such as shell tables, and we never release them.
    We could use subxacts while generating those kinds of DDL commands and release the subxact once the commands are generated; then we can send the commands as today.

    However, this would increase the chances of deadlocks. Even today, if the user is concurrently issuing DDLs against Citus tables during metadata sync, either the whole metadata sync or the user's DDL might fail because of a deadlock.
    However, if we release the locks we acquire while generating DDLs in the earlier phases of metadata sync, then later in the sync, if we need to acquire the same kind of lock, we'll have the same deadlock risk at that later stage of metadata sync.
    Today, if we can successfully acquire Citus table locks in the earlier phases while generating DDLs, then we don't have the same risk in the later stages at least.

    As of today, those are the code-paths in metadata sync that are leaking huge number of locks until the end of the metadata sync.

    • SendDependencyCreationCommands: Its call to GetAllDependencyCreateDDLCommands locks every distributed table it builds DDL for.

    • SendInterTableRelationshipCommands: Its calls to ShouldSyncTableMetadata, IsTableOwnedByExtension, and InterTableRelationshipOfRelationCommandList lock every table again when building the foreign key and partition attach commands.

    • SendDeletionCommandsForReplicatedTablePlacements: It calls DeleteAllReplicatedTablePlacementsFromNodeGroupViaMetadataContext, which takes LockShardDistributionMetadata(shardId, ExclusiveLock) for each deleted placement.

  4. Having mentioned today's known deadlock issues and the risk mentioned for future improvements 2 & 3, in general, during metadata sync, we should consider retrying DDL commands a few times for at least some SQL error codes, such as the one indicating a deadlock.


I'd start with 4, then 1, 2, and 3.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions