jerryshao commented on code in PR #13499:
URL: https://github.com/apache/gravitino/pull/13499#discussion_r4142870322
##########
core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java:
##########
@@ -378,14 +378,15 @@ public boolean dropMetalake(NameIdentifier ident, boolean
force)
"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) {
+ if (force) {
+ return store.delete(ident, EntityType.METALAKE, true);
+ }
+ try {
+ return store.delete(ident, EntityType.METALAKE, false);
Review Comment:
[Question] Switching the non-force metalake drop from `cascade = true` to
`cascade = false` also drops the catalog-scoped sweeps that the old call ran.
Before this PR, non-force reached `MetalakeMetaService.deleteMetalake(ident,
true)` after the manager's own emptiness check, so the cascade branch ran all
30 cleanups -- including `softDeleteTableMetasByMetalakeId`, filesets, topics,
functions, models, views and semantic models. With `false` it now takes the
non-cascade branch, which runs only `metalakeScopedCleanups`
(MetalakeMetaService.java:254-257). Sharing `metalakeScopedCleanups` across
both branches is clearly right and keeps users, groups, roles, tags and
policies covered; my question is only about the catalog-scoped half.
In normal operation this is a no-op difference: the drop is rejected unless
`listCatalogPOsByMetalakeId` comes back empty, and a catalog delete
soft-deletes its own children in the same transaction, so there should be no
live row under a dead catalog. I could not construct a sequence that leaves one
-- the legacy-timeline GC only hard-deletes rows that are already soft-deleted.
So: was losing that sweep a deliberate simplification, or was it relied on as a
backstop for rows the existence of `OrphanedMetadataObjectRelationMapper`
suggests do occur?
Verified by: compared the pre-PR call (`store.delete(ident, METALAKE,
true)`, removed in this hunk) against both branches of
`MetalakeMetaService.deleteMetalake` (lines 222-258) and counted the cleanup
lists -- old cascade 30 entries, new `catalogScopedCleanups` 15 +
`metalakeScopedCleanups` 15.
##########
core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java:
##########
@@ -272,21 +273,50 @@ public <E extends Entity & HasIdentifier> CatalogEntity
updateCatalog(
metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
baseMetricName = "deleteCatalog")
public boolean deleteCatalog(NameIdentifier identifier, boolean cascade) {
+ return deleteCatalog(identifier, cascade, Set.of());
+ }
+
+ /**
+ * Delete a catalog after checking under the catalog row lock that every
remaining schema is among
+ * those classified as safe to discard by the manager.
+ *
+ * @param identifier the catalog identifier
+ * @param allowedSchemaIds IDs of schema entities that may be deleted with
the catalog
+ * @return true if the catalog was deleted
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "deleteCatalog")
+ public boolean deleteCatalogWithAllowedSchemas(
+ NameIdentifier identifier, Set<Long> allowedSchemaIds) {
+ return deleteCatalog(identifier, false, allowedSchemaIds);
+ }
+
+ private boolean deleteCatalog(
+ NameIdentifier identifier, boolean cascade, Set<Long> allowedSchemaIds) {
NameIdentifierUtil.checkCatalog(identifier);
String catalogName = identifier.name();
// Read the whole row, not just the ID, because the delete below needs the
version we saw.
CatalogPO catalogPO = getCatalogPOByName(identifier.namespace().level(0),
catalogName);
long catalogId = catalogPO.getCatalogId();
- if (cascade) {
+ if (cascade || !allowedSchemaIds.isEmpty()) {
SessionUtils.doMultipleWithCommit(
() -> {
// Delete the parent first, then its children. The parent delete
locks the catalog row,
// and schema writes lock that same row before they touch a
schema, so no schema can be
// added or removed after this point. Anything that goes wrong
later in this
// transaction rolls this soft delete back with it.
deleteCatalogWithVersion(identifier, catalogPO);
+ if (!cascade) {
+ List<SchemaPO> schemaPOs = listSchemaPOsForCascade(catalogId);
Review Comment:
[Nit] The schema list is read twice in the same transaction. This line calls
`listSchemaPOsForCascade(catalogId)`, and
`deleteSchemasWithVersions(identifier, catalogId)` on line 320 calls it again
internally (CatalogMetaService.java:536-537). Since the catalog row is already
locked, the second read is guaranteed to return the same rows -- it is a
redundant `SELECT` on the drop path.
`MetalakeMetaService` already has the shape that avoids this:
`deleteSchemasWithVersions(ident, listSchemaPOsForCascade(metalakeId))` takes
the list as a parameter (MetalakeMetaService.java:231). Adding the same
`List<SchemaPO>` overload here would let the check reuse the rows it just read,
and would make it textually obvious that the check and the delete operate on
one snapshot.
Verified by: read `listSchemaPOsForCascade`
(CatalogMetaService.java:554-557) and `deleteSchemasWithVersions` (536-547);
both go through `SchemaMetaMapper.listSchemaPOsByCatalogId`.
##########
core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java:
##########
@@ -272,21 +273,50 @@ public <E extends Entity & HasIdentifier> CatalogEntity
updateCatalog(
metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
baseMetricName = "deleteCatalog")
public boolean deleteCatalog(NameIdentifier identifier, boolean cascade) {
+ return deleteCatalog(identifier, cascade, Set.of());
+ }
+
+ /**
+ * Delete a catalog after checking under the catalog row lock that every
remaining schema is among
+ * those classified as safe to discard by the manager.
+ *
+ * @param identifier the catalog identifier
+ * @param allowedSchemaIds IDs of schema entities that may be deleted with
the catalog
+ * @return true if the catalog was deleted
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "deleteCatalog")
+ public boolean deleteCatalogWithAllowedSchemas(
+ NameIdentifier identifier, Set<Long> allowedSchemaIds) {
+ return deleteCatalog(identifier, false, allowedSchemaIds);
+ }
+
+ private boolean deleteCatalog(
+ NameIdentifier identifier, boolean cascade, Set<Long> allowedSchemaIds) {
NameIdentifierUtil.checkCatalog(identifier);
String catalogName = identifier.name();
// Read the whole row, not just the ID, because the delete below needs the
version we saw.
CatalogPO catalogPO = getCatalogPOByName(identifier.namespace().level(0),
catalogName);
long catalogId = catalogPO.getCatalogId();
- if (cascade) {
+ if (cascade || !allowedSchemaIds.isEmpty()) {
SessionUtils.doMultipleWithCommit(
() -> {
// Delete the parent first, then its children. The parent delete
locks the catalog row,
// and schema writes lock that same row before they touch a
schema, so no schema can be
// added or removed after this point. Anything that goes wrong
later in this
// transaction rolls this soft delete back with it.
deleteCatalogWithVersion(identifier, catalogPO);
+ if (!cascade) {
+ List<SchemaPO> schemaPOs = listSchemaPOsForCascade(catalogId);
+ if (schemaPOs.stream()
+ .anyMatch(schema ->
!allowedSchemaIds.contains(schema.getSchemaId()))) {
+ throw new NonEmptyEntityException(
+ "Entity %s has sub-entities, you should remove
sub-entities first", identifier);
+ }
Review Comment:
[Question] The new guard is schema-granular only, so a concurrently created
*grandchild* is still silently removed. Once the allowlist check passes, the
sibling operations in this same `doMultipleWithCommit` unconditionally sweep
everything under the catalog by id -- `softDeleteTableMetasByCatalogId` (line
325), filesets (337), topics (345), functions, models (376), views (383),
semantic models. A table created in an allowed schema after the manager
classified it is deleted with no error.
This is reachable in practice, not theoretical: `containsUserCreatedSchemas`
returns `false` for a Hive catalog whose only schema is `default`
(CatalogManager.java:1539-1541), for a Kafka catalog's single schema
(1534-1535) and for PostgreSQL's `public` (1536-1538) -- regardless of what
those schemas contain. So a non-force drop of such a catalog cascade-deletes
every table entity inside the built-in schema, including ones created a moment
ago by another server.
I am fairly sure this is pre-existing (main did the same via the
unconditional `store.delete(ident, CATALOG, true)`), so I am asking rather than
asserting: is closing this deliberately out of scope for this PR? If so it
would be worth naming in the `Why are the changes needed?` section alongside
the metalake fan-out follow-up, because the user-facing note ("rejects a
concurrently created child") reads as broader than what the allowlist actually
guarantees.
Verified by: read the whole cascade operation list in `deleteCatalog`
(CatalogMetaService.java:304-395) and the classification rules in
`containsUserCreatedSchemas` (CatalogManager.java:1520-1589); the allowlist is
compared against schema ids only, never against the schemas' contents.
##########
core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java:
##########
@@ -311,6 +316,23 @@ public boolean delete(NameIdentifier ident,
Entity.EntityType entityType, boolea
}
}
+ @Override
+ public boolean deleteCatalogWithAllowedSchemas(NameIdentifier ident,
Set<Long> allowedSchemaIds)
+ throws IOException {
+ if (!(backend instanceof SupportsConditionalCatalogDelete)) {
+ throw new UnsupportedOperationException(
+ "Atomic catalog delete with allowed schemas is not supported by this
backend");
Review Comment:
[Nit] This message reaches the user but omits the remedy. `CatalogManager`
rethrows `UnsupportedOperationException` unwrapped
(CatalogManager.java:1483-1485) and `ExceptionHandlers` turns it into an
`unsupportedOperation` error response (ExceptionHandlers.java:1247-1250), so an
operator running a custom `RelationalBackend` sees only "Atomic catalog delete
with allowed schemas is not supported by this backend" -- with no hint that
`force` would work.
The sibling message in `CatalogManager.java:1431-1439` gets this right ("Use
the force option, or use an entity store that implements ..."). Worth mirroring
that here, and naming the catalog, so the two paths give the same advice.
Verified by: followed the exception from this throw site through
`CatalogManager.dropCatalog`'s catch clause to `BaseExceptionHandler.handle`.
##########
core/src/main/java/org/apache/gravitino/catalog/CatalogManager.java:
##########
@@ -1414,7 +1415,38 @@ public boolean dropCatalog(NameIdentifier ident, boolean
force)
// the cache with stale data between invalidate and delete.
Map<String, String> catalogProperties =
copyProperties(catalogWrapper.catalog().entity().getProperties());
- boolean deleted = store.delete(ident, EntityType.CATALOG, true);
+ boolean deleted;
+ if (force) {
+ deleted = store.delete(ident, EntityType.CATALOG, true);
+ } else {
+ try {
+ if (schemaEntities.isEmpty()) {
+ deleted = store.delete(ident, EntityType.CATALOG, false);
+ } else {
+ Set<Long> allowedSchemaIds =
+
schemaEntities.stream().map(SchemaEntity::id).collect(Collectors.toSet());
+ if (!(store instanceof SupportsConditionalCatalogDelete)) {
+ // Fail closed: an unconditional cascade could delete a
schema created after
+ // the classification above.
+ throw new UnsupportedOperationException(
+ String.format(
+ "Catalog %s still has built-in, imported, or
externally removed "
+ + "schemas, and entity store %s cannot delete
it atomically with "
+ + "them. Use the force option, or use an
entity store that "
+ + "implements %s",
+ ident,
+ store.getClass().getName(),
+
SupportsConditionalCatalogDelete.class.getSimpleName()));
+ }
+ deleted =
+ ((SupportsConditionalCatalogDelete) store)
+ .deleteCatalogWithAllowedSchemas(ident,
allowedSchemaIds);
+ }
+ } catch (NonEmptyEntityException e) {
+ throw new NonEmptyCatalogException(
+ "Catalog %s has schemas, please drop them first or use
force option", ident);
Review Comment:
[Nit] The caught `NonEmptyEntityException` is discarded, so the message
naming which schema blocked the drop is lost. The store-side exception carries
the identifier of the offending entity ("Entity %s has sub-entities ...",
CatalogMetaService.java:316-317), which is the only place that information
exists -- the replacement message just says the catalog "has schemas".
`NonEmptyCatalogException` only declares `(String message, Object... args)`
(NonEmptyCatalogException.java:34-36), while its parent
`NonEmptyEntityException` already has the cause-taking form
(NonEmptyEntityException.java:46-48). So either add the matching constructor
and pass `e` through, or at minimum `LOG.debug("...", e)` before rethrowing --
otherwise the identifier is only recoverable by reproducing the race.
Verified by: read the message produced at CatalogMetaService.java:316-317
and the declared constructors on both exception classes.
##########
core/src/main/java/org/apache/gravitino/SupportsConditionalCatalogDelete.java:
##########
@@ -0,0 +1,40 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.gravitino;
+
+import java.io.IOException;
+import java.util.Set;
+
+/**
+ * An optional capability for deleting a catalog after checking its remaining
schemas atomically.
+ */
+public interface SupportsConditionalCatalogDelete {
+
+ /**
+ * Delete a catalog only if all remaining schemas have IDs in the allowlist.
The check and
+ * deletion must run in one transaction and be serialized with schema
creation.
+ *
+ * @param ident the catalog identifier
+ * @param allowedSchemaIds IDs of schemas that may be deleted with the
catalog
+ * @return true if the catalog was deleted
+ * @throws IOException if the store operation fails
+ */
+ boolean deleteCatalogWithAllowedSchemas(NameIdentifier ident, Set<Long>
allowedSchemaIds)
Review Comment:
[Important] The Javadoc does not state the one behaviour `CatalogManager`
actually depends on: that a disallowed schema must be signalled by throwing
`NonEmptyEntityException`, not by returning `false`.
As written, the contract is `@return true if the catalog was deleted` plus
`@throws IOException`. An implementer following only that documentation would
reasonably return `false` when the allowlist check fails.
`CatalogManager.dropCatalog` reads `false` as "nothing to delete"
(CatalogManager.java:1450 gates all cleanup on `deleted`, and line 1472 returns
it), so a non-force drop blocked by a concurrently created schema would come
back to the client as a successful no-op `false` -- i.e. "the catalog did not
exist" -- instead of the `NonEmptyCatalogException` the catch at
CatalogManager.java:1445-1447 is there to produce. That is exactly the
silent-wrong-answer the PR sets out to remove, reintroduced through a
third-party store.
Both implementations in this PR do throw (`CatalogMetaService.java:316-317`,
and `TestMemoryEntityStore.InMemoryEntityStore`), so this is a documentation
fix, not a code bug -- but this interface is the extension point the PR
description tells custom stores to implement, so the contract should be
explicit: add `@throws NonEmptyEntityException if any remaining schema is
outside {@code allowedSchemaIds}` and say that `false` means only "the catalog
was already gone".
Verified by: read the new interface in full, then traced how
`CatalogManager.dropCatalog` consumes the return value
(CatalogManager.java:1441-1472) and how both implementations signal rejection
(CatalogMetaService.java:312-319, TestMemoryEntityStore.java:180-199).
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]