jerryshao commented on code in PR #13198:
URL: https://github.com/apache/gravitino/pull/13198#discussion_r4121058676
##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +256,174 @@ 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 expected id of the entity being updated
+ * @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
+ * @throws UnsupportedOperationException if the store cannot read versions
+ */
+ @Nullable
+ protected EntityVersion observeRegistration(NameIdentifier ident,
Entity.EntityType type) {
+ try {
+ return store.getVersion(ident, type);
+ } catch (NoSuchEntityException e) {
+ return null;
+ } catch (IOException e) {
+ throw new RuntimeException("Failed to read the registration of " +
ident, e);
+ }
+ }
+
+ /**
+ * Deletes the registration observed before an external-catalog call, and
only that one.
+ *
+ * <p>The external drop has already succeeded when this runs, so a conflict
is not reported to the
+ * caller: a failed request would invite a retry, and by now the name may
belong to a newly
+ * created object that the retry would drop. Instead:
+ *
+ * <ul>
+ * <li>A registration with another id belongs to a newer incarnation. It
is kept, and a
+ * reconcile pass, not this drop, decides what happens to it.
+ * <li>The same id with a newer version means a concurrent store write to
the entity that was
+ * just dropped. The fenced delete is retried, because giving up would
leave a registration
+ * whose external object is gone, and a retried drop would no longer
reach the store.
+ * </ul>
+ *
+ * @param ident the entity identifier
+ * @param type the entity type
+ * @param cascade whether to delete the children as well
+ * @param observed the id and version read before the external call, or null
when there was no
+ * registration to delete
+ * @return true if the observed registration was deleted
+ * @throws OptimisticLockException if the observed entity kept changing on
every attempt
+ * @throws UnsupportedOperationException if the store cannot delete with a
version check
+ */
+ protected boolean deleteObservedRegistration(
Review Comment:
[Important] This helper is now the only fenced post-external-drop delete,
but two other paths that delete a registration after an external drop has
already succeeded still delete by name, and they are on the Iceberg REST
service, not a corner case.
`IcebergTableHookDispatcher.dropTable` calls `dispatcher.dropTable(...)` and
then `bestEffortReconcileTableEntity`, which calls `deleteTableEntity` ->
`store.delete(ident, Entity.EntityType.TABLE)` with no `expected`
(`IcebergTableHookDispatcher.java:288-291`).
`IcebergViewHookDispatcher.dropView` does the same via `deleteViewEntity`
(`IcebergViewHookDispatcher.java:254-257`). The comment right above each call
names this PR's race explicitly: "without a distributed TreeLock, another node
may recreate the same table between the drop above and the EntityStore delete,
leaving a stale Gravitino entity if we blindly delete"
(`IcebergTableHookDispatcher.java:105-107`). The mitigation there is a
`tableExists` probe followed by `importTableEntity`
(`IcebergTableHookDispatcher.java:263-267`), which is the ABA this PR closes
elsewhere: if the name was re-created in between, the by-name delete removes
the *new* registration with its owner, tags and grants, and the re-import then
registers th
e object under a fresh id, so the attachments are gone rather than merely
stale.
This module is already in the change set — the PR updated
`TestIcebergNamespaceHookDispatcher` for the cleaner's new 4-argument delete —
and `EntityStore.getVersion` plus the fenced `delete` are available to it.
Suggest observing the registration before `dispatcher.dropTable(...)` /
`dispatcher.dropView(...)` and passing it to the fenced delete, the same shape
as `TableOperationDispatcher.dropTable`. If that is deliberately deferred, it
is worth saying so in the description, because as it stands the Iceberg REST
drop path keeps the behaviour the rest of the PR removes.
Verified by: read `IcebergTableHookDispatcher.java:100-300` and
`IcebergViewHookDispatcher.java:100-268` in full at HEAD `bdc3ca0`; grepped
every `store.put`/`store.delete` under
`core/src/main/java/org/apache/gravitino/catalog/` and
`iceberg/iceberg-rest-server/src/main/java/`; confirmed the two Iceberg call
sites are unchanged by this PR with `git diff origin/main...HEAD --name-only`,
which lists only the three Iceberg *test* files.
##########
core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java:
##########
@@ -424,15 +425,7 @@ public boolean dropTable(NameIdentifier ident) {
// Gravitino-only metadata. A true out-of-band drop can therefore
leave a stale
// registration that requires separate cleanup.
if (droppedFromCatalog) {
- try {
- store.delete(ident, TABLE);
- } catch (OptimisticLockException e) {
- throw e;
- } catch (NoSuchEntityException e) {
- LOG.warn("The table to be dropped does not exist in the store:
{}", ident, e);
- } catch (Exception e) {
- throw new RuntimeException(e);
- }
+ deleteObservedRegistration(ident, TABLE, false, observed);
Review Comment:
[Important] The fence here is id-only, and the import paths can hand a *new*
object the id this drop observed, so the fence passes and deletes the new
registration. That is the second hazard the description's "Why" section names
("Upserting a new object into a stale row can also reuse the old ID"), fixed
for creates via `putCreatedEntity` but not for imports.
`importTable` still does `store.put(tableEntity, true /* overwrite */)`
(`TableOperationDispatcher.java:562`), and the overwrite is an upsert on the
natural key that does **not** assign `table_id`:
`insertTableMetaOnDuplicateKeyUpdate` updates `table_name`, the parent ids,
`audit_info`, the versions and `deleted_at`, and leaves `table_id` at the
stored value (`TableMetaBaseSQLProvider.java:241-251`). The pre-existing test
`testNaturalKeyOverwriteUsesPersistedTableId` asserts exactly this for
MySQL/H2: after `insertTable(replacement, true)` with a fresh random id,
`stored.id()` still equals `original.id()`
(`TestTableMetaService.java:272-307`).
Failure scenario: store row for `t` has id X. Node A starts `dropTable(t)`,
observes X, and the external drop succeeds. `t` is then re-created out of band
in the source catalog. Node B serves `loadTable(t)`, finds no
`gravitino.identifier`, takes `uid = idGenerator.nextId()`
(`TableOperationDispatcher.java:541`) and upserts — the row keeps id X. Node A
resumes: `deleteObservedRegistration` sees id X under the name,
`checkExpectedIdentity` passes, and the live table's freshly imported
registration is deleted along with its columns, owner, tags and grants. The
version no longer guards this, since `OccWriteSupport.checkExpectedIdentity`
compares ids only and the CAS uses the row's current version.
`importSchema` (`SchemaOperationDispatcher.java:681`), `importTopic`
(`TopicOperationDispatcher.java:311`) and `importView`
(`ViewOperationDispatcher.java:578`) have the same shape. Suggest routing the
import paths through `putCreatedEntity`, or at least not upserting onto a row
whose id differs from the entity being imported, so an import always produces a
row whose id identifies the incarnation it describes.
Verified by: read `importTable` at `TableOperationDispatcher.java:513-571`
and `insertTable` at `TableMetaService.java:111-193` in full at HEAD `bdc3ca0`;
read the upsert SQL at `TableMetaBaseSQLProvider.java:224-252`; read
`testNaturalKeyOverwriteUsesPersistedTableId`, which is skipped only for
PostgreSQL (`TestTableMetaService.java:271-308`); grepped `store.put(` across
`core/src/main/java/org/apache/gravitino/catalog/` — the four `overwrite =
true` call sites left are the four import paths.
##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +256,174 @@ 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 expected id of the entity being updated
+ * @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
+ * @throws UnsupportedOperationException if the store cannot read versions
+ */
+ @Nullable
+ protected EntityVersion observeRegistration(NameIdentifier ident,
Entity.EntityType type) {
+ try {
+ return store.getVersion(ident, type);
+ } catch (NoSuchEntityException e) {
+ return null;
+ } catch (IOException e) {
+ throw new RuntimeException("Failed to read the registration of " +
ident, e);
+ }
+ }
+
+ /**
+ * Deletes the registration observed before an external-catalog call, and
only that one.
+ *
+ * <p>The external drop has already succeeded when this runs, so a conflict
is not reported to the
+ * caller: a failed request would invite a retry, and by now the name may
belong to a newly
+ * created object that the retry would drop. Instead:
+ *
+ * <ul>
+ * <li>A registration with another id belongs to a newer incarnation. It
is kept, and a
+ * reconcile pass, not this drop, decides what happens to it.
+ * <li>The same id with a newer version means a concurrent store write to
the entity that was
+ * just dropped. The fenced delete is retried, because giving up would
leave a registration
+ * whose external object is gone, and a retried drop would no longer
reach the store.
+ * </ul>
+ *
+ * @param ident the entity identifier
+ * @param type the entity type
+ * @param cascade whether to delete the children as well
+ * @param observed the id and version read before the external call, or null
when there was no
+ * registration to delete
+ * @return true if the observed registration was deleted
+ * @throws OptimisticLockException if the observed entity kept changing on
every attempt
+ * @throws UnsupportedOperationException if the store cannot delete with a
version check
+ */
+ protected boolean deleteObservedRegistration(
+ NameIdentifier ident,
+ Entity.EntityType type,
+ boolean cascade,
+ @Nullable EntityVersion observed) {
+ if (observed == null) {
+ LOG.warn(
+ "No {} registration was found for {} before the external drop;
leaving the store alone",
+ type.name().toLowerCase(Locale.ROOT),
+ ident);
+ return false;
+ }
+ for (int attempt = 1; ; attempt++) {
+ try {
+ return store.delete(ident, type, cascade, observed);
+ } catch (OptimisticLockException e) {
+ EntityVersion current = observeRegistration(ident, type);
+ if (current == null) {
+ LOG.warn(
+ "The {} registration of {} was removed concurrently while the
external drop ran",
+ type.name().toLowerCase(Locale.ROOT),
+ ident);
+ return false;
+ }
+ if (current.id() != observed.id()) {
+ LOG.warn(
+ "The {} registration of {} changed while the external drop ran
(expected {}, found"
+ + " {}); it is kept instead of being deleted under the new
incarnation",
+ type.name().toLowerCase(Locale.ROOT),
+ ident,
+ observed,
+ current);
+ return false;
+ }
+ if (attempt >= MAX_FENCED_DELETE_ATTEMPTS) {
+ throw e;
Review Comment:
[Important] This `throw` contradicts the contract written 40 lines above it,
and it re-opens the hazard that paragraph is there to prevent.
The javadoc states: "The external drop has already succeeded when this runs,
so a conflict is not reported to the caller: a failed request would invite a
retry, and by now the name may belong to a newly created object that the retry
would drop" (`OperationDispatcher.java:305-307`). Every other exit from the
loop honours that — a different id, a vanished row, and `observed == null` all
return `false`. Only the attempt-limit path throws, and
`OptimisticLockException` is mapped to a 409 by
`ExceptionHandlers.java:1177-1180` (`Utils.optimisticLockConflict`), which is
precisely the retryable status the comment argues against.
Failure scenario: `dropTable(t)` observes id X and the external drop
succeeds, so table `t` is gone from the catalog. Another writer keeps bumping
X's `current_version` (a tag or owner attachment, a statistic update — anything
that goes through `updateTable`), so the CAS in `deleteTableWithVersion` loses
three times. The client gets a 409. It retries `dropTable(t)`; meanwhile
someone creates a new `t`, and the retry drops that new external table. Note
the retry does not even buy the intended protection here: after three attempts
the registration is left behind anyway, so throwing gets the orphaned row *and*
the retryable error.
Suggest logging at WARN and returning `false` on exhaustion, matching the
other branches, and leaving the leftover row to the reconcile pass the javadoc
already defers to. Whichever way you go, the javadoc and the `@throws` tag at
line 323 should agree with the code.
Verified by: read `OperationDispatcher.java:302-370` at HEAD `bdc3ca0`;
confirmed the 409 mapping at
`server/src/main/java/org/apache/gravitino/server/web/rest/ExceptionHandlers.java:1177`;
grepped `core/src/main/java/org/apache/gravitino/catalog/` for a `catch
(OptimisticLockException` around the four `deleteObservedRegistration` call
sites (`TableOperationDispatcher.java:428` and `:477`,
`SchemaOperationDispatcher.java:597`, `TopicOperationDispatcher.java:258`,
`ViewOperationDispatcher.java:346`) — none catches it, so it reaches the REST
layer.
##########
core/src/main/java/org/apache/gravitino/catalog/SchemaOperationDispatcher.java:
##########
@@ -556,32 +560,41 @@ 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);
+ // The cascade removes every child registered under the observed
row, including one
+ // registered after the observation, which the identity fence on
the schema cannot
+ // tell apart. A child can only be created while the schema exists
in the catalog, so
+ // if it exists again the name was re-created meanwhile: keep its
registrations.
+ if (schemaRecreatedInCatalog(catalogIdent, ident)) {
Review Comment:
[Question] Is this probe safe on an asynchronous catalog? Elsewhere in these
dispatchers the code assumes it is not — `internalCreateTable` skips re-reading
the table because "some catalog APIs are asynchronous"
(`TableOperationDispatcher.java:666-667`), and `internalCreateTopic` repeats it
(`TopicOperationDispatcher.java:376-377`).
If `schemaExists` can still answer `true` shortly after a `dropSchema` that
returned `true`, this branch treats a perfectly ordinary successful drop as a
re-creation: the registration and all child registrations are kept,
`droppedFromStore` stays `false`, and the follow-up
`SchemaEntityCleaner.deleteOrphanedSchemaEntities` call at line 601 uses the
same predicate, so its first candidate is `ident` itself, `schemaExists` breaks
the walk, and it cleans nothing either. The request still returns `true`
because `droppedFromCatalog` is `true`, so the caller sees success while the
store keeps the whole stale subtree until some later drop happens to re-probe
it.
I could not find a catalog in-tree where I can show that happening, which is
why this is a question rather than a finding. If you have a catalog in mind
where the probe is reliable, it would help to say so in the comment; if not,
gating the skip on `droppedFromCatalog == false` (the case where the name
plausibly belongs to something else) would keep the protection for the
cascade-over-a-re-created-schema case without making a successful drop depend
on read-after-write.
Verified by: read `dropSchema` at `SchemaOperationDispatcher.java:525-614`
and `schemaRecreatedInCatalog` at `:694-703` at HEAD `bdc3ca0`; traced the
follow-on cleanup through `SchemaEntityCleaner.deleteOrphanedSchemaEntities`
(`SchemaEntityCleaner.java:58-98`), where `includeSelf = false` still starts
the walk at the dropped schema's own scope list; compared against the
async-catalog comments at `TableOperationDispatcher.java:667` and
`TopicOperationDispatcher.java:377`.
##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +256,174 @@ 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 expected id of the entity being updated
+ * @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
+ * @throws UnsupportedOperationException if the store cannot read versions
+ */
+ @Nullable
+ protected EntityVersion observeRegistration(NameIdentifier ident,
Entity.EntityType type) {
+ try {
+ return store.getVersion(ident, type);
+ } catch (NoSuchEntityException e) {
+ return null;
+ } catch (IOException e) {
+ throw new RuntimeException("Failed to read the registration of " +
ident, e);
+ }
+ }
+
+ /**
+ * Deletes the registration observed before an external-catalog call, and
only that one.
+ *
+ * <p>The external drop has already succeeded when this runs, so a conflict
is not reported to the
+ * caller: a failed request would invite a retry, and by now the name may
belong to a newly
+ * created object that the retry would drop. Instead:
+ *
+ * <ul>
+ * <li>A registration with another id belongs to a newer incarnation. It
is kept, and a
+ * reconcile pass, not this drop, decides what happens to it.
+ * <li>The same id with a newer version means a concurrent store write to
the entity that was
+ * just dropped. The fenced delete is retried, because giving up would
leave a registration
+ * whose external object is gone, and a retried drop would no longer
reach the store.
+ * </ul>
+ *
+ * @param ident the entity identifier
+ * @param type the entity type
+ * @param cascade whether to delete the children as well
+ * @param observed the id and version read before the external call, or null
when there was no
+ * registration to delete
+ * @return true if the observed registration was deleted
+ * @throws OptimisticLockException if the observed entity kept changing on
every attempt
+ * @throws UnsupportedOperationException if the store cannot delete with a
version check
+ */
+ protected boolean deleteObservedRegistration(
+ NameIdentifier ident,
+ Entity.EntityType type,
+ boolean cascade,
+ @Nullable EntityVersion observed) {
+ if (observed == null) {
+ LOG.warn(
+ "No {} registration was found for {} before the external drop;
leaving the store alone",
+ type.name().toLowerCase(Locale.ROOT),
+ ident);
+ return false;
+ }
+ for (int attempt = 1; ; attempt++) {
+ try {
+ return store.delete(ident, type, cascade, observed);
+ } catch (OptimisticLockException e) {
+ EntityVersion current = observeRegistration(ident, type);
+ if (current == null) {
+ LOG.warn(
+ "The {} registration of {} was removed concurrently while the
external drop ran",
+ type.name().toLowerCase(Locale.ROOT),
+ ident);
+ return false;
+ }
+ if (current.id() != observed.id()) {
+ LOG.warn(
+ "The {} registration of {} changed while the external drop ran
(expected {}, found"
+ + " {}); it is kept instead of being deleted under the new
incarnation",
+ type.name().toLowerCase(Locale.ROOT),
+ ident,
+ observed,
+ current);
+ return false;
+ }
+ if (attempt >= MAX_FENCED_DELETE_ATTEMPTS) {
+ throw e;
+ }
+ LOG.info(
+ "The {} registration of {} was updated concurrently; retrying the
delete (attempt {})",
+ type.name().toLowerCase(Locale.ROOT),
+ ident,
+ attempt + 1);
+ } catch (NoSuchEntityException e) {
+ LOG.warn("The {} to be dropped does not exist in the store: {}", type,
ident, e);
+ return false;
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ }
+ }
+
+ /**
+ * Stores the registration of an entity that the external catalog has just
created.
+ *
+ * <p>A successful external create means a registration observed before the
create with another id
+ * is stale: its object was dropped out of band, or by a drop on another
server that has not
+ * reached the store yet. A registration first observed after the create may
belong to a newer
+ * object, so it must be kept. An upsert by name would keep that row's id
and hand its owner,
+ * tags, and grants to the new object, and the pending drop's identity fence
would then pass and
+ * delete the new registration. The stale row is replaced instead: deleted
fenced on its own id,
+ * then the new registration is inserted.
+ *
+ * @param entity the registration of the newly created object
+ * @param cascade whether a stale registration is deleted with its children
+ * @param observed the registration read before the external create, or null
if none existed
+ * @param <E> the entity type
+ * @throws IOException if a store operation fails
+ * @throws OptimisticLockException if the stale registration changed while
it was replaced
+ */
+ protected <E extends Entity & HasIdentifier> void putCreatedEntity(
+ E entity, boolean cascade, @Nullable EntityVersion observed) throws
IOException {
+ try {
+ store.put(entity, false /* overwrite */);
+ return;
+ } catch (EntityAlreadyExistsException e) {
+ // A registration already sits under the name; decide below whether it
is stale.
+ }
+
+ NameIdentifier ident = entity.nameIdentifier();
+ EntityVersion existing = observeRegistration(ident, entity.type());
+ if (existing != null && existing.id() == entity.id()) {
+ // Another node already imported this object. Keep any updates it has
made since then.
+ return;
+ }
+ if (existing != null) {
+ if (observed == null || existing.id() != observed.id()) {
+ throw new OptimisticLockException(
+ "The registration of %s changed during create; keeping the newer
registration %s",
+ ident, existing);
+ }
+ LOG.warn(
+ "Replacing the stale {} registration {} of {} with the newly created
object {}",
+ entity.type().name().toLowerCase(Locale.ROOT),
+ existing,
+ ident,
+ entity.id());
+ store.delete(ident, entity.type(), cascade, observed);
+ }
+ store.put(entity, false /* overwrite */);
Review Comment:
[Nit] The stale-registration replacement is two independent store calls
where the old code was one atomic upsert, and for schemas the first call
cascades.
`createSchema` passes `cascade = true`
(`SchemaOperationDispatcher.java:230`), so line 422 removes the stale schema
row *and* every child registration under it before line 424 inserts the new
row. If the insert then fails — an `EntityAlreadyExistsException` from a racer
that took the name in between, or any `IOException` — the caller only logs and
returns a combined schema with no entity
(`SchemaOperationDispatcher.java:231-236`), leaving the just-created schema
with no registration at all. `loadSchema` will re-import it eventually, so this
self-heals for the schema row, and the cascaded children were stale by
construction, which is why this is a nit rather than a finding.
`EntityStore.executeInTransaction` (`EntityStore.java:400`) would make the
delete and the insert atomic and is the obvious fit here, though I see no other
dispatcher using it, so it may not be the convention you want in this layer.
Verified by: read `putCreatedEntity` at `OperationDispatcher.java:390-426`
at HEAD `bdc3ca0`, the four call sites (`SchemaOperationDispatcher.java:230`,
`TableOperationDispatcher.java:712`, `TopicOperationDispatcher.java:400`,
`ViewOperationDispatcher.java:433`) and the `catch (Exception e)` handler each
one wraps it in; confirmed `store.put(entity, false)` surfaces a duplicate name
as `EntityAlreadyExistsException` through `ExceptionUtils.checkSQLException` in
`TableMetaService.insertTable` (`TableMetaService.java:189-192`); grepped
`executeInTransaction` across
`core/src/main/java/org/apache/gravitino/catalog/` — no current users.
--
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]