yuqi1129 commented on code in PR #13198:
URL: https://github.com/apache/gravitino/pull/13198#discussion_r4123728371
##########
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:
I kept the two calls to keep this fix small. RelationalEntityStore does not
support executeInTransaction yet. In d220e5b436, I documented the gap and added
a test showing that a later load restores the schema row if the insert fails.
Making the two calls atomic needs a separate change.
--
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]