From efdfa04f08e27e4275d581292764202f621c8935 Mon Sep 17 00:00:00 2001 From: Ayush Saxena Date: Fri, 24 Jul 2026 12:17:10 +0530 Subject: [PATCH] Fix concurrent table commits failing with a fatal 400 instead of a retryable 409 --- CHANGELOG.md | 1 + .../catalog/iceberg/CatalogHandlerUtils.java | 13 ++++- .../iceberg/IcebergCatalogHandler.java | 9 +++- .../AbstractLocalIcebergCatalogTest.java | 39 +++++++++++++++ ...ntTest.java => CommitTransactionTest.java} | 50 ++++++++++++++++++- 5 files changed, 109 insertions(+), 3 deletions(-) rename runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/{CommitTransactionEventTest.java => CommitTransactionTest.java} (85%) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7344e63eff2..63b17e1ffce 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -64,6 +64,7 @@ request adding CHANGELOG notes for breaking (!) changes and possibly other secti - Polaris does not include the `expiration-time` property anymore when vending credentials. This property is not consumed by any known client and duplicates the properties specific to each storage provider, such as `s3.session-token-expires-at-ms` for S3. +- Concurrent table commits that hit a stale sequence number now return a retryable `409` instead of a fatal `400`, for both single-table commits and `commitTransaction`. ### New Features - Added Kafka PolarisEventListener for publishing events to Kafka. diff --git a/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/CatalogHandlerUtils.java b/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/CatalogHandlerUtils.java index 8aedda19255..8c00688b987 100644 --- a/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/CatalogHandlerUtils.java +++ b/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/CatalogHandlerUtils.java @@ -46,6 +46,7 @@ import org.apache.iceberg.DataOperations; import org.apache.iceberg.MetadataUpdate; import org.apache.iceberg.MetadataUpdate.UpgradeFormatVersion; +import org.apache.iceberg.RetryableValidationException; import org.apache.iceberg.Schema; import org.apache.iceberg.Snapshot; import org.apache.iceberg.SnapshotRef; @@ -513,7 +514,17 @@ public TableMetadata commit(TableOperations ops, UpdateTableRequest request) { } TableMetadata.Builder newMetadataBuilder = TableMetadata.buildFrom(newBase); - request.updates().forEach((update) -> update.applyTo(newMetadataBuilder)); + try { + request.updates().forEach((update) -> update.applyTo(newMetadataBuilder)); + } catch (RetryableValidationException e) { + // Validation failed because the commit includes stale values (e.g. sequence + // number or first-row-id behind the current table state). This is not a conflict. + // Server-side retry won't help since the stale values are in the request itself. + // Wrap as CommitFailedException so the client can retry with refreshed metadata. + throw new ValidationFailureException( + new CommitFailedException( + e, "Validation failed, please retry: %s", e.getMessage())); + } TableMetadata updated = newMetadataBuilder.build(); if (updated.changes().isEmpty()) { // do not commit if the metadata has not changed diff --git a/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/IcebergCatalogHandler.java b/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/IcebergCatalogHandler.java index ddde1341117..a4a3402751e 100644 --- a/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/IcebergCatalogHandler.java +++ b/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/IcebergCatalogHandler.java @@ -50,6 +50,7 @@ import org.apache.iceberg.CatalogUtil; import org.apache.iceberg.MetadataUpdate; import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.RetryableValidationException; import org.apache.iceberg.SortOrder; import org.apache.iceberg.Table; import org.apache.iceberg.TableMetadata; @@ -1531,7 +1532,13 @@ public void commitTransaction(CommitTransactionRequest commitTransactionRequest) } // Apply updates to builder - singleUpdate.applyTo(metadataBuilder); + try { + singleUpdate.applyTo(metadataBuilder); + } catch (RetryableValidationException e) { + // Surface as a retryable 409, matching CatalogHandlerUtils.commit. + throw new CommitFailedException( + e, "Validation failed, please retry: %s", e.getMessage()); + } } // Update currentMetadata to reflect this change for subsequent requirement validation diff --git a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java index 14a57593968..8da58c833c6 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java @@ -20,6 +20,7 @@ import static java.nio.charset.StandardCharsets.UTF_8; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Fail.fail; import static org.awaitility.Awaitility.await; import static org.mockito.ArgumentMatchers.any; @@ -31,6 +32,7 @@ import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.when; @@ -73,6 +75,7 @@ import org.apache.iceberg.MetadataUpdate; import org.apache.iceberg.NullOrder; import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.RetryableValidationException; import org.apache.iceberg.RowDelta; import org.apache.iceberg.Schema; import org.apache.iceberg.Snapshot; @@ -95,6 +98,7 @@ import org.apache.iceberg.exceptions.NoSuchNamespaceException; import org.apache.iceberg.exceptions.NotFoundException; import org.apache.iceberg.exceptions.ServiceFailureException; +import org.apache.iceberg.exceptions.ValidationException; import org.apache.iceberg.inmemory.InMemoryFileIO; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.FileIO; @@ -676,6 +680,41 @@ public void testConcurrentWritesWithRollbackNonEmptyTable() { } } + @Test + public void commitSurfacesRetryableValidationFailureAsRetryableCommitConflict() { + Schema schema = new Schema(Types.NestedField.required(1, "id", Types.LongType.get())); + TableMetadata base = + TableMetadata.newTableMetadata( + schema, PartitionSpec.unpartitioned(), "file:///tmp/t", Map.of()); + TableOperations ops = mock(TableOperations.class); + when(ops.current()).thenReturn(base); + + // Stands in for the RetryableValidationException addSnapshot raises under a concurrent commit. + MetadataUpdate retryableUpdate = + new MetadataUpdate() { + @Override + public void applyTo(TableMetadata.Builder metadataBuilder) { + throw new RetryableValidationException( + "Cannot add snapshot with sequence number 6 older than last sequence number 6"); + } + + @Override + public void applyTo(ViewMetadata.Builder viewMetadataBuilder) { + throw new UnsupportedOperationException(); + } + }; + + UpdateTableRequest request = new UpdateTableRequest(List.of(), List.of(retryableUpdate)); + CatalogHandlerUtils catalogHandlerUtils = new CatalogHandlerUtils(5, false); + + assertThatThrownBy(() -> catalogHandlerUtils.commit(ops, request)) + .isInstanceOf(CommitFailedException.class) + // RetryableValidationException must not leak: it is a ValidationException, which maps to + // 400. + .isNotInstanceOf(ValidationException.class) + .hasMessageContaining("Validation failed, please retry"); + } + @Test public void testConcurrentWritesWithRollbackWithNonReplaceSnapshotInBetween() { LocalIcebergCatalog catalog = this.catalog(); diff --git a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/CommitTransactionEventTest.java b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/CommitTransactionTest.java similarity index 85% rename from runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/CommitTransactionEventTest.java rename to runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/CommitTransactionTest.java index e996e16cea7..24bb97389c9 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/CommitTransactionEventTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/CommitTransactionTest.java @@ -29,14 +29,18 @@ import java.util.Map; import java.util.UUID; import org.apache.iceberg.MetadataUpdate; +import org.apache.iceberg.RetryableValidationException; import org.apache.iceberg.TableMetadata; import org.apache.iceberg.UpdateRequirement; import org.apache.iceberg.catalog.Namespace; import org.apache.iceberg.catalog.TableIdentifier; +import org.apache.iceberg.exceptions.CommitFailedException; +import org.apache.iceberg.exceptions.ValidationException; import org.apache.iceberg.rest.requests.CommitTransactionRequest; import org.apache.iceberg.rest.requests.CreateNamespaceRequest; import org.apache.iceberg.rest.requests.CreateTableRequest; import org.apache.iceberg.rest.requests.UpdateTableRequest; +import org.apache.iceberg.view.ViewMetadata; import org.apache.polaris.core.admin.model.Catalog; import org.apache.polaris.core.admin.model.CatalogProperties; import org.apache.polaris.core.admin.model.CreateCatalogRequest; @@ -51,7 +55,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -public class CommitTransactionEventTest { +public class CommitTransactionTest { private static final String namespace = "ns"; private static final String catalog = "test-catalog"; private static final String propertyName = "custom-property-1"; @@ -127,6 +131,50 @@ void testEventsForUnSuccessfulTransaction() { .isInstanceOf(IllegalStateException.class); } + @Test + void commitTransactionSurfacesRetryableValidationFailureAsRetryableCommitConflict() { + TestServices testServices = createTestServices(); + createCatalogAndNamespace(testServices, Map.of(), catalogLocation); + String tableName = "retryable-conflict-table"; + createTable(testServices, tableName, catalogLocation); + + // Stands in for the RetryableValidationException addSnapshot raises under a concurrent commit. + MetadataUpdate retryableUpdate = + new MetadataUpdate() { + @Override + public void applyTo(TableMetadata.Builder metadataBuilder) { + throw new RetryableValidationException( + "Cannot add snapshot with sequence number 6 older than last sequence number 6"); + } + + @Override + public void applyTo(ViewMetadata.Builder viewMetadataBuilder) { + throw new UnsupportedOperationException(); + } + }; + CommitTransactionRequest request = + new CommitTransactionRequest( + List.of( + UpdateTableRequest.create( + TableIdentifier.of(namespace, tableName), + List.of(), + List.of(retryableUpdate)))); + + assertThatThrownBy( + () -> + testServices + .restApi() + .commitTransaction( + catalog, + request, + IDEMPOTENCY_KEY, + testServices.realmContext(), + testServices.securityContext())) + .isInstanceOf(CommitFailedException.class) + .isNotInstanceOf(ValidationException.class) + .hasMessageContaining("Validation failed, please retry"); + } + @Test void testLoadTableResponsesInCommitTransaction() { TestServices testServices = createTestServices();