jerryshao commented on code in PR #13198:
URL: https://github.com/apache/gravitino/pull/13198#discussion_r4080628719
##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +249,102 @@ protected StringIdentifier
getStringIdFromProperties(Map<String, String> propert
}
}
+ /**
+ * Wraps an updater so the store rejects the update, before writing
anything, when the row under
+ * the name is not the entity the external catalog reported.
+ *
+ * <p>The updater runs inside the store's update, after the current row is
read and before the
+ * version-checked write. Throwing here therefore aborts the transaction
with nothing written,
+ * which is what a post-write id comparison cannot do.
+ *
+ * @param expectedId the id read from the external catalog
+ * @param updater the update to apply when the ids match
+ * @param <E> the entity type
+ * @return the guarded updater
+ */
+ protected static <E extends Entity & HasIdentifier> Function<E, E>
requireEntityId(
+ long expectedId, Function<E, E> updater) {
+ return entity -> {
+ if (entity.id() != expectedId) {
+ throw new EntityIdMismatchException(entity.id(), expectedId);
+ }
+ return updater.apply(entity);
+ };
+ }
+
+ /**
+ * Reads the id and store version of a registration before an
external-catalog call, so a later
+ * store write can be fenced on it.
+ *
+ * @param ident the entity identifier
+ * @param type the entity type
+ * @return the observed id and version, or null when nothing is registered
under the name or the
+ * store cannot read versions
+ */
+ @Nullable
+ protected EntityVersion observeRegistration(NameIdentifier ident,
Entity.EntityType type) {
+ try {
+ return store.getVersion(ident, type);
+ } catch (NoSuchEntityException | UnsupportedOperationException e) {
+ return null;
Review Comment:
[Important] `observeRegistration` collapses two different situations into
`null`: "nothing is registered under this name" (`NoSuchEntityException`) and
"this store cannot read versions" (`UnsupportedOperationException`).
`deleteObservedRegistration` then returns early on `observed == null` (line
316), so for a store that does not implement `getVersion` the unconditional
fallback at line 327 is unreachable — the only way to reach it would be a store
that implements `getVersion` but not the version-checked `delete`.
That matters because both `EntityStore.getVersion` and
`RelationalBackend.getVersion` default to throwing
`UnsupportedOperationException`, and the backend class is instantiated
reflectively from config (`RelationalEntityStore.java:176-190`,
`Class.forName(className)`). A `RelationalBackend` that does not override
`getVersion` would therefore make every external-backed drop —
`TableOperationDispatcher.java:432` and `:481`,
`SchemaOperationDispatcher.java:586`, `TopicOperationDispatcher.java:258`,
`ViewOperationDispatcher.java:348` — stop deleting the store registration
altogether and only log a warning, instead of falling back to today's
unconditional delete as the fallback clearly intends. It also silently flips
`dropTopic`'s return value to `false`, since `droppedFromStore` is what that
method returns for managed topics.
Suggest distinguishing the two cases, e.g. have `observeRegistration` signal
"unsupported" separately (a sentinel `EntityVersion`, an `Optional`, or a
`supportsVersions()` probe) so `deleteObservedRegistration` can take the
unconditional path for it while still skipping the delete when there is
genuinely no registration.
Related, on the same helper: the `catch (UnsupportedOperationException)` at
line 326 wraps the whole `store.delete(ident, type, cascade, expected)` call,
so a `UnsupportedOperationException` raised anywhere deeper in the store would
be downgraded to an unconditional delete by name — exactly the write this PR is
trying to prevent. Narrowing the fallback to a capability check rather than a
catch would avoid that too.
Verified by: read `OperationDispatcher.java:285-336` at HEAD `cc52211`; read
the `getVersion` defaults in `EntityStore.java:241-243` and
`RelationalBackend.java:154-156`; read
`RelationalEntityStore.createRelationalEntityBackend` at lines 176-190; and
grepped `implements EntityStore|implements RelationalBackend` across the repo
(only `RelationalEntityStore` and `JDBCBackend` in main, plus two test stores).
##########
core/src/main/java/org/apache/gravitino/catalog/SchemaOperationDispatcher.java:
##########
@@ -556,33 +559,31 @@ public boolean dropSchema(NameIdentifier ident, boolean
cascade) throws NonEmpty
schemaProperties = new HashMap<>(schemaEntity.properties());
}
+ // For managed schema, we don't need to drop the schema from the
store again.
+ boolean isManagedSchema = isManagedEntity(catalogIdent,
Capability.Scope.SCHEMA);
+ // Read the registration before the external call, so the store
delete below can only
+ // remove the row this drop started with and never one re-created
under the same name.
+ EntityVersion observed = isManagedSchema ? null :
observeRegistration(ident, SCHEMA);
boolean droppedFromCatalog =
doWithCatalog(
catalogIdent,
c -> c.doWithSchemaOps(s -> s.dropSchema(ident, cascade)),
NonEmptySchemaException.class,
RuntimeException.class);
- // For managed schema, we don't need to drop the schema from the
store again.
- boolean isManagedSchema = isManagedEntity(catalogIdent,
Capability.Scope.SCHEMA);
if (isManagedSchema) {
if (droppedFromCatalog) {
secretManager.deleteSecretsFromProperties(schemaProperties);
}
return droppedFromCatalog;
}
- // A non-cascading drop preserves a missing registration because the
source schema
- // may have been renamed. An explicit cascading drop also removes
stale metadata.
+ // A non-cascading false result may mean the external schema was
renamed, so preserve
+ // its registration. An explicit cascade also removes stale
metadata, but only if the
+ // registration is still the one observed before the external call.
boolean droppedFromStore = false;
if (droppedFromCatalog || cascade) {
- try {
- droppedFromStore = store.delete(ident, SCHEMA, true);
- } catch (NoSuchEntityException e) {
- LOG.warn("The schema to be dropped does not exist in the store:
{}", ident, e);
- } catch (Exception e) {
- throw new RuntimeException(e);
- }
+ droppedFromStore = deleteObservedRegistration(ident, SCHEMA, true,
observed);
Review Comment:
[Question] The cascade here is fenced on the schema row only.
`SchemaMetaService.deleteSchema(identifier, cascade, expected)` compares
`expected` against the schema PO's `(schemaId, currentVersion)` and then, on
`cascade`, deletes the children unconditionally.
Inserting a table into a schema does not bump that schema's
`current_version`, so a table registered after `observeRegistration` at line
566 and before this call passes the fence and is deleted, even though no drop
ever observed it. The schema re-creation case is genuinely covered (a new
schema gets a new id, so the check fires), so this is the same ABA one level
down rather than a regression — is it in scope for this PR, or deliberately
left to the reconcile pass the helper's javadoc mentions?
Verified by: read `SchemaMetaService.deleteSchema` at
`SchemaMetaService.java:302-320` and the cascade branch that follows, plus
`SchemaOperationDispatcher.java:562-587`, at HEAD `cc52211`.
##########
core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java:
##########
@@ -311,6 +312,25 @@ public boolean delete(NameIdentifier ident,
Entity.EntityType entityType, boolea
}
}
+ @Override
+ public EntityVersion getVersion(NameIdentifier ident, Entity.EntityType
entityType)
+ throws IOException {
+ // Always read through: a cached entity does not carry the store version,
and a stale cache
+ // entry must never be the basis of a version check.
+ return backend.getVersion(ident, entityType);
+ }
+
+ @Override
+ public boolean delete(
+ NameIdentifier ident, Entity.EntityType entityType, boolean cascade,
EntityVersion expected)
+ throws IOException {
+ try {
+ return backend.delete(ident, entityType, cascade, expected);
+ } finally {
+ cache.invalidate(ident, entityType);
Review Comment:
[Important] This is the only mutating path in `RelationalEntityStore` that
calls `cache.invalidate(...)` directly instead of the `invalidateCache(...)`
helper, so it skips `cacheInvalidationEpoch.incrementAndGet()`.
That epoch exists precisely to close this window: `batchGet` samples the
epoch before reading the backend (line 273) and refuses to write a fetched
entity back into the cache if the epoch moved in between (lines 279, 288-294).
Because the version-checked delete never advances the epoch, a `batchGet` that
started before it can repopulate the cache with the just-deleted entity after
`cache.invalidate` ran, and the stale copy then survives until the TTL. Since
this new method is exactly the one the drop paths now use
(`OperationDispatcher.deleteObservedRegistration` -> `store.delete(ident, type,
cascade, expected)`), the practical effect is that a successfully dropped
table/schema/topic/view can stay readable from the cache — which undoes part of
what the PR is fixing.
Every other mutating method here uses the helper: lines 165, 222, 233, 311
(the unconditional `delete`), 344 (`deleteAndGet`), 452-453, 473-475, 501, 534,
557.
Suggested fix: `invalidateCache(ident, entityType);`
Verified by: read `RelationalEntityStore.java` at HEAD `cc52211` — the
`invalidateCache` definition at lines 603-607, the epoch checks in `batchGet`
at 273-294, and grepped every `invalidateCache|cache.invalidate` occurrence in
the file; line 330 is the only direct `cache.invalidate` call outside the
helper.
--
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]