This is an automated email from the ASF dual-hosted git repository.
roryqi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new d1480113b5 [#13452] fix(iceberg): discard builder changes when
filtering snapshots by refs (#13453)
d1480113b5 is described below
commit d1480113b5de78981d2b4e9391067fcef542fdd8
Author: Xinyi Lu <[email protected]>
AuthorDate: Wed Sep 23 18:46:36 2026 -0700
[#13452] fix(iceberg): discard builder changes when filtering snapshots by
refs (#13453)
### What changes were proposed in this pull request?
Call `discardChanges()` in
`IcebergTableOperations.filterSnapshotsByRefs` after
`suppressHistoricalSnapshots()` and before `build()`.
### Why are the changes needed?
`suppressHistoricalSnapshots()` also removes the statistics and
partition statistics
attached to the suppressed snapshots, and records those removals as
`MetadataUpdate`
changes on the builder. `TableMetadata.Builder.build()` refuses to set a
metadata
location while changes are pending, so `GET
.../tables/{table}?snapshots=refs` fails with
IllegalArgumentException: Cannot set metadata location with changes to
table metadata: N changes
for any table whose unreferenced snapshots carry statistics (e.g. tables
written by
Trino, which attaches statistics on every INSERT). This surfaced after
#13290 started
setting the metadata location explicitly. The filtered metadata is a
read-only view of
the table, not a commit, so the response must not carry pending updates.
Fix: #13452
### Does this PR introduce _any_ user-facing change?
no
(Please list the user-facing changes introduced by your change,
including
1. Change in user-facing APIs.
2. Addition or removal of property keys.)
### How was this patch tested?
UT
Co-authored-by: Jerry Shao <[email protected]>
---
.../service/rest/IcebergTableOperations.java | 1 +
.../service/rest/TestIcebergTableOperations.java | 148 +++++++++++++++++++++
2 files changed, 149 insertions(+)
diff --git
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/rest/IcebergTableOperations.java
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/rest/IcebergTableOperations.java
index a5d2998b13..552eb2726e 100644
---
a/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/rest/IcebergTableOperations.java
+++
b/iceberg/iceberg-rest-server/src/main/java/org/apache/gravitino/iceberg/service/rest/IcebergTableOperations.java
@@ -641,6 +641,7 @@ public class IcebergTableOperations {
TableMetadata.buildFrom(metadata)
.withMetadataLocation(metadata.metadataFileLocation())
.suppressHistoricalSnapshots()
+ .discardChanges()
.build();
LoadTableResponse.Builder builder =
LoadTableResponse.builder()
diff --git
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
index db0610e59f..d11ed1e1af 100644
---
a/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
+++
b/iceberg/iceberg-rest-server/src/test/java/org/apache/gravitino/iceberg/service/rest/TestIcebergTableOperations.java
@@ -20,6 +20,7 @@
package org.apache.gravitino.iceberg.service.rest;
import com.fasterxml.jackson.databind.JsonNode;
+import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.ImmutableSet;
import java.util.Arrays;
@@ -68,12 +69,16 @@ import
org.apache.gravitino.listener.api.event.IcebergUpdateTableFailureEvent;
import org.apache.gravitino.listener.api.event.IcebergUpdateTablePreEvent;
import org.apache.gravitino.server.ServerConfig;
import org.apache.gravitino.server.authorization.GravitinoAuthorizerProvider;
+import org.apache.iceberg.GenericBlobMetadata;
+import org.apache.iceberg.GenericStatisticsFile;
+import org.apache.iceberg.ImmutableGenericPartitionStatisticsFile;
import org.apache.iceberg.MetadataUpdate;
import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Schema;
import org.apache.iceberg.Snapshot;
import org.apache.iceberg.SnapshotParser;
import org.apache.iceberg.SnapshotRef;
+import org.apache.iceberg.StatisticsFile;
import org.apache.iceberg.TableMetadata;
import org.apache.iceberg.TableMetadataParser;
import org.apache.iceberg.UpdateRequirement;
@@ -664,6 +669,14 @@ public class TestIcebergTableOperations extends
IcebergNamespaceTestBase {
return loadTableResponse.tableMetadata();
}
+ private Response doUpdateTable(
+ Namespace ns, String name, TableMetadata base, List<MetadataUpdate>
metadataUpdates) {
+ List<UpdateRequirement> requirements =
UpdateRequirements.forUpdateTable(base, metadataUpdates);
+ UpdateTableRequest updateTableRequest = new
UpdateTableRequest(requirements, metadataUpdates);
+ return getTableClientBuilder(ns, Optional.of(name))
+ .post(Entity.entity(updateTableRequest,
MediaType.APPLICATION_JSON_TYPE));
+ }
+
private void verifyUpdateTableFail(Namespace ns, String name, int status,
TableMetadata base) {
Response response = doUpdateTable(ns, name, base);
Assertions.assertEquals(status, response.getStatus());
@@ -1258,6 +1271,71 @@ public class TestIcebergTableOperations extends
IcebergNamespaceTestBase {
Assertions.assertEquals("org.apache.iceberg.aws.s3.S3FileIO",
filtered.config().get("io-impl"));
}
+ @ParameterizedTest
+
@MethodSource("org.apache.gravitino.iceberg.service.rest.IcebergRestTestUtil#testNamespaces")
+ void testLoadTableSnapshotsRefsWithStatisticsOnHistoricalSnapshot(Namespace
namespace) {
+ verifyCreateNamespaceSucc(namespace);
+ String tableName = "snapshots_refs_stats_foo1";
+ verifyCreateTableSucc(namespace, tableName, true);
+
+ LoadTableResponse before =
+ doLoadTableWithSnapshots(namespace, tableName,
"all").readEntity(LoadTableResponse.class);
+ TableMetadata base = before.tableMetadata();
+ Set<Long> referencedSnapshotIds =
+
base.refs().values().stream().map(SnapshotRef::snapshotId).collect(Collectors.toSet());
+ long historicalSnapshotId =
+ base.snapshots().stream()
+ .map(Snapshot::snapshotId)
+ .filter(id -> !referencedSnapshotIds.contains(id))
+ .findFirst()
+ .orElseThrow(() -> new AssertionError("expected an unreferenced
snapshot"));
+ long currentSnapshotId = base.currentSnapshot().snapshotId();
+
+ // Attach statistics and partition statistics to the historical snapshot,
and statistics to the
+ // current one so we can check that only the suppressed snapshot's files
are dropped.
+ Response updateResponse =
+ doUpdateTable(
+ namespace,
+ tableName,
+ base,
+ ImmutableList.of(
+ new
MetadataUpdate.SetStatistics(statisticsFile(historicalSnapshotId)),
+ new MetadataUpdate.SetPartitionStatistics(
+ ImmutableGenericPartitionStatisticsFile.builder()
+ .snapshotId(historicalSnapshotId)
+
.path("s3://bucket/db/tbl/metadata/partition-stats-historical.parquet")
+ .fileSizeInBytes(10L)
+ .build()),
+ new
MetadataUpdate.SetStatistics(statisticsFile(currentSnapshotId))));
+ Assertions.assertEquals(Status.OK.getStatusCode(),
updateResponse.getStatus());
+
+ Response allResponse = doLoadTableWithSnapshots(namespace, tableName,
"all");
+ Assertions.assertEquals(Status.OK.getStatusCode(),
allResponse.getStatus());
+ TableMetadata all =
allResponse.readEntity(LoadTableResponse.class).tableMetadata();
+ Assertions.assertEquals(2, all.statisticsFiles().size());
+ Assertions.assertEquals(1, all.partitionStatisticsFiles().size());
+
+ Response refsResponse = doLoadTableWithSnapshots(namespace, tableName,
"refs");
+ Assertions.assertEquals(
+ Status.OK.getStatusCode(),
+ refsResponse.getStatus(),
+ "snapshots=refs must not fail when a suppressed snapshot has
statistics");
+ TableMetadata refs =
refsResponse.readEntity(LoadTableResponse.class).tableMetadata();
+
+ Assertions.assertEquals(
+ referencedSnapshotIds,
+
refs.snapshots().stream().map(Snapshot::snapshotId).collect(Collectors.toSet()));
+ Assertions.assertEquals(
+ ImmutableSet.of(currentSnapshotId),
+
refs.statisticsFiles().stream().map(StatisticsFile::snapshotId).collect(Collectors.toSet()),
+ "statistics of suppressed snapshots are dropped, statistics of kept
snapshots remain");
+ Assertions.assertTrue(
+ refs.partitionStatisticsFiles().isEmpty(),
+ "partition statistics of the suppressed snapshot must be dropped");
+ Assertions.assertEquals(all.metadataFileLocation(),
refs.metadataFileLocation());
+ Assertions.assertEquals(all.lastUpdatedMillis(), refs.lastUpdatedMillis());
+ }
+
@Test
void testFilterSnapshotsByRefsPreservesMetadataLocationAndHistory() {
TableMetadata base =
@@ -1302,6 +1380,76 @@ public class TestIcebergTableOperations extends
IcebergNamespaceTestBase {
"snapshot-log must be kept intact for lazy snapshot loading");
}
+ @Test
+ void testFilterSnapshotsByRefsDiscardsStatisticsRemovalChanges() {
+ TableMetadata base =
+ TableMetadata.newTableMetadata(
+ tableSchema, PartitionSpec.unpartitioned(), "s3://bucket/db/tbl",
ImmutableMap.of());
+ Snapshot first = snapshot(1L, null, 1000L);
+ Snapshot second = snapshot(2L, 1L, 2000L);
+ TableMetadata withHistory =
+ TableMetadata.buildFrom(
+ TableMetadata.buildFrom(base).setBranchSnapshot(first,
"main").build())
+ .setBranchSnapshot(second, "main")
+ // statistics on the historical snapshot are what
suppressHistoricalSnapshots() removes
+ .setStatistics(statisticsFile(1L))
+ .setPartitionStatistics(
+ ImmutableGenericPartitionStatisticsFile.builder()
+ .snapshotId(1L)
+
.path("s3://bucket/db/tbl/metadata/partition-stats-1.parquet")
+ .fileSizeInBytes(10L)
+ .build())
+ .setStatistics(statisticsFile(2L))
+ .build();
+ String metadataLocation =
"s3://bucket/db/tbl/metadata/00003-abc.metadata.json";
+ TableMetadata metadata =
+ TableMetadataParser.fromJson(metadataLocation,
TableMetadataParser.toJson(withHistory));
+ Assertions.assertEquals(2, metadata.statisticsFiles().size());
+ Assertions.assertEquals(1, metadata.partitionStatisticsFiles().size());
+ LoadTableResponse original =
LoadTableResponse.builder().withTableMetadata(metadata).build();
+
+ // Before the fix this threw IllegalArgumentException:
+ // "Cannot set metadata location with changes to table metadata: 2 changes"
+ LoadTableResponse filtered =
+ Assertions.assertDoesNotThrow(() ->
IcebergTableOperations.filterSnapshotsByRefs(original));
+
+ Assertions.assertEquals(
+ ImmutableSet.of(2L),
+ filtered.tableMetadata().snapshots().stream()
+ .map(Snapshot::snapshotId)
+ .collect(Collectors.toSet()));
+ Assertions.assertEquals(
+ ImmutableSet.of(2L),
+ filtered.tableMetadata().statisticsFiles().stream()
+ .map(StatisticsFile::snapshotId)
+ .collect(Collectors.toSet()),
+ "only the suppressed snapshot's statistics are dropped");
+
Assertions.assertTrue(filtered.tableMetadata().partitionStatisticsFiles().isEmpty());
+ Assertions.assertTrue(
+ filtered.tableMetadata().changes().isEmpty(),
+ "a load response must not carry pending metadata updates");
+ Assertions.assertEquals(metadataLocation, filtered.metadataLocation());
+ Assertions.assertEquals(
+ metadata.lastUpdatedMillis(),
filtered.tableMetadata().lastUpdatedMillis());
+ Assertions.assertEquals(metadata.previousFiles(),
filtered.tableMetadata().previousFiles());
+ Assertions.assertEquals(metadata.snapshotLog(),
filtered.tableMetadata().snapshotLog());
+ }
+
+ private static StatisticsFile statisticsFile(long snapshotId) {
+ return new GenericStatisticsFile(
+ snapshotId,
+ String.format("s3://bucket/db/tbl/metadata/stats-%d.puffin",
snapshotId),
+ 100L,
+ 42L,
+ ImmutableList.of(
+ new GenericBlobMetadata(
+ "apache-datasketches-theta-v1",
+ snapshotId,
+ 1L,
+ ImmutableList.of(1),
+ ImmutableMap.of())));
+ }
+
private static Snapshot snapshot(long snapshotId, Long parentId, long
timestampMs) {
String json =
String.format(