This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 242a736c8d [#13435] fix(core): re-disable a user-disabled metalake
when force drop fails (#13436)
242a736c8d is described below
commit 242a736c8d69d9851aea8b4517ce58b86e59521a
Author: YangJie <[email protected]>
AuthorDate: Mon Sep 28 05:45:46 2026 -0400
[#13435] fix(core): re-disable a user-disabled metalake when force drop
fails (#13436)
### What changes were proposed in this pull request?
`dropCatalogsUnderMetalake` now tracks whether it temporarily enabled a
disabled metalake, and restores the disabled state (best effort) before
rethrowing on failure, mirroring the enable's own best-effort semantics.
A vanished metalake (`NoSuchMetalakeException`) still takes the benign
already-gone path.
### Why are the changes needed?
A force drop that failed during child cleanup left a user-disabled
metalake enabled with its catalogs partially dropped, violating the
disable the user requested.
Fix: #13435
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Added
`TestMetalakeManager.testFailedForceDropKeepsDisabledMetalakeDisabled`,
which pins that a force drop failing mid-cleanup leaves the metalake
disabled. It fails against the pre-fix code.
---
.../apache/gravitino/metalake/MetalakeManager.java | 146 ++++++++++++----
.../gravitino/metalake/TestMetalakeManager.java | 190 +++++++++++++++++++++
2 files changed, 299 insertions(+), 37 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
b/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
index df02f48466..961906bee0 100644
--- a/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
+++ b/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
@@ -29,6 +29,7 @@ import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
import org.apache.gravitino.Entity.EntityType;
import org.apache.gravitino.EntityAlreadyExistsException;
@@ -361,37 +362,46 @@ public class MetalakeManager implements
MetalakeDispatcher, Closeable {
// Mirror CatalogManager.dropCatalog → ops.dropSchema: force-drop children
through the real
// drop path so FilesetCatalogOperations / CatalogManager clean
write-through secrets (and
// managed storage). Do this before the metalake root lock to avoid
nesting tree locks.
+ boolean temporarilyEnabled = false;
if (force) {
- dropCatalogsUnderMetalake(ident);
+ temporarilyEnabled = dropCatalogsUnderMetalake(ident);
}
- return TreeLockUtils.doWithRootTreeLock(
- LockType.WRITE,
- () -> {
- try {
- boolean inUse = metalakeInUse(store, ident);
- if (inUse && !force) {
- throw new MetalakeInUseException(
- "Metalake %s is in use, please disable it first or use force
option", ident);
- }
+ try {
+ return TreeLockUtils.doWithRootTreeLock(
+ LockType.WRITE,
+ () -> {
+ try {
+ boolean inUse = metalakeInUse(store, ident);
+ if (inUse && !force) {
+ throw new MetalakeInUseException(
+ "Metalake %s is in use, please disable it first or use
force option", ident);
+ }
- List<CatalogEntity> catalogEntities =
- store.list(Namespace.of(ident.name()), CatalogEntity.class,
EntityType.CATALOG);
- if (!catalogEntities.isEmpty() && !force) {
- throw new NonEmptyMetalakeException(
- "Metalake %s has catalogs, please drop them first or use
force option", ident);
- }
+ List<CatalogEntity> catalogEntities =
+ store.list(Namespace.of(ident.name()), CatalogEntity.class,
EntityType.CATALOG);
+ if (!catalogEntities.isEmpty() && !force) {
+ throw new NonEmptyMetalakeException(
+ "Metalake %s has catalogs, please drop them first or use
force option", ident);
+ }
- return store.delete(ident, EntityType.METALAKE, true);
- } catch (NoSuchMetalakeException | NoSuchEntityException e) {
- // Another server may have completed the drop after the initial
existence check.
- // Dropping an already-removed metalake remains an idempotent
false result.
- return false;
+ return store.delete(ident, EntityType.METALAKE, true);
+ } catch (NoSuchMetalakeException | NoSuchEntityException e) {
+ // Another server may have completed the drop after the initial
existence check.
+ // Dropping an already-removed metalake remains an idempotent
false result.
+ return false;
- } catch (IOException e) {
- throw new RuntimeException(e);
- }
- });
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ });
+ } catch (RuntimeException e) {
+ // Phase 1 briefly re-enabled a user-disabled metalake to clean up its
catalogs. If the
+ // metalake delete above then fails, the metalake still exists, so
restore its disabled
+ // state rather than leaving the user's disable silently undone.
+ restoreDisabledState(ident, temporarilyEnabled);
+ throw e;
+ }
}
/**
@@ -401,32 +411,93 @@ public class MetalakeManager implements
MetalakeDispatcher, Closeable {
*
* <p>Callers typically {@code disableMetalake} before force-drop. {@link
* CatalogManager#dropCatalog} requires catalog {@code
metalake-in-use=true}, so a disabled
- * metalake is briefly re-enabled for child cleanup. The metalake entity is
deleted immediately
- * afterward, so the temporary enable is not restored.
+ * metalake is briefly re-enabled for child cleanup. If the cleanup fails it
is re-disabled here
+ * (best effort); if the cleanup succeeds this returns whether the metalake
was temporarily
+ * enabled, so {@code dropMetalake} can re-disable it should the metalake
delete then fail. Either
+ * way a user-disabled metalake does not stay enabled after a failed force
drop.
+ *
+ * @param metalakeIdent the metalake whose child catalogs are force-dropped
+ * @return {@code true} if this temporarily enabled a user-disabled
metalake, so the caller must
+ * restore the disabled state if the subsequent metalake delete fails
*/
- private void dropCatalogsUnderMetalake(NameIdentifier metalakeIdent) {
+ private boolean dropCatalogsUnderMetalake(NameIdentifier metalakeIdent) {
if (catalogManager == null) {
- return;
+ return false;
}
+ // True once this call has flipped the metalake to in-use. Set under the
metalake write lock
+ // right after the store update (before the non-atomic catalog-status
update), so a failure in
+ // either the enable itself or the later drops re-disables a metalake this
operation enabled,
+ // while a concurrent enable leaves it false so we never re-disable one we
did not enable.
+ AtomicBoolean weEnabled = new AtomicBoolean(false);
try {
- if (!metalakeInUse(store, metalakeIdent)) {
- enableMetalake(metalakeIdent);
- }
- List<CatalogEntity> catalogs =
- store.list(Namespace.of(metalakeIdent.name()), CatalogEntity.class,
EntityType.CATALOG);
- for (CatalogEntity catalog : catalogs) {
- catalogManager.dropCatalog(
- NameIdentifier.of(metalakeIdent.name(), catalog.name()), true /*
force */);
+ try {
+ enableMetalakeIfDisabled(metalakeIdent, weEnabled);
+ List<CatalogEntity> catalogs =
+ store.list(Namespace.of(metalakeIdent.name()),
CatalogEntity.class, EntityType.CATALOG);
+ for (CatalogEntity catalog : catalogs) {
+ catalogManager.dropCatalog(
+ NameIdentifier.of(metalakeIdent.name(), catalog.name()), true /*
force */);
+ }
+ // Report whether we enabled the metalake so the caller can re-disable
it if the metalake
+ // delete that follows fails.
+ return weEnabled.get();
+ } catch (NoSuchMetalakeException e) {
+ // Metalake is already gone; dropMetalake will return false. Nothing
to restore.
+ throw e;
+ } catch (IOException | RuntimeException e) {
+ restoreDisabledState(metalakeIdent, weEnabled.get());
+ throw e;
}
} catch (NoSuchMetalakeException e) {
// Metalake is already gone; dropMetalake will return false.
+ return false;
} catch (IOException e) {
throw new RuntimeException(e);
}
}
+ /**
+ * Best effort: undo a temporary enable this force drop performed, so a
failed force drop keeps a
+ * user-disabled metalake disabled. Only re-disables when {@code weEnabled}
is true, i.e. this
+ * operation is the one that enabled the metalake, so it never overwrites a
concurrent user's
+ * enable.
+ */
+ private void restoreDisabledState(NameIdentifier metalakeIdent, boolean
weEnabled) {
+ if (!weEnabled) {
+ return;
+ }
+ try {
+ disableMetalake(metalakeIdent);
+ } catch (Exception restoreFailure) {
+ LOG.warn(
+ "Failed to restore the disabled state of metalake {} after a failed
force drop; "
+ + "the metalake may remain enabled",
+ metalakeIdent,
+ restoreFailure);
+ }
+ }
+
@Override
public void enableMetalake(NameIdentifier ident) throws
NoSuchMetalakeException {
+ enableMetalakeIfDisabled(ident, new AtomicBoolean());
+ }
+
+ /**
+ * Enables the metalake only if it is currently disabled, deciding and
writing atomically under
+ * the metalake write lock.
+ *
+ * <p>{@code enabledByUs} is set to true the moment this call flips the
metalake to in-use, before
+ * the non-atomic catalog-status update. A force-drop caller reads it to
decide whether it must
+ * re-disable the metalake on failure: keying that on this flag (rather than
a prior {@code
+ * metalakeInUse} read) means a concurrent enable cannot make the caller
claim, and later undo, an
+ * enable it did not perform, and a failure partway through the enable still
leaves the caller
+ * able to re-disable.
+ *
+ * @param ident the metalake to enable
+ * @param enabledByUs set to true iff this call performed the disabled to
in-use flip
+ */
+ private void enableMetalakeIfDisabled(NameIdentifier ident, AtomicBoolean
enabledByUs)
+ throws NoSuchMetalakeException {
TreeLockUtils.doWithTreeLock(
ident,
LockType.WRITE,
@@ -453,6 +524,7 @@ public class MetalakeManager implements MetalakeDispatcher,
Closeable {
return builder.build();
});
+ enabledByUs.set(true);
// The only problem is that we can't make sure we can change all
catalog properties
// in a transaction. If any catalog property update fails, the
metalake is already
diff --git
a/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
b/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
index a9a69ec126..d27b096c55 100644
--- a/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
+++ b/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
@@ -352,6 +352,196 @@ public class TestMetalakeManager {
store.close();
}
+ @Test
+ public void testFailedForceDropKeepsDisabledMetalakeDisabled() throws
Exception {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ Object originalEnvCatalogManager =
+ FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager",
true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager",
catalogManager, true);
+ try {
+ MetalakeManager manager =
+ new MetalakeManager(entityStore, new RandomIdGenerator(),
catalogManager);
+ NameIdentifier ident = NameIdentifier.of("metalake_force_drop_failure");
+ manager.createMetalake(ident, "comment", ImmutableMap.of());
+ manager.disableMetalake(ident);
+ Assertions.assertFalse(MetalakeManager.metalakeInUse(entityStore,
ident));
+
+ entityStore.put(
+ CatalogEntity.builder()
+ .withId(new RandomIdGenerator().nextId())
+ .withName("catalog1")
+ .withNamespace(Namespace.of(ident.name()))
+ .withType(Catalog.Type.RELATIONAL)
+ .withProvider("hive")
+ .withAuditInfo(
+
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
+ .build(),
+ false);
+
+ Mockito.doThrow(new RuntimeException("catalog drop failed"))
+ .when(catalogManager)
+ .dropCatalog(Mockito.any(), Mockito.anyBoolean());
+
+ Assertions.assertThrows(RuntimeException.class, () ->
manager.dropMetalake(ident, true));
+
+ // A failed force drop must not leave a user-disabled metalake
re-enabled.
+ Assertions.assertFalse(
+ MetalakeManager.metalakeInUse(entityStore, ident),
+ "failed force drop must keep the user-disabled metalake disabled");
+ } finally {
+ FieldUtils.writeField(
+ GravitinoEnv.getInstance(), "catalogManager",
originalEnvCatalogManager, true);
+ }
+ }
+
+ @Test
+ public void testForceDropRestoresDisabledMetalakeWhenDeleteFails() throws
Exception {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ Object originalEnvCatalogManager =
+ FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager",
true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager",
catalogManager, true);
+ try {
+ EntityStore spyStore = Mockito.spy(entityStore);
+ MetalakeManager manager =
+ new MetalakeManager(spyStore, new RandomIdGenerator(),
catalogManager);
+ NameIdentifier ident =
NameIdentifier.of("metalake_force_drop_delete_failure");
+ manager.createMetalake(ident, "comment", ImmutableMap.of());
+ manager.disableMetalake(ident);
+ Assertions.assertFalse(MetalakeManager.metalakeInUse(spyStore, ident));
+
+ // No catalogs, so phase-1 catalog cleanup succeeds and re-enables the
metalake; make the
+ // phase-2 metalake delete fail so the temporary enable outlives the
cleanup handler.
+ Mockito.doThrow(new IOException("metalake delete failed"))
+ .when(spyStore)
+ .delete(ident, EntityType.METALAKE, true);
+
+ Assertions.assertThrows(RuntimeException.class, () ->
manager.dropMetalake(ident, true));
+
+ Assertions.assertFalse(
+ MetalakeManager.metalakeInUse(spyStore, ident),
+ "a force drop whose metalake delete fails must keep the
user-disabled metalake disabled");
+ } finally {
+ FieldUtils.writeField(
+ GravitinoEnv.getInstance(), "catalogManager",
originalEnvCatalogManager, true);
+ }
+ }
+
+ @Test
+ public void testFailedForceDropDoesNotDisableMetalakeItDidNotEnable() throws
Exception {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ Object originalEnvCatalogManager =
+ FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager",
true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager",
catalogManager, true);
+ try {
+ MetalakeManager manager =
+ new MetalakeManager(entityStore, new RandomIdGenerator(),
catalogManager);
+ NameIdentifier ident =
NameIdentifier.of("metalake_force_drop_already_enabled");
+ manager.createMetalake(ident, "comment", ImmutableMap.of());
+ // Left enabled: this stands in for a metalake a concurrent user has
enabled. The force drop
+ // does not enable it, so a failure must not re-disable it and clobber
that user's enable.
+ Assertions.assertTrue(MetalakeManager.metalakeInUse(entityStore, ident));
+
+ entityStore.put(
+ CatalogEntity.builder()
+ .withId(new RandomIdGenerator().nextId())
+ .withName("catalog1")
+ .withNamespace(Namespace.of(ident.name()))
+ .withType(Catalog.Type.RELATIONAL)
+ .withProvider("hive")
+ .withAuditInfo(
+
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
+ .build(),
+ false);
+
+ Mockito.doThrow(new RuntimeException("catalog drop failed"))
+ .when(catalogManager)
+ .dropCatalog(Mockito.any(), Mockito.anyBoolean());
+
+ Assertions.assertThrows(RuntimeException.class, () ->
manager.dropMetalake(ident, true));
+
+ Assertions.assertTrue(
+ MetalakeManager.metalakeInUse(entityStore, ident),
+ "a failed force drop must not disable a metalake it did not enable");
+ } finally {
+ FieldUtils.writeField(
+ GravitinoEnv.getInstance(), "catalogManager",
originalEnvCatalogManager, true);
+ }
+ }
+
+ @Test
+ public void testFailedForceDropRestoresWhenTemporaryEnableFailsMidway()
throws Exception {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ Object originalEnvCatalogManager =
+ FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager",
true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager",
catalogManager, true);
+ try {
+ MetalakeManager manager =
+ new MetalakeManager(entityStore, new RandomIdGenerator(),
catalogManager);
+ NameIdentifier ident =
NameIdentifier.of("metalake_force_drop_enable_midfail");
+ manager.createMetalake(ident, "comment", ImmutableMap.of());
+ manager.disableMetalake(ident);
+ Assertions.assertFalse(MetalakeManager.metalakeInUse(entityStore,
ident));
+
+ entityStore.put(
+ CatalogEntity.builder()
+ .withId(new RandomIdGenerator().nextId())
+ .withName("catalog1")
+ .withNamespace(Namespace.of(ident.name()))
+ .withType(Catalog.Type.RELATIONAL)
+ .withProvider("hive")
+ .withAuditInfo(
+
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
+ .build(),
+ false);
+
+ // The temporary enable flips the metalake in-use, then fails while
propagating that to the
+ // catalog. A metalake this force drop enabled must not be left enabled.
+ Mockito.doThrow(new RuntimeException("catalog in-use update failed"))
+ .when(catalogManager)
+ .setMetalakeInUseStatus(Mockito.any(), Mockito.anyBoolean());
+
+ Assertions.assertThrows(RuntimeException.class, () ->
manager.dropMetalake(ident, true));
+
+ Assertions.assertFalse(
+ MetalakeManager.metalakeInUse(entityStore, ident),
+ "a temporary enable that fails partway must not leave the metalake
enabled");
+ } finally {
+ FieldUtils.writeField(
+ GravitinoEnv.getInstance(), "catalogManager",
originalEnvCatalogManager, true);
+ }
+ }
+
+ @Test
+ public void testForceDropDeleteFailureDoesNotDisableAlreadyEnabledMetalake()
throws Exception {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ Object originalEnvCatalogManager =
+ FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager",
true);
+ FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager",
catalogManager, true);
+ try {
+ EntityStore spyStore = Mockito.spy(entityStore);
+ MetalakeManager manager =
+ new MetalakeManager(spyStore, new RandomIdGenerator(),
catalogManager);
+ NameIdentifier ident =
NameIdentifier.of("metalake_force_drop_enabled_delete_fail");
+ manager.createMetalake(ident, "comment", ImmutableMap.of());
+ // Left enabled (a concurrent user's enable); no catalogs, so phase 1
succeeds without this
+ // operation enabling it, then the phase-2 metalake delete fails.
+ Assertions.assertTrue(MetalakeManager.metalakeInUse(spyStore, ident));
+
+ Mockito.doThrow(new IOException("metalake delete failed"))
+ .when(spyStore)
+ .delete(ident, EntityType.METALAKE, true);
+
+ Assertions.assertThrows(RuntimeException.class, () ->
manager.dropMetalake(ident, true));
+
+ Assertions.assertTrue(
+ MetalakeManager.metalakeInUse(spyStore, ident),
+ "a delete-failed force drop must not disable a metalake it did not
enable");
+ } finally {
+ FieldUtils.writeField(
+ GravitinoEnv.getInstance(), "catalogManager",
originalEnvCatalogManager, true);
+ }
+ }
+
private void testProperties(Map<String, String> expectedProps, Map<String,
String> testProps) {
expectedProps.forEach(
(k, v) -> {