yuqi1129 commented on code in PR #11892:
URL: https://github.com/apache/gravitino/pull/11892#discussion_r3519506331
##########
catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceTableOperations.java:
##########
@@ -170,6 +171,91 @@ public void testLoadDeclaredTableSchemaFromLocation()
throws Exception {
Assertions.assertFalse(loadedTable.properties().containsKey(LANCE_TABLE_DECLARED));
}
+ /**
+ * Reproduces the concurrent repair-on-load race seen in {@code
LanceSparkRESTServiceIT}. When two
+ * loads repair the same table at once, the optimistic-locked {@code
store.update} of the slower
+ * one matches zero rows and {@code TableMetaService} surfaces it as {@code
IOException("Failed to
+ * update the entity")}. {@code repairTableMetadata} currently rethrows it
as a fatal {@code
+ * RuntimeException} (HTTP 500) instead of tolerating the concurrent update.
A correct fix should
+ * treat the lost race as benign and return a usable table, so this test
asserts the desired
+ * behavior and currently fails against the buggy code.
+ */
+ @Test
+ public void testLoadTableSurvivesConcurrentRepairVersionRace() throws
Exception {
+ NameIdentifier ident = NameIdentifier.of("schema", "table");
+ String location = tempDir.resolve("concurrent-repair-table").toString();
+ TableEntity tableEntity =
+ tableEntity(
+ ident,
+ List.of(),
+ Map.of(
+ Table.PROPERTY_LOCATION,
+ location,
+ LANCE_TABLE_DECLARED,
+ "true",
+ LANCE_STORAGE_OPTIONS_PREFIX + "endpoint",
+ "http://endpoint"));
+ // The winner of the race already repaired the table to the dataset schema
and version.
+ TableEntity alreadyRepairedTableEntity =
+ tableEntity(
+ ident,
+ List.of(
+ ColumnEntity.builder()
+ .withId(11L)
+ .withName("id")
+ .withDataType(Types.IntegerType.get())
+ .withPosition(0)
+ .withAuditInfo(AuditInfo.EMPTY)
+ .build(),
+ ColumnEntity.builder()
+ .withId(12L)
+ .withName("name")
+ .withDataType(Types.StringType.get())
+ .withPosition(1)
+ .withAuditInfo(AuditInfo.EMPTY)
+ .build()),
+ Map.of(Table.PROPERTY_LOCATION, location, LANCE_TABLE_VERSION,
"8"));
+ when(store.get(eq(ident), eq(Entity.EntityType.TABLE),
eq(TableEntity.class)))
+ .thenReturn(tableEntity);
+ when(idGenerator.nextId()).thenReturn(10L, 11L);
+
+ // First repair attempt loses the optimistic-lock CAS (a concurrent load
already bumped the
+ // version): TableMetaService surfaces exactly this IOException. The retry
re-reads the winner's
+ // already-repaired entity, against which the idempotent updater succeeds.
+ when(store.update(eq(ident), eq(TableEntity.class),
eq(Entity.EntityType.TABLE), any()))
+ .thenThrow(new IOException("Failed to update the entity: " + ident))
+ .thenAnswer(
+ invocation -> {
+ @SuppressWarnings("unchecked")
+ Function<TableEntity, TableEntity> updater =
invocation.getArgument(3);
+ return updater.apply(alreadyRepairedTableEntity);
+ });
+
+ Dataset dataset = mock(Dataset.class);
+ when(dataset.getSchema())
+ .thenReturn(
+ new Schema(
+ List.of(
+ Field.nullable("id", new ArrowType.Int(32, true)),
+ Field.nullable("name", new ArrowType.Utf8()))));
+ when(dataset.version()).thenReturn(8L);
+ Mockito.doReturn(dataset)
+ .when(lanceTableOps)
+ .openDataset(location, Map.of("endpoint", "http://endpoint"));
+
+ // A lost repair race must not fail the load: the bounded CAS retry
recovers and returns the
+ // repaired table. Fails today because repairTableMetadata rethrows the
first conflict as
+ // RuntimeException("Failed to repair table").
Review Comment:
Fixed in d79c7b2963: rewrote the inline comments to describe the intended
guarantee instead of present-tense failure.
##########
catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceTableOperations.java:
##########
@@ -170,6 +171,91 @@ public void testLoadDeclaredTableSchemaFromLocation()
throws Exception {
Assertions.assertFalse(loadedTable.properties().containsKey(LANCE_TABLE_DECLARED));
}
+ /**
+ * Reproduces the concurrent repair-on-load race seen in {@code
LanceSparkRESTServiceIT}. When two
+ * loads repair the same table at once, the optimistic-locked {@code
store.update} of the slower
+ * one matches zero rows and {@code TableMetaService} surfaces it as {@code
IOException("Failed to
+ * update the entity")}. {@code repairTableMetadata} currently rethrows it
as a fatal {@code
+ * RuntimeException} (HTTP 500) instead of tolerating the concurrent update.
A correct fix should
+ * treat the lost race as benign and return a usable table, so this test
asserts the desired
+ * behavior and currently fails against the buggy code.
+ */
Review Comment:
Fixed in d79c7b2963: rewrote the test Javadoc to describe the pre-fix
behavior instead of saying the test currently fails.
##########
catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java:
##########
@@ -547,13 +552,37 @@ 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) {
+ lastConflict = e;
+ LOG.debug(
+ "Optimistic-lock conflict updating table {} metadata (attempt
{}/{}), retrying",
+ ident,
+ attempt,
+ REPAIR_UPDATE_MAX_ATTEMPTS,
+ e);
+ }
Review Comment:
Fixed in d79c7b2963: the log now says `retrying` only when another attempt
remains and `retries exhausted` on the final attempt.
--
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]