From 76fb8cfcf333321695d4a9a2dd144f38d8b87e31 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 18:36:45 +0800 Subject: [PATCH 1/2] [core] Fix JDBC catalog drop, alter and create failure handling MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit dropDatabaseImpl deleted only the JDBC row sets and left the warehouse database directory behind, unlike the file system catalog, so after a cascade drop re-creating the database and a same-name table failed because the schema directory already existed. Delete the database directory too. alterDatabaseImpl never checked that the database exists: an ALTER on a missing database inserted property rows, silently materializing a phantom database. Throw DatabaseNotExistException first. createTableImplWithLock cleaned the committed schema directory only when insertTable returned false, but the realistic failure is the SQLException it throws; clean the directory in the catch as well — but only until the table row is committed, since the table keeps working when a later step like the property sync fails. Assisted-by: GLM-5.3 --- .../org/apache/paimon/jdbc/JdbcCatalog.java | 28 ++++++-- .../apache/paimon/jdbc/JdbcCatalogTest.java | 68 +++++++++++++++++++ 2 files changed, 90 insertions(+), 6 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java b/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java index e73db8a5c007..ba03e1d4d997 100644 --- a/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java +++ b/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java @@ -303,6 +303,9 @@ public void dropDatabase(String name, boolean ignoreIfNotExists, boolean cascade @Override protected void dropDatabaseImpl(String name) { + // Delete the database directory in the warehouse like the file system catalog, + // otherwise a re-created same-name database cannot re-create its tables + fileIO.deleteDirectoryQuietly(newDatabasePath(name)); // Delete table from paimon_tables execute(connections, JdbcUtils.DELETE_TABLES_SQL, catalogKey, name); // Delete properties from paimon_database_properties @@ -320,7 +323,11 @@ protected void dropDatabaseImpl(String name) { } @Override - protected void alterDatabaseImpl(String name, List changes) { + protected void alterDatabaseImpl(String name, List changes) + throws DatabaseNotExistException { + if (!JdbcUtils.databaseExists(connections, catalogKey, name)) { + throw new DatabaseNotExistException(name); + } Pair, Set> setPropertiesToRemoveKeys = PropertyChange.getSetPropertiesToRemoveKeys(changes); Map setProperties = setPropertiesToRemoveKeys.getLeft(); @@ -532,17 +539,20 @@ protected void createTableImpl(Identifier identifier, Schema schema) { } private void createTableImplWithLock(Identifier identifier, Schema schema) { + boolean registered = false; try { // create table file SchemaManager schemaManager = getSchemaManager(identifier); TableSchema tableSchema = schemaManager.createTable(schema); // Update schema metadata Path path = getTableLocation(identifier); - if (JdbcUtils.insertTable( - connections, - catalogKey, - identifier.getDatabaseName(), - identifier.getTableName())) { + registered = + JdbcUtils.insertTable( + connections, + catalogKey, + identifier.getDatabaseName(), + identifier.getTableName()); + if (registered) { LOG.debug("Successfully committed to new table: {}", identifier); } else { try { @@ -564,6 +574,12 @@ private void createTableImplWithLock(Identifier identifier, Schema schema) { collectTableProperties(tableSchema)); } } catch (Exception e) { + // the schema directory may already be committed: without this cleanup a failed + // registration leaves it behind and blocks re-creation. Once registered the + // table works even if a later step fails, so the directory must stay. + if (!registered) { + fileIO.deleteDirectoryQuietly(getTableLocation(identifier)); + } throw new RuntimeException("Failed to create table " + identifier.getFullName(), e); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java b/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java index 55808e1e426e..8e51341cc5a0 100644 --- a/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java @@ -25,6 +25,7 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.CatalogTestBase; import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.catalog.PropertyChange; import org.apache.paimon.fs.Path; import org.apache.paimon.options.CatalogOptions; import org.apache.paimon.options.Options; @@ -839,6 +840,69 @@ public void testCreateViewStoresCrossDatabaseQueryWithoutResolvingReferencedTabl assertThat(catalog.getView(identifier).query()).isEqualTo(query); } + @Test + public void testCreateTableSurvivesLatePropertySyncFailure() throws Exception { + // the property-sync step runs after the table row is committed: its failure + // must not delete the working table's directory + JdbcCatalog jdbcCatalog = initCatalogWithSync(true); + String databaseName = "late_sync_db"; + jdbcCatalog.createDatabase(databaseName, false); + dropPropertyTable(jdbcCatalog); + Identifier identifier = Identifier.create(databaseName, "t"); + + assertThatThrownBy( + () -> + jdbcCatalog.createTable( + identifier, + Schema.newBuilder() + .column("k", DataTypes.INT()) + .option("comment", "synced") + .build(), + false)) + .isInstanceOf(RuntimeException.class); + + java.nio.file.Path dir = + java.nio.file.Paths.get(jdbcCatalog.newDatabasePath(databaseName).toUri()); + assertThat(dir).exists(); + // the table is registered and readable: only the property rows are missing + assertThat(jdbcCatalog.listTables(databaseName)).contains("t"); + assertThat(jdbcCatalog.getTable(identifier)).isNotNull(); + } + + @Test + public void testDropDatabaseCascadeDeletesWarehouseDirectory() throws Exception { + String databaseName = "drop_dir_db"; + catalog.createDatabase(databaseName, false); + Identifier identifier = Identifier.create(databaseName, "t"); + catalog.createTable( + identifier, Schema.newBuilder().column("k", DataTypes.INT()).build(), false); + java.nio.file.Path dir = + java.nio.file.Paths.get( + ((JdbcCatalog) catalog).newDatabasePath(databaseName).toUri()); + assertThat(dir).exists(); + + // leaving the directory behind blocks re-creating the table after re-creating + // the database + catalog.dropDatabase(databaseName, false, true); + assertThat(dir).doesNotExist(); + } + + @Test + public void testAlterDatabaseOnMissingDatabaseThrows() throws Exception { + assertThatThrownBy( + () -> + catalog.alterDatabase( + "missing_db", + Collections.singletonList( + PropertyChange.setProperty("k", "v")), + false)) + .isInstanceOf(Catalog.DatabaseNotExistException.class); + // no phantom database materialized by the property insert + assertThatThrownBy(() -> catalog.getDatabase("missing_db")) + .isInstanceOf(Catalog.DatabaseNotExistException.class); + assertThat(catalog.listDatabases()).doesNotContain("missing_db"); + } + @Test public void testDropDatabaseCleansViewMetadata() throws Exception { String databaseName = "drop_view_db"; @@ -1506,4 +1570,8 @@ public void testAlterView() throws Exception { false)) .isInstanceOf(Catalog.DialectNotExistException.class); } + + private static void dropPropertyTable(JdbcCatalog catalog) throws Exception { + JdbcUtils.execute(catalog.getConnections(), "DROP TABLE paimon_table_properties"); + } } From 3d15caf494289b8e9b378ca71f35e6052ec18e8f Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Mon, 28 Sep 2026 02:19:22 +0800 Subject: [PATCH 2/2] fix: do not delete pre-existing on-disk table when create fails createTableImplWithLock cleaned the table directory on any pre-registration failure. When createTable fails because an on-disk schema already exists (table present on the file system but missing from the JDBC catalog, a state repairTable is meant to recover), that recursively deleted the pre-existing table's data. Guard the cleanup with a schemaCreated flag so only a directory created by this call is removed. Add a test pinning that the pre-existing directory survives and stays recoverable via repairTable. --- .../org/apache/paimon/jdbc/JdbcCatalog.java | 13 ++++--- .../apache/paimon/jdbc/JdbcCatalogTest.java | 35 +++++++++++++++++++ 2 files changed, 44 insertions(+), 4 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java b/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java index ba03e1d4d997..173022f0544d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java +++ b/paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java @@ -540,10 +540,14 @@ protected void createTableImpl(Identifier identifier, Schema schema) { private void createTableImplWithLock(Identifier identifier, Schema schema) { boolean registered = false; + boolean schemaCreated = false; try { // create table file SchemaManager schemaManager = getSchemaManager(identifier); TableSchema tableSchema = schemaManager.createTable(schema); + // this call created the schema directory; a pre-existing one makes createTable throw + // above, so schemaCreated stays false for that case and its directory is left intact + schemaCreated = true; // Update schema metadata Path path = getTableLocation(identifier); registered = @@ -574,10 +578,11 @@ private void createTableImplWithLock(Identifier identifier, Schema schema) { collectTableProperties(tableSchema)); } } catch (Exception e) { - // the schema directory may already be committed: without this cleanup a failed - // registration leaves it behind and blocks re-creation. Once registered the - // table works even if a later step fails, so the directory must stay. - if (!registered) { + // Clean up only a directory this call created but failed to register (the leaked-dir + // case the flag guards). A pre-existing on-disk table (schemaCreated == false, missing + // from the catalog and recoverable via repairTable) must survive a failed create, and + // once registered the table works even if a later step fails. + if (schemaCreated && !registered) { fileIO.deleteDirectoryQuietly(getTableLocation(identifier)); } throw new RuntimeException("Failed to create table " + identifier.getFullName(), e); diff --git a/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java b/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java index 8e51341cc5a0..54fc811915fa 100644 --- a/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java @@ -29,6 +29,7 @@ import org.apache.paimon.fs.Path; import org.apache.paimon.options.CatalogOptions; import org.apache.paimon.options.Options; +import org.apache.paimon.schema.FileSystemSchemaManager; import org.apache.paimon.schema.Schema; import org.apache.paimon.schema.SchemaChange; import org.apache.paimon.table.Table; @@ -869,6 +870,40 @@ public void testCreateTableSurvivesLatePropertySyncFailure() throws Exception { assertThat(jdbcCatalog.getTable(identifier)).isNotNull(); } + @Test + public void testCreateTablePreservesUnregisteredOnDiskTableOnFailure() throws Exception { + // A create-table failure caused by a pre-existing on-disk table (schema file present but + // no paimon_tables row) must NOT delete that directory: the state is recoverable via + // repairTable, and deleting it on a failed create would destroy the table's data. + JdbcCatalog jdbcCatalog = (JdbcCatalog) catalog; + String databaseName = "reg_fail_db"; + jdbcCatalog.createDatabase(databaseName, false); + Identifier identifier = Identifier.create(databaseName, "t"); + + // Seed an on-disk table that is missing from the JDBC catalog. createTable then fails + // inside schemaManager.createTable ("schema in filesystem exists") before registration. + Path tableLocation = jdbcCatalog.getTableLocation(identifier); + new FileSystemSchemaManager(jdbcCatalog.fileIO(), tableLocation) + .createTable(Schema.newBuilder().column("k", DataTypes.INT()).build()); + java.nio.file.Path dir = java.nio.file.Paths.get(tableLocation.toUri()); + assertThat(dir).exists(); + assertThat(jdbcCatalog.listTables(databaseName)).doesNotContain("t"); + + assertThatThrownBy( + () -> + jdbcCatalog.createTable( + identifier, + Schema.newBuilder().column("k", DataTypes.INT()).build(), + false)) + .isInstanceOf(RuntimeException.class); + + // the pre-existing directory must survive the failed create and stay recoverable + assertThat(dir).exists(); + jdbcCatalog.repairTable(identifier); + assertThat(jdbcCatalog.listTables(databaseName)).contains("t"); + assertThat(jdbcCatalog.getTable(identifier)).isNotNull(); + } + @Test public void testDropDatabaseCascadeDeletesWarehouseDirectory() throws Exception { String databaseName = "drop_dir_db";