This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 2aec1315ab [Cherry-pick to branch-1.3] [#13452] fix(iceberg): discard
builder changes when filtering snapshots by refs (#13453) (#13494)
2aec1315ab is described below
commit 2aec1315ab4a8056d778cce532ccd285285eb151
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Sep 24 12:06:40 2026 +0800
[Cherry-pick to branch-1.3] [#13452] fix(iceberg): discard builder changes
when filtering snapshots by refs (#13453) (#13494)
**Cherry-pick Information:**
- Original commit: d1480113b5de78981d2b4e9391067fcef542fdd8
- Target branch: `branch-1.3`
- Status: ✅ Clean cherry-pick (no conflicts)
Co-authored-by: Xinyi Lu <[email protected]>
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 5355bd0439..9f04543fa8 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
@@ -569,6 +569,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 5a0c395206..5f1618099b 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;
@@ -66,12 +67,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;
@@ -585,6 +590,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());
@@ -1179,6 +1192,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 =
@@ -1223,6 +1301,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(