This is an automated email from the ASF dual-hosted git repository.

jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/branch-1.3 by this push:
     new ad9a9c4d1e [Cherry-pick to branch-1.3] [#12761] fix(clickhouse): 
Propagate clustered table rename (#12763) (#12832)
ad9a9c4d1e is described below

commit ad9a9c4d1ec3d66cba054bfea01df6956ff1cff3
Author: StormSpirit <[email protected]>
AuthorDate: Wed Sep 2 14:51:24 2026 +0800

    [Cherry-pick to branch-1.3] [#12761] fix(clickhouse): Propagate clustered 
table rename (#12763) (#12832)
    
    ### What changes were proposed in this pull request?
    
    This PR manually backports #12763 to `branch-1.3` and supersedes the
    unresolved automated cherry-pick PR #12790.
    
    - Override the ClickHouse table rename path so tables with trusted
    Gravitino cluster metadata use `RENAME TABLE ... ON CLUSTER ...`, while
    unmarked external tables keep local rename behavior.
    - Read the existing table metadata from `system.tables` with an exact
    current-database and table-name predicate, and reject a present but
    blank cluster marker before executing DDL.
    - Preserve the existing ClickHouse exception mapping and identifier
    quoting contract without changing the common JDBC rename path.
    - Expand the ClickHouse cluster fixture to three independently
    addressable nodes and add lifecycle coverage for rename propagation,
    query text, all-node state, comment metadata, and cleanup.
    - Resolve the branch-1.3 unit-test conflict by retaining the target
    branch tests and only the rename-related tests from #12763;
    source-branch-only index parameter tests are not included.
    
    ### Why are the changes needed?
    
    The common JDBC rename path generates a local `RENAME TABLE old_name TO
    new_name` statement. For ClickHouse tables created through Gravitino
    with `ON CLUSTER`, that statement renames only the JDBC connection node
    and leaves the old name on the other nodes, silently splitting cluster
    metadata and making later DDL inconsistent.
    
    ClickHouse does not include `ON CLUSTER` in `SHOW CREATE TABLE`, so the
    connector must use the trusted cluster marker embedded in the stored
    comment when generating the rename statement.
    
    The automated backport PR #12790 contains unresolved conflict markers in
    the ClickHouse unit test. This PR resolves that conflict against the
    current `branch-1.3` baseline without carrying unrelated source-branch
    changes.
    
    The source-branch-only `StringIdentifier.ID_KEY` property assertions are
    not carried into this `branch-1.3` backport because the existing
    branch-1.3 ClickHouse flow does not expose that derived property for
    this test path. The cluster integration test still verifies that the raw
    StringIdentifier and cluster marker remain in `system.tables.comment` on
    every node.
    
    Fix: #12761
    
    ### Does this PR introduce _any_ user-facing change?
    
    Yes. Renaming a Gravitino-created ClickHouse table with trusted cluster
    metadata now propagates the rename to every configured cluster node.
    Local tables and unmarked external tables retain local rename behavior.
    This PR does not add, remove, or change catalog-facing APIs or property
    keys.
    
    ### How was this patch tested?
    
    - `./gradlew :catalogs-contrib:catalog-jdbc-clickhouse:spotlessCheck` —
    passed.
    - `./gradlew rat` — passed.
    - `./gradlew :catalogs-contrib:catalog-jdbc-clickhouse:test -PskipITs
    -PskipDockerTests=true` — passed.
    - `./gradlew :catalogs-contrib:catalog-jdbc-clickhouse:test --tests
    
'org.apache.gravitino.catalog.clickhouse.integration.test.CatalogClickHouseIT'
    -PskipDockerTests=false` — passed.
    - `./gradlew :catalogs-contrib:catalog-jdbc-clickhouse:test --tests
    
'org.apache.gravitino.catalog.clickhouse.integration.test.CatalogClickHouseClusterIT'
    -PskipDockerTests=false` — passed with 18 tests, 0 skipped, 0 failures,
    and 0 errors.
    - `./gradlew :catalogs-contrib:catalog-jdbc-clickhouse:build -x test` —
    passed.
    
    Signed-off-by: jiangxt2 <[email protected]>
---
 .../operations/ClickHouseClusterUtils.java         |  14 ++
 .../operations/ClickHouseTableOperations.java      | 176 +++++++++++++++------
 .../test/CatalogClickHouseClusterIT.java           | 143 ++++++++++++++++-
 .../integration/test/CatalogClickHouseIT.java      |  15 ++
 .../TestClickHouseTableOperationsCluster.java      |  37 +++++
 .../TestClickHouseTableOperationsUnit.java         | 121 +++++++++++++-
 .../src/test/resources/remote_servers.xml          |  10 +-
 .../integration/test/container/ContainerSuite.java |  57 ++++---
 8 files changed, 498 insertions(+), 75 deletions(-)

diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseClusterUtils.java
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseClusterUtils.java
index b225102650..010c6dc4a5 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseClusterUtils.java
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseClusterUtils.java
@@ -141,6 +141,20 @@ public final class ClickHouseClusterUtils {
     return storedComment.substring(idx + CLUSTER_META_PREFIX.length());
   }
 
+  /**
+   * Returns whether {@code storedComment} contains the Gravitino cluster 
metadata marker.
+   *
+   * <p>This deliberately distinguishes an absent marker from a present marker 
with a blank value.
+   * Callers performing cluster-wide DDL can therefore keep unmarked external 
objects local while
+   * rejecting corrupted Gravitino metadata before mutation.
+   *
+   * @param storedComment The raw comment as stored in ClickHouse.
+   * @return {@code true} if the Gravitino cluster metadata marker is present.
+   */
+  public static boolean hasClusterMetadata(String storedComment) {
+    return storedComment != null && 
storedComment.contains(CLUSTER_META_PREFIX);
+  }
+
   /**
    * Returns the user-visible portion of the stored comment, stripping any 
embedded cluster metadata
    * suffix. Returns {@code null} if {@code storedComment} is {@code null}.
diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseTableOperations.java
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseTableOperations.java
index 8d7afb9ed5..e0e67c73bd 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseTableOperations.java
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/main/java/org/apache/gravitino/catalog/clickhouse/operations/ClickHouseTableOperations.java
@@ -690,68 +690,116 @@ public class ClickHouseTableOperations extends 
JdbcTableOperations {
     }
   }
 
+  @Override
+  public void rename(String databaseName, String oldTableName, String 
newTableName)
+      throws NoSuchTableException {
+    LOG.info(
+        "Attempting to rename table {}/{} to {}/{}",
+        databaseName,
+        oldTableName,
+        databaseName,
+        newTableName);
+    try (Connection connection = getConnection(databaseName)) {
+      TablePropertiesWithClusterMetadata metadata =
+          loadTablePropertiesWithClusterMetadata(connection, oldTableName);
+      JdbcConnectorUtils.executeUpdate(
+          connection,
+          generateRenameTableSql(
+              oldTableName, newTableName, metadata.hasClusterMetadata(), 
metadata.clusterName()));
+      LOG.info(
+          "Renamed table {}/{} to {}/{}", databaseName, oldTableName, 
databaseName, newTableName);
+    } catch (final SQLException se) {
+      throw exceptionMapper.toGravitinoException(se);
+    }
+  }
+
+  @VisibleForTesting
+  String generateRenameTableSql(
+      String oldTableName,
+      String newTableName,
+      boolean hasClusterMetadata,
+      @Nullable String clusterName) {
+    if (hasClusterMetadata) {
+      Preconditions.checkArgument(
+          StringUtils.isNotBlank(clusterName),
+          "ClickHouse cluster metadata for table %s is missing a cluster name",
+          oldTableName);
+      return "RENAME TABLE %s TO %s ON CLUSTER %s"
+          .formatted(
+              quoteIdentifier(oldTableName),
+              quoteIdentifier(newTableName),
+              quoteIdentifier(clusterName));
+    }
+    return "RENAME TABLE %s TO %s"
+        .formatted(quoteIdentifier(oldTableName), 
quoteIdentifier(newTableName));
+  }
+
   @Override
   protected Map<String, String> getTableProperties(Connection connection, 
String tableName)
       throws SQLException {
+    return loadTablePropertiesWithClusterMetadata(connection, 
tableName).properties();
+  }
+
+  private TablePropertiesWithClusterMetadata 
loadTablePropertiesWithClusterMetadata(
+      Connection connection, String tableName) throws SQLException {
     try (PreparedStatement statement =
-        connection.prepareStatement("select * from system.tables where name = 
? ")) {
+        connection.prepareStatement(
+            "SELECT comment, engine, engine_full FROM system.tables "
+                + "WHERE database = currentDatabase() AND name = ?")) {
       statement.setString(1, tableName);
       try (ResultSet resultSet = statement.executeQuery()) {
-        while (resultSet.next()) {
-          String name = resultSet.getString("name");
-          if (Objects.equals(name, tableName)) {
-            Map<String, String> tableProperties = new HashMap<>();
-
-            // Extract cluster name embedded in the COMMENT at create time.
-            // SHOW CREATE TABLE does not include ON CLUSTER (see 
ClickHouseClusterUtils).
-            String storedComment = resultSet.getString(COMMENT);
-            String clusterName = 
ClickHouseClusterUtils.extractClusterFromComment(storedComment);
-            tableProperties.put(
-                COMMENT, 
ClickHouseClusterUtils.stripClusterMetadata(storedComment));
-            String engine = resultSet.getString(CLICKHOUSE_ENGINE_KEY);
-            String engineFull = resultSet.getString("engine_full");
-            tableProperties.put(GRAVITINO_ENGINE_KEY, engine);
-            if (StringUtils.isNotBlank(clusterName)) {
-              tableProperties.put(ClusterConstants.ON_CLUSTER, 
String.valueOf(true));
-              tableProperties.put(ClusterConstants.CLUSTER_NAME, clusterName);
-            } else {
-              tableProperties.put(ClusterConstants.ON_CLUSTER, 
String.valueOf(false));
-            }
+        if (!resultSet.next()) {
+          throw new NoSuchTableException(
+              "Table %s does not exist in %s.", tableName, 
connection.getCatalog());
+        }
 
-            if (StringUtils.equalsIgnoreCase(engine, 
ENGINE.DISTRIBUTED.getValue())) {
-              Matcher distributedEngineMatcher =
-                  
DISTRIBUTED_ENGINE_PATTERN.matcher(StringUtils.trimToEmpty(engineFull));
-              if (distributedEngineMatcher.matches()) {
-                String distributedClusterName = 
unquote(distributedEngineMatcher.group(1));
-                tableProperties.put(ClusterConstants.CLUSTER_NAME, 
distributedClusterName);
-                tableProperties.put(
-                    DistributedTableConstants.REMOTE_DATABASE,
-                    unquote(distributedEngineMatcher.group(2)));
-                tableProperties.put(
-                    DistributedTableConstants.REMOTE_TABLE,
-                    unquote(distributedEngineMatcher.group(3)));
-                tableProperties.put(
-                    DistributedTableConstants.SHARDING_KEY,
-                    StringUtils.trim(distributedEngineMatcher.group(4)));
-              }
-            } else if (StringUtils.equalsIgnoreCase(engine, 
ENGINE.GRAPHITEMERGETREE.getValue())) {
-              String graphiteConfig = extractGraphiteConfig(engineFull);
-              if (StringUtils.isNotBlank(graphiteConfig)) {
-                tableProperties.put(TableConstants.GRAPHITE_CONFIG, 
graphiteConfig);
-              }
-            } else if (isGenericEngineParameterEngine(engine)) {
-              String engineParams = extractEngineParams(engine, engineFull);
-              if (StringUtils.isNotBlank(engineParams)) {
-                tableProperties.put(TableConstants.ENGINE_PARAMETERS, 
engineParams);
-              }
-            }
+        Map<String, String> tableProperties = new HashMap<>();
+
+        // Extract cluster name embedded in the COMMENT at create time.
+        // SHOW CREATE TABLE does not include ON CLUSTER (see 
ClickHouseClusterUtils).
+        String storedComment = resultSet.getString(COMMENT);
+        boolean hasClusterMetadata = 
ClickHouseClusterUtils.hasClusterMetadata(storedComment);
+        String clusterName = 
ClickHouseClusterUtils.extractClusterFromComment(storedComment);
+        tableProperties.put(COMMENT, 
ClickHouseClusterUtils.stripClusterMetadata(storedComment));
+        String engine = resultSet.getString(CLICKHOUSE_ENGINE_KEY);
+        String engineFull = resultSet.getString("engine_full");
+        tableProperties.put(GRAVITINO_ENGINE_KEY, engine);
+        if (StringUtils.isNotBlank(clusterName)) {
+          tableProperties.put(ClusterConstants.ON_CLUSTER, 
String.valueOf(true));
+          tableProperties.put(ClusterConstants.CLUSTER_NAME, clusterName);
+        } else {
+          tableProperties.put(ClusterConstants.ON_CLUSTER, 
String.valueOf(false));
+        }
 
-            return Collections.unmodifiableMap(tableProperties);
+        if (StringUtils.equalsIgnoreCase(engine, 
ENGINE.DISTRIBUTED.getValue())) {
+          Matcher distributedEngineMatcher =
+              
DISTRIBUTED_ENGINE_PATTERN.matcher(StringUtils.trimToEmpty(engineFull));
+          if (distributedEngineMatcher.matches()) {
+            String distributedClusterName = 
unquote(distributedEngineMatcher.group(1));
+            tableProperties.put(ClusterConstants.CLUSTER_NAME, 
distributedClusterName);
+            tableProperties.put(
+                DistributedTableConstants.REMOTE_DATABASE,
+                unquote(distributedEngineMatcher.group(2)));
+            tableProperties.put(
+                DistributedTableConstants.REMOTE_TABLE, 
unquote(distributedEngineMatcher.group(3)));
+            tableProperties.put(
+                DistributedTableConstants.SHARDING_KEY,
+                StringUtils.trim(distributedEngineMatcher.group(4)));
+          }
+        } else if (StringUtils.equalsIgnoreCase(engine, 
ENGINE.GRAPHITEMERGETREE.getValue())) {
+          String graphiteConfig = extractGraphiteConfig(engineFull);
+          if (StringUtils.isNotBlank(graphiteConfig)) {
+            tableProperties.put(TableConstants.GRAPHITE_CONFIG, 
graphiteConfig);
+          }
+        } else if (isGenericEngineParameterEngine(engine)) {
+          String engineParams = extractEngineParams(engine, engineFull);
+          if (StringUtils.isNotBlank(engineParams)) {
+            tableProperties.put(TableConstants.ENGINE_PARAMETERS, 
engineParams);
           }
         }
 
-        throw new NoSuchTableException(
-            "Table %s does not exist in %s.", tableName, 
connection.getCatalog());
+        return new TablePropertiesWithClusterMetadata(
+            Collections.unmodifiableMap(tableProperties), hasClusterMetadata, 
clusterName);
       }
     }
   }
@@ -1606,6 +1654,32 @@ public class ClickHouseTableOperations extends 
JdbcTableOperations {
     private Transform[] partitioning = Transforms.EMPTY_TRANSFORM;
   }
 
+  private static final class TablePropertiesWithClusterMetadata {
+    private final Map<String, String> properties;
+    private final boolean hasClusterMetadata;
+    @Nullable private final String clusterName;
+
+    private TablePropertiesWithClusterMetadata(
+        Map<String, String> properties, boolean hasClusterMetadata, @Nullable 
String clusterName) {
+      this.properties = properties;
+      this.hasClusterMetadata = hasClusterMetadata;
+      this.clusterName = clusterName;
+    }
+
+    private Map<String, String> properties() {
+      return properties;
+    }
+
+    private boolean hasClusterMetadata() {
+      return hasClusterMetadata;
+    }
+
+    @Nullable
+    private String clusterName() {
+      return clusterName;
+    }
+  }
+
   @VisibleForTesting
   String[][] parseIndexFields(String expression) {
     return ClickHouseTableSqlUtils.parseIndexFields(expression);
diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseClusterIT.java
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseClusterIT.java
index 4ac96e8f31..fb7bb74c56 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseClusterIT.java
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseClusterIT.java
@@ -31,11 +31,15 @@ import com.google.common.collect.Maps;
 import java.nio.file.Paths;
 import java.sql.Connection;
 import java.sql.DriverManager;
+import java.sql.PreparedStatement;
 import java.sql.ResultSet;
+import java.sql.SQLException;
 import java.sql.Statement;
+import java.time.Duration;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import org.apache.commons.lang3.StringUtils;
@@ -43,7 +47,9 @@ import org.apache.gravitino.Catalog;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
 import org.apache.gravitino.Schema;
+import org.apache.gravitino.StringIdentifier;
 import 
org.apache.gravitino.catalog.clickhouse.integration.test.service.ClickHouseService;
+import 
org.apache.gravitino.catalog.clickhouse.operations.ClickHouseClusterUtils;
 import org.apache.gravitino.catalog.jdbc.config.JdbcConfig;
 import org.apache.gravitino.client.GravitinoMetalake;
 import org.apache.gravitino.integration.test.container.ClickHouseContainer;
@@ -64,6 +70,7 @@ import 
org.apache.gravitino.rel.expressions.transforms.Transforms;
 import org.apache.gravitino.rel.indexes.Index;
 import org.apache.gravitino.rel.indexes.Indexes;
 import org.apache.gravitino.rel.types.Types;
+import org.awaitility.Awaitility;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeAll;
@@ -98,6 +105,7 @@ public class CatalogClickHouseClusterIT extends BaseIT {
   private Catalog catalog;
   private ClickHouseService clickHouseService;
   private ClickHouseContainer clickHouseClusterContainer;
+  private List<ClickHouseContainer> clickHouseClusterContainers;
   private final TestDatabaseName TEST_DB_NAME = 
TestDatabaseName.CLICKHOUSE_CLUSTER_CLICKHOUSE_IT;
 
   @BeforeAll
@@ -106,6 +114,8 @@ public class CatalogClickHouseClusterIT extends BaseIT {
         Paths.get("src", "test", "resources", 
"remote_servers.xml").toAbsolutePath().toString();
     containerSuite.startClickHouseClusterContainer(TEST_DB_NAME, 
remoteServersConfig);
     clickHouseClusterContainer = 
containerSuite.getClickHouseClusterContainer();
+    clickHouseClusterContainers = 
containerSuite.getClickHouseClusterContainers();
+    Assertions.assertEquals(3, clickHouseClusterContainers.size());
 
     clickHouseService = new ClickHouseService(clickHouseClusterContainer, 
TEST_DB_NAME);
     createMetalake();
@@ -162,8 +172,10 @@ public class CatalogClickHouseClusterIT extends BaseIT {
   }
 
   private void createSchema() {
-    Schema createdSchema =
-        catalog.asSchemas().createSchema(schemaName, null, 
Collections.emptyMap());
+    Map<String, String> properties = new HashMap<>();
+    properties.put(CLUSTER_NAME, ClickHouseContainer.DEFAULT_CLUSTER_NAME);
+    properties.put(ON_CLUSTER, String.valueOf(true));
+    Schema createdSchema = catalog.asSchemas().createSchema(schemaName, null, 
properties);
     Schema loadSchema = catalog.asSchemas().loadSchema(schemaName);
     Assertions.assertEquals(createdSchema.name(), loadSchema.name());
   }
@@ -688,6 +700,56 @@ public class CatalogClickHouseClusterIT extends BaseIT {
     }
   }
 
+  @Test
+  public void testRenameTableOnClusterPropagatesToEveryNode() throws 
SQLException {
+    String oldName = GravitinoITUtils.genRandomName("ck_cluster_rename_old");
+    String newName = GravitinoITUtils.genRandomName("ck_cluster_rename_new");
+    NameIdentifier oldIdent = NameIdentifier.of(schemaName, oldName);
+    NameIdentifier newIdent = NameIdentifier.of(schemaName, newName);
+    TableCatalog tableCatalog = catalog.asTableCatalog();
+
+    try {
+      tableCatalog.createTable(
+          oldIdent,
+          createColumns(),
+          "cluster rename comment",
+          clusterMergeTreeProperties(),
+          Transforms.EMPTY_TRANSFORM,
+          Distributions.NONE,
+          getSortOrders("col_3"),
+          Indexes.EMPTY_INDEXES);
+
+      Table renamed = tableCatalog.alterTable(oldIdent, 
TableChange.rename(newName));
+      Assertions.assertEquals(newName, renamed.name());
+      Assertions.assertEquals(String.valueOf(true), 
renamed.properties().get(ON_CLUSTER));
+      Assertions.assertEquals(
+          ClickHouseContainer.DEFAULT_CLUSTER_NAME, 
renamed.properties().get(CLUSTER_NAME));
+
+      awaitTableStateOnEveryNode(oldName, false, newName, true);
+      assertRenameQueryUsesOnCluster(oldName, newName);
+      for (ClickHouseContainer container : clickHouseClusterContainers) {
+        String storedComment = loadStoredComment(container, newName);
+        Assertions.assertNotNull(StringIdentifier.fromComment(storedComment));
+        Assertions.assertEquals(
+            ClickHouseContainer.DEFAULT_CLUSTER_NAME,
+            ClickHouseClusterUtils.extractClusterFromComment(storedComment));
+      }
+
+      Assertions.assertTrue(tableCatalog.dropTable(newIdent));
+      awaitTableStateOnEveryNode(oldName, false, newName, false);
+    } finally {
+      try {
+        clickHouseService.executeQuery(
+            "DROP TABLE IF EXISTS `%s`.`%s` ON CLUSTER `%s` SYNC"
+                .formatted(schemaName, oldName, 
ClickHouseContainer.DEFAULT_CLUSTER_NAME));
+      } finally {
+        clickHouseService.executeQuery(
+            "DROP TABLE IF EXISTS `%s`.`%s` ON CLUSTER `%s` SYNC"
+                .formatted(schemaName, newName, 
ClickHouseContainer.DEFAULT_CLUSTER_NAME));
+      }
+    }
+  }
+
   /**
    * When a user updates the table comment via Gravitino, the cluster metadata 
embedded in the
    * ClickHouse COMMENT field must be preserved. Without re-embedding, the 
next loadTable call would
@@ -1103,4 +1165,81 @@ public class CatalogClickHouseClusterIT extends BaseIT {
       tableCatalog.dropTable(localIdent);
     }
   }
+
+  private void awaitTableStateOnEveryNode(
+      String oldTableName, boolean oldTableExists, String newTableName, 
boolean newTableExists) {
+    Awaitility.await()
+        .atMost(Duration.ofSeconds(30))
+        .pollInterval(Duration.ofMillis(250))
+        .untilAsserted(
+            () -> {
+              for (ClickHouseContainer container : 
clickHouseClusterContainers) {
+                Assertions.assertEquals(
+                    oldTableExists,
+                    tableExists(container, oldTableName),
+                    "Unexpected old-table state on " + 
container.getContainerIpAddress());
+                Assertions.assertEquals(
+                    newTableExists,
+                    tableExists(container, newTableName),
+                    "Unexpected new-table state on " + 
container.getContainerIpAddress());
+              }
+            });
+  }
+
+  private boolean tableExists(ClickHouseContainer container, String tableName) 
throws SQLException {
+    try (Connection connection =
+            DriverManager.getConnection(
+                container.getJdbcUrl(TEST_DB_NAME),
+                container.getUsername(),
+                container.getPassword());
+        PreparedStatement statement =
+            connection.prepareStatement(
+                "SELECT count() FROM system.tables WHERE database = ? AND name 
= ?")) {
+      statement.setString(1, schemaName);
+      statement.setString(2, tableName);
+      try (ResultSet resultSet = statement.executeQuery()) {
+        Assertions.assertTrue(resultSet.next());
+        return resultSet.getLong(1) == 1;
+      }
+    }
+  }
+
+  private String loadStoredComment(ClickHouseContainer container, String 
tableName)
+      throws SQLException {
+    try (Connection connection =
+            DriverManager.getConnection(
+                container.getJdbcUrl(TEST_DB_NAME),
+                container.getUsername(),
+                container.getPassword());
+        PreparedStatement statement =
+            connection.prepareStatement(
+                "SELECT comment FROM system.tables WHERE database = ? AND name 
= ?")) {
+      statement.setString(1, schemaName);
+      statement.setString(2, tableName);
+      try (ResultSet resultSet = statement.executeQuery()) {
+        Assertions.assertTrue(resultSet.next());
+        return resultSet.getString(1);
+      }
+    }
+  }
+
+  private void assertRenameQueryUsesOnCluster(String oldTableName, String 
newTableName) {
+    clickHouseService.executeQuery("SYSTEM FLUSH LOGS");
+    String query =
+        clickHouseService.executeQueryForResult(
+            String.format(
+                "SELECT query FROM system.query_log "
+                    + "WHERE type = 'QueryFinish' "
+                    + "AND startsWith(query, 'RENAME TABLE') "
+                    + "AND query LIKE '%%`%s`%%' "
+                    + "ORDER BY event_time DESC LIMIT 1",
+                oldTableName));
+
+    Assertions.assertNotNull(query, "The initiating RENAME query must be 
present in query_log");
+    Assertions.assertTrue(
+        query.contains("RENAME TABLE `%s` TO `%s`".formatted(oldTableName, 
newTableName)));
+    Assertions.assertTrue(
+        query.contains("ON CLUSTER 
`%s`".formatted(ClickHouseContainer.DEFAULT_CLUSTER_NAME)));
+    Assertions.assertEquals(1, StringUtils.countMatches(query, "ON CLUSTER"));
+  }
 }
diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseIT.java
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseIT.java
index f99ee59ed2..d715af3a25 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseIT.java
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/integration/test/CatalogClickHouseIT.java
@@ -1372,6 +1372,21 @@ public class CatalogClickHouseIT extends BaseIT {
         .asTableCatalog()
         .alterTable(NameIdentifier.of(schemaName, tableName), 
TableChange.rename(alertTableName));
 
+    clickhouseService.executeQuery("SYSTEM FLUSH LOGS");
+    String renameQuery =
+        clickhouseService.executeQueryForResult(
+            String.format(
+                "SELECT query FROM system.query_log "
+                    + "WHERE type = 'QueryFinish' "
+                    + "AND startsWith(query, 'RENAME TABLE') "
+                    + "AND query LIKE '%%`%s`%%' "
+                    + "ORDER BY event_time DESC LIMIT 1",
+                tableName));
+    Assertions.assertNotNull(renameQuery);
+    Assertions.assertTrue(
+        renameQuery.contains("RENAME TABLE `%s` TO `%s`".formatted(tableName, 
alertTableName)));
+    Assertions.assertFalse(renameQuery.contains("ON CLUSTER"));
+
     catalog
         .asTableCatalog()
         .alterTable(
diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsCluster.java
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsCluster.java
index 5978fe212c..aa1120cebf 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsCluster.java
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsCluster.java
@@ -21,6 +21,7 @@ package org.apache.gravitino.catalog.clickhouse.operations;
 import static 
org.apache.gravitino.catalog.clickhouse.ClickHouseTablePropertiesMetadata.GRAVITINO_ENGINE_KEY;
 import static 
org.apache.gravitino.catalog.clickhouse.operations.ClickHouseClusterUtils.CLUSTER_META_PREFIX;
 import static 
org.apache.gravitino.catalog.clickhouse.operations.ClickHouseClusterUtils.extractClusterFromComment;
+import static 
org.apache.gravitino.catalog.clickhouse.operations.ClickHouseClusterUtils.hasClusterMetadata;
 import static 
org.apache.gravitino.catalog.clickhouse.operations.ClickHouseClusterUtils.stripClusterMetadata;
 
 import java.util.HashMap;
@@ -253,6 +254,30 @@ class TestClickHouseTableOperationsCluster {
     Assertions.assertEquals("DROP TABLE `orders` ON CLUSTER `ck_cluster` 
SYNC", sql);
   }
 
+  @Test
+  void testGenerateRenameTableSqlWithoutClusterMetadata() {
+    String sql = ops.buildRenameSql("old-table", "new table", false, null);
+
+    Assertions.assertEquals("RENAME TABLE `old-table` TO `new table`", sql);
+  }
+
+  @Test
+  void testGenerateRenameTableSqlWithClusterMetadata() {
+    String sql = ops.buildRenameSql("old-table", "new table", true, 
"ck-cluster");
+
+    Assertions.assertEquals("RENAME TABLE `old-table` TO `new table` ON 
CLUSTER `ck-cluster`", sql);
+  }
+
+  @Test
+  void testGenerateRenameTableSqlRejectsBlankClusterMetadata() {
+    IllegalArgumentException exception =
+        Assertions.assertThrows(
+            IllegalArgumentException.class,
+            () -> ops.buildRenameSql("orders", "renamed_orders", true, " "));
+
+    Assertions.assertTrue(exception.getMessage().contains("missing a cluster 
name"));
+  }
+
   @Test
   void testGenerateDropTableSqlOnClusterWithoutClusterName() {
     // on-cluster=true but no cluster-name → fall back to plain DROP TABLE
@@ -364,6 +389,13 @@ class TestClickHouseTableOperationsCluster {
     Assertions.assertEquals("", stripClusterMetadata(stored));
   }
 
+  @Test
+  void testClusterMetadataPresenceDistinguishesAbsentAndBlankMarker() {
+    Assertions.assertFalse(hasClusterMetadata("plain comment"));
+    Assertions.assertTrue(hasClusterMetadata("comment" + CLUSTER_META_PREFIX));
+    Assertions.assertEquals("", extractClusterFromComment("comment" + 
CLUSTER_META_PREFIX));
+  }
+
   /** ALTER TABLE with ON CLUSTER=true should include ON CLUSTER in SQL. */
   @Test
   void testAlterTableWithOnCluster() {
@@ -827,6 +859,11 @@ class TestClickHouseTableOperationsCluster {
     String buildDropSql(String tableName, Map<String, String> properties) {
       return generateDropTableSql(tableName, properties);
     }
+
+    String buildRenameSql(
+        String oldTableName, String newTableName, boolean hasClusterMetadata, 
String clusterName) {
+      return generateRenameTableSql(oldTableName, newTableName, 
hasClusterMetadata, clusterName);
+    }
   }
 
   private static final class StubClickHouseTableOperations extends 
ClickHouseTableOperations {
diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsUnit.java
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsUnit.java
index ae24dbf649..da3b533de5 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsUnit.java
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/java/org/apache/gravitino/catalog/clickhouse/operations/TestClickHouseTableOperationsUnit.java
@@ -23,15 +23,19 @@ import static 
org.apache.gravitino.catalog.clickhouse.ClickHouseUtils.getSortOrd
 import java.sql.Connection;
 import java.sql.PreparedStatement;
 import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import javax.sql.DataSource;
 import 
org.apache.gravitino.catalog.clickhouse.ClickHouseConstants.TableConstants;
 import 
org.apache.gravitino.catalog.clickhouse.ClickHouseTablePropertiesMetadata.ENGINE;
 import 
org.apache.gravitino.catalog.clickhouse.converter.ClickHouseColumnDefaultValueConverter;
 import 
org.apache.gravitino.catalog.clickhouse.converter.ClickHouseExceptionConverter;
 import 
org.apache.gravitino.catalog.clickhouse.converter.ClickHouseTypeConverter;
 import org.apache.gravitino.catalog.jdbc.JdbcColumn;
+import org.apache.gravitino.exceptions.GravitinoRuntimeException;
 import org.apache.gravitino.exceptions.NoSuchTableException;
 import org.apache.gravitino.rel.expressions.FunctionExpression;
 import org.apache.gravitino.rel.expressions.NamedReference;
@@ -86,9 +90,13 @@ public class TestClickHouseTableOperationsUnit {
   }
 
   private ExposedClickHouseTableOperations newOps() {
+    return newOps(null);
+  }
+
+  private ExposedClickHouseTableOperations newOps(DataSource dataSource) {
     ExposedClickHouseTableOperations ops = new 
ExposedClickHouseTableOperations();
     ops.initialize(
-        null,
+        dataSource,
         new ClickHouseExceptionConverter(),
         new ClickHouseTypeConverter(),
         new ClickHouseColumnDefaultValueConverter(),
@@ -101,7 +109,6 @@ public class TestClickHouseTableOperationsUnit {
     PreparedStatement statement = Mockito.mock(PreparedStatement.class);
     ResultSet resultSet = Mockito.mock(ResultSet.class);
     Mockito.when(resultSet.next()).thenReturn(true);
-    Mockito.when(resultSet.getString("name")).thenReturn("test_table");
     Mockito.when(resultSet.getString("COMMENT")).thenReturn("");
     Mockito.when(resultSet.getString("ENGINE")).thenReturn(engine);
     Mockito.when(resultSet.getString("engine_full")).thenReturn(engineFull);
@@ -263,6 +270,82 @@ public class TestClickHouseTableOperationsUnit {
     Assertions.assertTrue(exception.getMessage().contains("db_name"));
   }
 
+  @Test
+  void testGetTablePropertiesScopesMetadataToCurrentDatabase() throws 
Exception {
+    ExposedClickHouseTableOperations ops = newOps();
+    Connection connection = Mockito.mock(Connection.class);
+    PreparedStatement statement = Mockito.mock(PreparedStatement.class);
+    ResultSet resultSet = Mockito.mock(ResultSet.class);
+    ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class);
+
+    
Mockito.when(connection.prepareStatement(sqlCaptor.capture())).thenReturn(statement);
+    Mockito.when(statement.executeQuery()).thenReturn(resultSet);
+    Mockito.when(resultSet.next()).thenReturn(true);
+    Mockito.when(resultSet.getString("COMMENT")).thenReturn("table comment");
+    
Mockito.when(resultSet.getString("ENGINE")).thenReturn(ENGINE.MERGETREE.getValue());
+    Mockito.when(resultSet.getString("engine_full")).thenReturn("MergeTree 
ORDER BY id");
+
+    ops.callGetTableProperties(connection, "same_name");
+
+    Assertions.assertEquals(
+        "SELECT comment, engine, engine_full FROM system.tables "
+            + "WHERE database = currentDatabase() AND name = ?",
+        sqlCaptor.getValue());
+    Mockito.verify(statement).setString(1, "same_name");
+  }
+
+  @Test
+  void testRenameUsesTrustedClusterMetadata() throws Exception {
+    RenameMocks mocks = renameMocks("comment\n[Gravitino] 
ch.cluster=ck_cluster", "MergeTree");
+    ExposedClickHouseTableOperations ops = newOps(mocks.dataSource);
+
+    ops.rename("db_name", "old-table", "new table");
+
+    Mockito.verify(mocks.connection).setCatalog("db_name");
+    Mockito.verify(mocks.updateStatement)
+        .executeUpdate("RENAME TABLE `old-table` TO `new table` ON CLUSTER 
`ck_cluster`");
+  }
+
+  @Test
+  void testRenameDoesNotPromoteUnmarkedDistributedTableToClusterScope() throws 
Exception {
+    // The Distributed engine contains a cluster name, but only the Gravitino 
comment marker may
+    // authorize cluster-wide DDL.
+    RenameMocks mocks =
+        renameMocks("external table", "Distributed('ck_cluster', 'db', 
'remote', id)");
+    ExposedClickHouseTableOperations ops = newOps(mocks.dataSource);
+
+    ops.rename("db_name", "old_table", "new_table");
+
+    Mockito.verify(mocks.updateStatement).executeUpdate("RENAME TABLE 
`old_table` TO `new_table`");
+  }
+
+  @Test
+  void testRenameRejectsCorruptedClusterMetadataBeforeMutation() throws 
Exception {
+    RenameMocks mocks = renameMocks("comment\n[Gravitino] ch.cluster= ", 
"MergeTree");
+    ExposedClickHouseTableOperations ops = newOps(mocks.dataSource);
+
+    IllegalArgumentException exception =
+        Assertions.assertThrows(
+            IllegalArgumentException.class, () -> ops.rename("db_name", 
"old_table", "new_table"));
+
+    Assertions.assertTrue(exception.getMessage().contains("missing a cluster 
name"));
+    Mockito.verify(mocks.connection, Mockito.never()).createStatement();
+  }
+
+  @Test
+  void testRenameMapsSqlException() throws Exception {
+    RenameMocks mocks = renameMocks("comment\n[Gravitino] 
ch.cluster=ck_cluster", "MergeTree");
+    SQLException sqlException = new SQLException("rename failed");
+    
Mockito.when(mocks.updateStatement.executeUpdate(Mockito.anyString())).thenThrow(sqlException);
+    ExposedClickHouseTableOperations ops = newOps(mocks.dataSource);
+
+    GravitinoRuntimeException exception =
+        Assertions.assertThrows(
+            GravitinoRuntimeException.class, () -> ops.rename("db_name", 
"old_table", "new_table"));
+
+    Assertions.assertSame(sqlException, exception.getCause());
+  }
+
   // 
---------------------------------------------------------------------------
   // extractEngineParams
   // 
---------------------------------------------------------------------------
@@ -475,4 +558,38 @@ public class TestClickHouseTableOperationsUnit {
             IllegalArgumentException.class, () -> 
newOps().callGenerateCreateTableSql(properties));
     Assertions.assertTrue(exception.getMessage().contains("balanced"));
   }
+
+  private RenameMocks renameMocks(String storedComment, String engineFull) 
throws Exception {
+    DataSource dataSource = Mockito.mock(DataSource.class);
+    Connection connection = Mockito.mock(Connection.class);
+    PreparedStatement metadataStatement = 
Mockito.mock(PreparedStatement.class);
+    ResultSet metadataResult = Mockito.mock(ResultSet.class);
+    Statement updateStatement = Mockito.mock(Statement.class);
+
+    Mockito.when(dataSource.getConnection()).thenReturn(connection);
+    
Mockito.when(connection.prepareStatement(Mockito.anyString())).thenReturn(metadataStatement);
+    Mockito.when(metadataStatement.executeQuery()).thenReturn(metadataResult);
+    Mockito.when(metadataResult.next()).thenReturn(true);
+    
Mockito.when(metadataResult.getString("COMMENT")).thenReturn(storedComment);
+    Mockito.when(metadataResult.getString("ENGINE"))
+        .thenReturn(
+            engineFull.startsWith("Distributed")
+                ? ENGINE.DISTRIBUTED.getValue()
+                : ENGINE.MERGETREE.getValue());
+    
Mockito.when(metadataResult.getString("engine_full")).thenReturn(engineFull);
+    Mockito.when(connection.createStatement()).thenReturn(updateStatement);
+    return new RenameMocks(dataSource, connection, updateStatement);
+  }
+
+  private static final class RenameMocks {
+    private final DataSource dataSource;
+    private final Connection connection;
+    private final Statement updateStatement;
+
+    private RenameMocks(DataSource dataSource, Connection connection, 
Statement updateStatement) {
+      this.dataSource = dataSource;
+      this.connection = connection;
+      this.updateStatement = updateStatement;
+    }
+  }
 }
diff --git 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/resources/remote_servers.xml
 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/resources/remote_servers.xml
index 0b791ea3be..ce38abdc61 100644
--- 
a/catalogs-contrib/catalog-jdbc-clickhouse/src/test/resources/remote_servers.xml
+++ 
b/catalogs-contrib/catalog-jdbc-clickhouse/src/test/resources/remote_servers.xml
@@ -20,7 +20,15 @@
         <gravitino_cluster>
             <shard>
                 <replica>
-                    <host>gravitino-ci-clickhouse-cluster</host>
+                    <host>gravitino-ci-clickhouse-cluster-1</host>
+                    <port>9000</port>
+                </replica>
+                <replica>
+                    <host>gravitino-ci-clickhouse-cluster-2</host>
+                    <port>9000</port>
+                </replica>
+                <replica>
+                    <host>gravitino-ci-clickhouse-cluster-3</host>
                     <port>9000</port>
                 </replica>
             </shard>
diff --git 
a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java
 
b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java
index b2b0ae50da..3b513bf800 100644
--- 
a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java
+++ 
b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/container/ContainerSuite.java
@@ -35,6 +35,7 @@ import java.nio.file.FileSystems;
 import java.nio.file.Files;
 import java.nio.file.Paths;
 import java.nio.file.attribute.PosixFilePermission;
+import java.util.ArrayList;
 import java.util.EnumMap;
 import java.util.List;
 import java.util.Map;
@@ -62,6 +63,8 @@ public class ContainerSuite implements Closeable {
   private static final String CONTAINER_NETWORK_GATEWAY = "10.20.30.1";
   private static final String CONTAINER_NETWORK_IPRANGE = "10.20.30.0/28";
   private static final String NETWORK_NAME = "gravitino-ci-network";
+  private static final int CLICKHOUSE_CLUSTER_SIZE = 3;
+  private static final String CLICKHOUSE_CLUSTER_HOST_PREFIX = 
"gravitino-ci-clickhouse-cluster-";
 
   private static Network network = null;
   private static volatile HiveContainer hiveContainer;
@@ -85,6 +88,7 @@ public class ContainerSuite implements Closeable {
   private static volatile OceanBaseContainer oceanBaseContainer;
   private static volatile ClickHouseContainer clickHouseContainer;
   private static volatile ClickHouseContainer clickHouseClusterContainer;
+  private static volatile List<ClickHouseContainer> 
clickHouseClusterContainers = List.of();
   private static volatile ZooKeeperContainer zooKeeperContainer;
 
   private static volatile GravitinoLocalStackContainer 
gravitinoLocalStackContainer;
@@ -583,30 +587,35 @@ public class ContainerSuite implements Closeable {
 
   public void startClickHouseClusterContainer(
       TestDatabaseName testDatabaseName, String remoteServersTemplatePath) {
-    if (clickHouseClusterContainer == null) {
+    if (clickHouseClusterContainers.isEmpty()) {
       synchronized (ContainerSuite.class) {
-        if (clickHouseClusterContainer == null) {
+        if (clickHouseClusterContainers.isEmpty()) {
           initIfNecessary();
           startZooKeeperContainer();
           String zkHost = zooKeeperContainer.getContainerIpAddress();
           String resolvedConfigPath = 
prepareRemoteServersConfig(remoteServersTemplatePath, zkHost);
-          ClickHouseContainer.Builder clickHouseBuilder =
-              ClickHouseContainer.builder()
-                  .withHostName("gravitino-ci-clickhouse-cluster")
-                  .withEnvVars(
-                      ImmutableMap.<String, String>builder()
-                          .put("CLICKHOUSE_PASSWORD", 
ClickHouseContainer.PASSWORD)
-                          .build())
-                  .withRemoteServersConfig(resolvedConfigPath)
-                  .withExposePorts(
-                      ImmutableSet.of(
-                          ClickHouseContainer.CLICKHOUSE_PORT,
-                          ClickHouseContainer.CLICKHOUSE_NATIVE_PORT))
-                  .withNetwork(network);
-
-          ClickHouseContainer container = 
closer.register(clickHouseBuilder.build());
-          container.start();
-          clickHouseClusterContainer = container;
+          List<ClickHouseContainer> containers = new 
ArrayList<>(CLICKHOUSE_CLUSTER_SIZE);
+          for (int node = 1; node <= CLICKHOUSE_CLUSTER_SIZE; node++) {
+            ClickHouseContainer.Builder clickHouseBuilder =
+                ClickHouseContainer.builder()
+                    .withHostName(CLICKHOUSE_CLUSTER_HOST_PREFIX + node)
+                    .withEnvVars(
+                        ImmutableMap.<String, String>builder()
+                            .put("CLICKHOUSE_PASSWORD", 
ClickHouseContainer.PASSWORD)
+                            .build())
+                    .withRemoteServersConfig(resolvedConfigPath)
+                    .withExposePorts(
+                        ImmutableSet.of(
+                            ClickHouseContainer.CLICKHOUSE_PORT,
+                            ClickHouseContainer.CLICKHOUSE_NATIVE_PORT))
+                    .withNetwork(network);
+
+            ClickHouseContainer container = 
closer.register(clickHouseBuilder.build());
+            container.start();
+            containers.add(container);
+          }
+          clickHouseClusterContainers = List.copyOf(containers);
+          clickHouseClusterContainer = clickHouseClusterContainers.get(0);
         }
       }
     }
@@ -669,6 +678,16 @@ public class ContainerSuite implements Closeable {
     return clickHouseClusterContainer;
   }
 
+  /**
+   * Returns every ClickHouse node in the repository-managed cluster fixture.
+   *
+   * @return an immutable list whose first element is also returned by {@link
+   *     #getClickHouseClusterContainer()}
+   */
+  public List<ClickHouseContainer> getClickHouseClusterContainers() {
+    return clickHouseClusterContainers;
+  }
+
   public ZooKeeperContainer getZooKeeperContainer() {
     return zooKeeperContainer;
   }

Reply via email to