Copilot commented on code in PR #11892:
URL: https://github.com/apache/gravitino/pull/11892#discussion_r3527662987


##########
catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java:
##########
@@ -547,13 +552,45 @@ private Table repairTableMetadata(NameIdentifier ident, 
Column[] columns, long d
     }
   }
 
+  /**
+   * Applies an idempotent update to the stored table, retrying when the 
optimistic-lock CAS is lost
+   * to a concurrent update. The repair-on-load path runs on every {@code 
loadTable}, so concurrent
+   * loads of the same table race on the version CAS; {@code store.update} 
surfaces the lost race as
+   * an {@link IOException}. Because the updater is idempotent, the loser 
re-reads the latest
+   * (already repaired) entity and retries instead of failing the whole load 
with a fatal error.
+   */
+  private TableEntity updateTableWithCasRetry(
+      NameIdentifier ident, Function<TableEntity, TableEntity> updater) throws 
IOException {
+    IOException lastConflict = null;
+    for (int attempt = 1; attempt <= REPAIR_UPDATE_MAX_ATTEMPTS; attempt++) {
+      try {
+        return store.update(ident, TableEntity.class, Entity.EntityType.TABLE, 
updater);
+      } catch (IOException e) {
+        // Only retry when the update matched 0 rows (lost optimistic-lock 
CAS). Other IO failures
+        // (DB outage, serialization errors, etc.) should fail fast.
+        String message = e.getMessage();
+        if (message == null || !message.startsWith("Failed to update the 
entity:")) {
+          throw e;
+        }

Review Comment:
   The retry decision is based on parsing an `IOException` message prefix, 
which is brittle (message text is not a stable API and may change, be 
localized, or differ by store implementation). Consider switching to a 
structured signal for optimistic-lock failure (e.g., a dedicated exception 
type, an error code, or a well-known checked exception from the 
`EntityStore`/meta layer), or at minimum centralize the match in a shared 
helper/constant owned by the layer that emits the exception.



##########
catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java:
##########
@@ -547,13 +552,45 @@ private Table repairTableMetadata(NameIdentifier ident, 
Column[] columns, long d
     }
   }
 
+  /**
+   * Applies an idempotent update to the stored table, retrying when the 
optimistic-lock CAS is lost
+   * to a concurrent update. The repair-on-load path runs on every {@code 
loadTable}, so concurrent
+   * loads of the same table race on the version CAS; {@code store.update} 
surfaces the lost race as
+   * an {@link IOException}. Because the updater is idempotent, the loser 
re-reads the latest
+   * (already repaired) entity and retries instead of failing the whole load 
with a fatal error.
+   */
+  private TableEntity updateTableWithCasRetry(
+      NameIdentifier ident, Function<TableEntity, TableEntity> updater) throws 
IOException {
+    IOException lastConflict = null;
+    for (int attempt = 1; attempt <= REPAIR_UPDATE_MAX_ATTEMPTS; attempt++) {
+      try {
+        return store.update(ident, TableEntity.class, Entity.EntityType.TABLE, 
updater);
+      } catch (IOException e) {
+        // Only retry when the update matched 0 rows (lost optimistic-lock 
CAS). Other IO failures
+        // (DB outage, serialization errors, etc.) should fail fast.
+        String message = e.getMessage();
+        if (message == null || !message.startsWith("Failed to update the 
entity:")) {
+          throw e;
+        }
+
+        lastConflict = e;
+        LOG.debug(
+            "Optimistic-lock conflict updating table {} metadata (attempt 
{}/{}), {}",
+            ident,
+            attempt,
+            REPAIR_UPDATE_MAX_ATTEMPTS,
+            attempt < REPAIR_UPDATE_MAX_ATTEMPTS ? "retrying" : "retries 
exhausted",
+            e);
+      }
+    }
+    throw lastConflict;

Review Comment:
   When retries are exhausted, the code rethrows the last conflict 
`IOException` as-is, which can hide the fact that multiple attempts were made 
and why the operation ultimately failed. Consider throwing a new `IOException` 
that includes context like attempts/maxAttempts/table ident (and set 
`lastConflict` as the cause) so operators can distinguish “persistent CAS 
contention” from a single conflict.



##########
catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceConcurrentRepairStress.java:
##########
@@ -0,0 +1,276 @@
+/*
+ *  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.catalog.lakehouse.lance;
+
+import static 
org.apache.gravitino.lance.common.utils.LanceConstants.LANCE_STORAGE_OPTIONS_PREFIX;
+import static 
org.apache.gravitino.lance.common.utils.LanceConstants.LANCE_TABLE_DECLARED;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.when;
+
+import java.io.IOException;
+import java.time.Instant;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.Callable;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Function;
+import org.apache.arrow.vector.types.pojo.ArrowType;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.apache.commons.lang3.tuple.Pair;
+import org.apache.gravitino.Config;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.Entity.EntityType;
+import org.apache.gravitino.EntityStore;
+import org.apache.gravitino.HasIdentifier;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.UserPrincipal;
+import org.apache.gravitino.catalog.ManagedSchemaOperations;
+import org.apache.gravitino.meta.AuditInfo;
+import org.apache.gravitino.meta.TableEntity;
+import org.apache.gravitino.rel.Table;
+import org.apache.gravitino.storage.IdGenerator;
+import org.apache.gravitino.utils.Executable;
+import org.apache.gravitino.utils.PrincipalUtils;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+import org.lance.Dataset;
+import org.mockito.Mockito;
+
+/**
+ * Real multi-threaded reproduction of the repair-on-load optimistic-lock race 
behind #11891. Unlike
+ * {@code 
TestLanceTableOperations#testLoadTableSurvivesConcurrentRepairVersionRace} (a
+ * deterministic mock that throws one scripted {@code IOException}), this test 
drives {@link
+ * LanceTableOperations#loadTable} from several threads at once against a 
{@link CasEntityStore}
+ * that models the production relational store's compare-and-set semantics 
faithfully: every {@code
+ * update} bumps a version guarded by the base version, so concurrent updates 
from the same base
+ * conflict and exactly one wins per generation — just like {@code 
TableMetaService.updateTable}
+ * ({@code UPDATE ... WHERE current_version = old}).
+ *
+ * <p>Before the CAS retry, the loser of the race got {@code 
IOException("Failed to update the
+ * entity")}, rethrown as a fatal {@code RuntimeException} (HTTP 500). This 
test asserts every
+ * concurrent load returns the repaired table instead.
+ */
+public class TestLanceConcurrentRepairStress {
+
+  @TempDir private java.nio.file.Path tempDir;
+
+  private static final NameIdentifier IDENT = NameIdentifier.of("schema", 
"table");
+
+  @Test
+  public void testConcurrentRepairLoadsSurviveCasContention() throws Exception 
{
+    // The reported race is "two loads repair the same table at once". Also 
drive a small herd to
+    // confirm the bounded CAS retry stays robust beyond the minimal two-load 
case.
+    runStress(2, 3000);
+    runStress(8, 2000);
+  }

Review Comment:
   This adds a high-iteration, multi-threaded stress test to the regular unit 
test suite (thousands of iterations across multiple threads). This is likely to 
increase CI runtime and may be flaky under constrained/variable build 
environments. Consider reducing iterations to the minimum that provides 
confidence, adding a JUnit `@Timeout`, and/or tagging it as a stress test (so 
it can be excluded from default runs and enabled in a dedicated stress/soak 
job).



-- 
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]

Reply via email to