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..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 @@ -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,24 @@ 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); - 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 +578,13 @@ private void createTableImplWithLock(Identifier identifier, Schema schema) { collectTableProperties(tableSchema)); } } catch (Exception e) { + // 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 55808e1e426e..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 @@ -25,9 +25,11 @@ 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; +import org.apache.paimon.schema.FileSystemSchemaManager; import org.apache.paimon.schema.Schema; import org.apache.paimon.schema.SchemaChange; import org.apache.paimon.table.Table; @@ -839,6 +841,103 @@ 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 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"; + 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 +1605,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"); + } }