From 2096c33cb65001ebcdfb271039d1b31371d98eaf Mon Sep 17 00:00:00 2001 From: Sai Dixith Date: Thu, 30 Jul 2026 15:41:14 +0530 Subject: [PATCH] Synchronizer: Add --skip-catalog-sync flag for DR-oriented syncs Allows skipping create/overwrite/remove of catalog objects on the target while still synchronizing catalog-roles and grants for catalogs that already exist there. Useful when target catalogs are pre-created with DR-specific storage locations that shouldn't be clobbered. --- .../sync/polaris/PolarisSynchronizer.java | 37 +- ...olarisSynchronizerSkipCatalogSyncTest.java | 335 ++++++++++++++++++ ...risSynchronizerSkipIcebergContentTest.java | 7 +- .../sync/polaris/SyncPolarisCommand.java | 13 +- 4 files changed, 388 insertions(+), 4 deletions(-) create mode 100644 polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipCatalogSyncTest.java diff --git a/polaris-synchronizer/api/src/main/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizer.java b/polaris-synchronizer/api/src/main/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizer.java index f7f96f80..d4360769 100644 --- a/polaris-synchronizer/api/src/main/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizer.java +++ b/polaris-synchronizer/api/src/main/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizer.java @@ -68,6 +68,7 @@ public class PolarisSynchronizer { private final SynchronizationReport report; private final boolean skipIcebergContent; + private final boolean skipCatalogSync; public PolarisSynchronizer( Logger clientLogger, @@ -78,7 +79,8 @@ public PolarisSynchronizer( ETagManager etagManager, boolean diffOnly, SynchronizationReport report, - boolean skipIcebergContent) { + boolean skipIcebergContent, + boolean skipCatalogSync) { this.clientLogger = clientLogger == null ? LoggerFactory.getLogger(PolarisSynchronizer.class) : clientLogger; this.haltOnFailure = haltOnFailure; @@ -89,6 +91,7 @@ public PolarisSynchronizer( this.diffOnly = diffOnly; this.report = report; this.skipIcebergContent = skipIcebergContent; + this.skipCatalogSync = skipCatalogSync; } /** @@ -632,7 +635,19 @@ public void syncCatalogs() { int syncsCompleted = 0; int totalSyncsToComplete = totalSyncsToComplete(catalogSyncPlan); + Set catalogNamesSkippedFromChildSync = new HashSet<>(); + for (Catalog catalog : catalogSyncPlan.entitiesToCreate()) { + if (skipCatalogSync) { + clientLogger.warn( + "Skipping creation of catalog {} because catalog synchronization is disabled. It does " + + "not exist on the target, so its catalog-roles and grants will not be synced either.", + catalog.getName()); + report.recordSuccess(EntityType.CATALOG, SyncOutcome.SKIPPED); + catalogNamesSkippedFromChildSync.add(catalog.getName()); + continue; + } + try { target.createCatalog(catalog); clientLogger.info( @@ -654,6 +669,15 @@ public void syncCatalogs() { } for (Catalog catalog : catalogSyncPlan.entitiesToOverwrite()) { + if (skipCatalogSync) { + clientLogger.info( + "Skipping overwrite of catalog {} because catalog synchronization is disabled. " + + "Catalog-roles and grants will still be synced against the existing target catalog.", + catalog.getName()); + report.recordSuccess(EntityType.CATALOG, SyncOutcome.SKIPPED); + continue; + } + try { target.dropCatalogCascade(catalog.getName()); target.createCatalog(catalog); @@ -676,6 +700,14 @@ public void syncCatalogs() { } for (Catalog catalog : catalogSyncPlan.entitiesToRemove()) { + if (skipCatalogSync) { + clientLogger.info( + "Skipping removal of catalog {} because catalog synchronization is disabled.", + catalog.getName()); + report.recordSuccess(EntityType.CATALOG, SyncOutcome.SKIPPED); + continue; + } + try { target.dropCatalogCascade(catalog.getName()); clientLogger.info( @@ -697,6 +729,9 @@ public void syncCatalogs() { } for (Catalog catalog : catalogSyncPlan.entitiesToSyncChildren()) { + if (catalogNamesSkippedFromChildSync.contains(catalog.getName())) { + continue; + } if (skipIcebergContent) { clientLogger.info( diff --git a/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipCatalogSyncTest.java b/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipCatalogSyncTest.java new file mode 100644 index 00000000..4aae7082 --- /dev/null +++ b/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipCatalogSyncTest.java @@ -0,0 +1,335 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.polaris.tools.sync.polaris; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import org.apache.iceberg.Table; +import org.apache.iceberg.catalog.Namespace; +import org.apache.iceberg.catalog.TableIdentifier; +import org.apache.polaris.core.admin.model.AwsStorageConfigInfo; +import org.apache.polaris.core.admin.model.Catalog; +import org.apache.polaris.core.admin.model.CatalogProperties; +import org.apache.polaris.core.admin.model.CatalogRole; +import org.apache.polaris.core.admin.model.GrantResource; +import org.apache.polaris.core.admin.model.Principal; +import org.apache.polaris.core.admin.model.PrincipalRole; +import org.apache.polaris.core.admin.model.PrincipalWithCredentials; +import org.apache.polaris.core.admin.model.PolarisCatalog; +import org.apache.polaris.core.admin.model.StorageConfigInfo; +import org.apache.polaris.tools.sync.polaris.catalog.NoOpETagManager; +import org.apache.polaris.tools.sync.polaris.planning.NoOpSyncPlanner; +import org.apache.polaris.tools.sync.polaris.planning.plan.SynchronizationPlan; +import org.apache.polaris.tools.sync.polaris.planning.plan.SynchronizationReport; +import org.apache.polaris.tools.sync.polaris.service.IcebergCatalogService; +import org.apache.polaris.tools.sync.polaris.service.PolarisService; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +/** + * Verifies that {@code --skip-catalog-sync} prevents {@link PolarisSynchronizer#syncCatalogs()} + * from creating, overwriting, or removing catalog objects on the target, while still synchronizing + * catalog-roles/grants for catalogs that already exist on the target. + */ +public class PolarisSynchronizerSkipCatalogSyncTest { + + private static Catalog newCatalog(String name) { + return new PolarisCatalog() + .name(name) + .type(Catalog.TypeEnum.INTERNAL) + .properties(new CatalogProperties()) + .storageConfigInfo( + new AwsStorageConfigInfo() + .storageType(StorageConfigInfo.StorageTypeEnum.S3) + .roleArn("roleArn") + .userArn("userArn") + .externalId("externalId") + .region("region")); + } + + private static final Catalog sourceOnlyCatalog = newCatalog("source-only-catalog"); + private static final Catalog overwriteCatalog = newCatalog("overwrite-catalog"); + private static final Catalog removeCatalog = newCatalog("remove-catalog"); + + /** Planner that stages one catalog for CREATE, one for OVERWRITE, and one for REMOVE. */ + private static class CreateOverwriteRemovePlanner extends NoOpSyncPlanner { + @Override + public SynchronizationPlan planCatalogSync( + List catalogsOnSource, List catalogsOnTarget) { + SynchronizationPlan plan = new SynchronizationPlan<>(); + plan.createEntity(sourceOnlyCatalog); + plan.overwriteEntity(overwriteCatalog); + plan.removeEntity(removeCatalog); + return plan; + } + + @Override + public SynchronizationPlan planCatalogRoleSync( + String catalogName, + List catalogRolesOnSource, + List catalogRolesOnTarget) { + return new SynchronizationPlan<>(); + } + + @Override + public SynchronizationPlan planNamespaceSync( + String catalogName, + Namespace namespace, + List namespacesOnSource, + List namespacesOnTarget) { + // NoOpSyncPlanner returns null here, which would NPE once syncNamespaces() runs. + return new SynchronizationPlan<>(); + } + } + + private static class CountingIcebergCatalogService implements IcebergCatalogService { + @Override + public List listNamespaces(Namespace parentNamespace) { + return List.of(); + } + + @Override + public Map loadNamespaceMetadata(Namespace namespace) { + return Map.of(); + } + + @Override + public void createNamespace(Namespace namespace, Map namespaceMetadata) {} + + @Override + public void setNamespaceProperties(Namespace namespace, Map namespaceProperties) {} + + @Override + public void dropNamespaceCascade(Namespace namespace) {} + + @Override + public List listTables(Namespace namespace) { + return List.of(); + } + + @Override + public Table loadTable(TableIdentifier tableIdentifier) { + throw new UnsupportedOperationException(); + } + + @Override + public void registerTable(TableIdentifier tableIdentifier, String metadataFileLocation) {} + + @Override + public void dropTableWithoutPurge(TableIdentifier tableIdentifier) {} + + @Override + public void close() {} + } + + /** Stub {@link PolarisService} that tracks catalog create/overwrite/remove and catalog-role sync calls. */ + private static class TrackingPolarisService implements PolarisService { + + private final List catalogs; + final List catalogsCreated = new ArrayList<>(); + final List catalogsDropped = new ArrayList<>(); + final List catalogRoleSyncsAttempted = new ArrayList<>(); + + TrackingPolarisService(List catalogs) { + this.catalogs = catalogs; + } + + @Override + public void initialize(Map properties) {} + + @Override + public List listPrincipals() { + return List.of(); + } + + @Override + public Principal getPrincipal(String principalName) { + throw new UnsupportedOperationException(); + } + + @Override + public PrincipalWithCredentials createPrincipal(Principal principal) { + throw new UnsupportedOperationException(); + } + + @Override + public void dropPrincipal(String principalName) {} + + @Override + public List listPrincipalRoles() { + return List.of(); + } + + @Override + public PrincipalRole getPrincipalRole(String principalRoleName) { + throw new UnsupportedOperationException(); + } + + @Override + public void createPrincipalRole(PrincipalRole principalRole) {} + + @Override + public void dropPrincipalRole(String principalRoleName) {} + + @Override + public List listPrincipalRolesAssigned(String principalName) { + return List.of(); + } + + @Override + public void assignPrincipalRole(String principalName, String principalRoleName) {} + + @Override + public void revokePrincipalRole(String principalName, String principalRoleName) {} + + @Override + public List listCatalogs() { + return catalogs; + } + + @Override + public Catalog getCatalog(String catalogName) { + throw new UnsupportedOperationException(); + } + + @Override + public void createCatalog(Catalog catalog) { + catalogsCreated.add(catalog.getName()); + } + + @Override + public void dropCatalogCascade(String catalogName) { + catalogsDropped.add(catalogName); + } + + @Override + public List listCatalogRoles(String catalogName) { + catalogRoleSyncsAttempted.add(catalogName); + return List.of(); + } + + @Override + public CatalogRole getCatalogRole(String catalogName, String catalogRoleName) { + throw new UnsupportedOperationException(); + } + + @Override + public void createCatalogRole(String catalogName, CatalogRole catalogRole) {} + + @Override + public void dropCatalogRole(String catalogName, String catalogRoleName) {} + + @Override + public List listAssigneePrincipalRolesForCatalogRole( + String catalogName, String catalogRoleName) { + return List.of(); + } + + @Override + public void assignCatalogRole( + String principalRoleName, String catalogName, String catalogRoleName) {} + + @Override + public void revokeCatalogRole( + String principalRoleName, String catalogName, String catalogRoleName) {} + + @Override + public List listGrants(String catalogName, String catalogRoleName) { + return List.of(); + } + + @Override + public void addGrant(String catalogName, String catalogRoleName, GrantResource grant) {} + + @Override + public void revokeGrant(String catalogName, String catalogRoleName, GrantResource grant) {} + + @Override + public IcebergCatalogService initializeIcebergCatalogService(String catalogName) { + return new CountingIcebergCatalogService(); + } + + @Override + public void close() {} + } + + @Test + public void testSkipCatalogSyncSkipsCreateOverwriteAndRemoveButStillSyncsRolesForExistingCatalogs() { + TrackingPolarisService source = + new TrackingPolarisService(List.of(sourceOnlyCatalog, overwriteCatalog)); + TrackingPolarisService target = new TrackingPolarisService(List.of(overwriteCatalog, removeCatalog)); + + PolarisSynchronizer synchronizer = + new PolarisSynchronizer( + null, + false, + new CreateOverwriteRemovePlanner(), + source, + target, + new NoOpETagManager(), + false, + new SynchronizationReport(), + true /* skipIcebergContent */, + true /* skipCatalogSync */); + + synchronizer.syncCatalogs(); + + // no catalog objects should be created, overwritten (dropped+recreated), or removed on target + Assertions.assertEquals(List.of(), target.catalogsCreated); + Assertions.assertEquals(List.of(), target.catalogsDropped); + + // the source-only catalog has no match on target, so it must be excluded from catalog-role sync + Assertions.assertFalse( + target.catalogRoleSyncsAttempted.contains(sourceOnlyCatalog.getName())); + + // the catalog that exists on both sides should still have its catalog-roles synced + Assertions.assertTrue(target.catalogRoleSyncsAttempted.contains(overwriteCatalog.getName())); + Assertions.assertTrue(source.catalogRoleSyncsAttempted.contains(overwriteCatalog.getName())); + } + + @Test + public void testCatalogSyncedWhenNotSkipped() { + TrackingPolarisService source = + new TrackingPolarisService(List.of(sourceOnlyCatalog, overwriteCatalog)); + TrackingPolarisService target = new TrackingPolarisService(List.of(overwriteCatalog, removeCatalog)); + + PolarisSynchronizer synchronizer = + new PolarisSynchronizer( + null, + false, + new CreateOverwriteRemovePlanner(), + source, + target, + new NoOpETagManager(), + false, + new SynchronizationReport(), + true /* skipIcebergContent */, + false /* skipCatalogSync */); + + synchronizer.syncCatalogs(); + + Assertions.assertEquals( + List.of(sourceOnlyCatalog.getName(), overwriteCatalog.getName()), target.catalogsCreated); + Assertions.assertEquals( + List.of(overwriteCatalog.getName(), removeCatalog.getName()), target.catalogsDropped); + Assertions.assertTrue(target.catalogRoleSyncsAttempted.contains(sourceOnlyCatalog.getName())); + Assertions.assertTrue(target.catalogRoleSyncsAttempted.contains(overwriteCatalog.getName())); + } +} diff --git a/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipIcebergContentTest.java b/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipIcebergContentTest.java index 34cc3546..8d2e313b 100644 --- a/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipIcebergContentTest.java +++ b/polaris-synchronizer/api/src/test/java/org/apache/polaris/tools/sync/polaris/PolarisSynchronizerSkipIcebergContentTest.java @@ -301,7 +301,8 @@ public void testSkipIcebergContentSkipsIcebergSyncButStillSyncsCatalogRoles() { new NoOpETagManager(), false, new SynchronizationReport(), - true); + true, + false); synchronizer.syncCatalogs(); @@ -326,6 +327,7 @@ public void testIcebergContentSyncedWhenNotSkipped() { new NoOpETagManager(), false, new SynchronizationReport(), + false, false); synchronizer.syncCatalogs(); @@ -358,7 +360,8 @@ public void testTableScopedGrantsStillAttemptedWhenIcebergContentSkipped() { new NoOpETagManager(), false, new SynchronizationReport(), - true); + true, + false); synchronizer.syncCatalogs(); diff --git a/polaris-synchronizer/cli/src/main/java/org/apache/polaris/tools/sync/polaris/SyncPolarisCommand.java b/polaris-synchronizer/cli/src/main/java/org/apache/polaris/tools/sync/polaris/SyncPolarisCommand.java index e761a86c..71d138e6 100644 --- a/polaris-synchronizer/cli/src/main/java/org/apache/polaris/tools/sync/polaris/SyncPolarisCommand.java +++ b/polaris-synchronizer/cli/src/main/java/org/apache/polaris/tools/sync/polaris/SyncPolarisCommand.java @@ -112,6 +112,16 @@ public class SyncPolarisCommand implements Callable { ) private boolean skipIcebergContent; + @CommandLine.Option( + names = {"--skip-catalog-sync"}, + description = "Skip creation, overwrite, and removal of catalogs themselves. Catalog-roles and grants " + + "will still be synchronized for catalogs that already exist on the target. Catalogs that only " + + "exist on the source will be skipped entirely, since there is no matching catalog on the target " + + "to synchronize catalog-roles and grants against. Useful for disaster-recovery setups where " + + "target catalogs are pre-created with different storage locations than the source." + ) + private boolean skipCatalogSync; + @CommandLine.Option( names = {"--strategy"}, defaultValue = "CREATE_ONLY", @@ -164,7 +174,8 @@ public Integer call() throws Exception { etagManager, diffOnly, report, - skipIcebergContent); + skipIcebergContent, + skipCatalogSync); synchronizer.syncPrincipalRoles(); if (shouldSyncPrincipals) { consoleLog.warn("Principal migration will reset credentials on the target Polaris instance. " +