Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 27 additions & 6 deletions paimon-core/src/main/java/org/apache/paimon/jdbc/JdbcCatalog.java
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -320,7 +323,11 @@ protected void dropDatabaseImpl(String name) {
}

@Override
protected void alterDatabaseImpl(String name, List<PropertyChange> changes) {
protected void alterDatabaseImpl(String name, List<PropertyChange> changes)
throws DatabaseNotExistException {
if (!JdbcUtils.databaseExists(connections, catalogKey, name)) {
throw new DatabaseNotExistException(name);
}
Pair<Map<String, String>, Set<String>> setPropertiesToRemoveKeys =
PropertyChange.getSetPropertiesToRemoveKeys(changes);
Map<String, String> setProperties = setPropertiesToRemoveKeys.getLeft();
Expand Down Expand Up @@ -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 {
Expand All @@ -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);
}
}
Expand Down
103 changes: 103 additions & 0 deletions paimon-core/src/test/java/org/apache/paimon/jdbc/JdbcCatalogTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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");
}
}
Loading