yuqi1129 commented on code in PR #13198:
URL: https://github.com/apache/gravitino/pull/13198#discussion_r4123726186
##########
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:
Fixed in d220e5b436. Iceberg table and view drops now read the stored ID
before the backend drop. The store delete checks that ID, so it keeps a new
object created under the same name. Added tests for both 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:
Fixed in d220e5b436. After three delete conflicts, we now log a warning and
keep the stored row instead of throwing an error. The external drop still
returns success. Updated the Javadoc and added a test for all three attempts
failing.
##########
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:
Fixed in d220e5b436. All four import paths now use putCreatedEntity instead
of overwrite. They read the stored ID before loading the external object and
only replace that old row. A newer row is kept. Added a test for import running
during a drop.
##########
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:
Agreed, true may mean the drop is not visible yet. In d220e5b436, I kept the
check and made the comment clear. We also keep the schema and its children if
the check fails. I did not limit the check to failed drops, since a schema can
be created again after a successful drop too. Added tests for both cases.
--
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]