This is an automated email from the ASF dual-hosted git repository.
CTTY pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new f2a76bb363 Core: Fix time-travel snapshot lookup to not assume
snapshot-log ordering (#17358) (#17360)
f2a76bb363 is described below
commit f2a76bb363d35d8957c75c8211aa070f174a4778
Author: Xiening Dai <[email protected]>
AuthorDate: Thu Jul 30 08:22:14 2026 -0700
Core: Fix time-travel snapshot lookup to not assume snapshot-log ordering
(#17358) (#17360)
nullableSnapshotIdAsOfTime previously relied on snapshot-log entries being
in chronological order, returning the last entry with timestamp <=
requested.
The spec does not guarantee ordering, and TableMetadata allows up to 1
minute
of clock skew between entries. This could return the wrong snapshot when
entries are out of order (e.g. due to WAP cherry-pick workflows).
Fix by finding the entry with the maximum timestamp <= requested timestamp,
which is correct regardless of iteration order.
Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
---
.../java/org/apache/iceberg/util/SnapshotUtil.java | 5 +-
.../org/apache/iceberg/util/TestSnapshotUtil.java | 67 ++++++++++++++++++++++
2 files changed, 71 insertions(+), 1 deletion(-)
diff --git a/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java
b/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java
index 370bbfed33..1067a119cc 100644
--- a/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java
+++ b/core/src/main/java/org/apache/iceberg/util/SnapshotUtil.java
@@ -390,9 +390,12 @@ public class SnapshotUtil {
public static Long nullableSnapshotIdAsOfTime(Table table, long
timestampMillis) {
Long snapshotId = null;
+ long bestTimestamp = Long.MIN_VALUE;
for (HistoryEntry logEntry : table.history()) {
- if (logEntry.timestampMillis() <= timestampMillis) {
+ if (logEntry.timestampMillis() <= timestampMillis
+ && logEntry.timestampMillis() > bestTimestamp) {
snapshotId = logEntry.snapshotId();
+ bestTimestamp = logEntry.timestampMillis();
}
}
diff --git a/core/src/test/java/org/apache/iceberg/util/TestSnapshotUtil.java
b/core/src/test/java/org/apache/iceberg/util/TestSnapshotUtil.java
index 96b56e5ffb..12a6936e73 100644
--- a/core/src/test/java/org/apache/iceberg/util/TestSnapshotUtil.java
+++ b/core/src/test/java/org/apache/iceberg/util/TestSnapshotUtil.java
@@ -22,13 +22,17 @@ import static
org.apache.iceberg.types.Types.NestedField.optional;
import static org.apache.iceberg.types.Types.NestedField.required;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
import java.io.File;
+import java.util.Arrays;
import java.util.Iterator;
import java.util.List;
import java.util.stream.StreamSupport;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.DataFiles;
+import org.apache.iceberg.HistoryEntry;
import org.apache.iceberg.MetadataTableType;
import org.apache.iceberg.MetadataTableUtils;
import org.apache.iceberg.PartitionSpec;
@@ -317,4 +321,67 @@ public class TestSnapshotUtil {
assertThat(SnapshotUtil.schemaFor(snapshotsTable,
firstSnapshotId).asStruct())
.isEqualTo(snapshotsTable.schema().asStruct());
}
+
+ @Test
+ public void snapshotIdAsOfTimeWithOutOfOrderHistory() {
+ long snapshotA = 1L;
+ long snapshotB = 2L;
+ long snapshotC = 3L;
+
+ long timeA = 1000L;
+ long timeB = 1010L;
+ long timeC = 1020L;
+
+ // history entries are out of chronological order: B appears before A
+ List<HistoryEntry> outOfOrderHistory =
+ Arrays.asList(
+ historyEntry(timeB, snapshotB),
+ historyEntry(timeA, snapshotA),
+ historyEntry(timeC, snapshotC));
+
+ Table mockTable = mock(Table.class);
+ when(mockTable.history()).thenReturn(outOfOrderHistory);
+
+ // querying at timeA should return snapshotA (max timestamp <= timeA)
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable,
timeA)).isEqualTo(snapshotA);
+
+ // querying at a point between A and B but closer to B should still return
snapshotA
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable, timeB -
1)).isEqualTo(snapshotA);
+
+ // querying at timeB should return snapshotB (max timestamp <= timeB)
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable,
timeB)).isEqualTo(snapshotB);
+
+ // querying at timeC should return snapshotC
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable,
timeC)).isEqualTo(snapshotC);
+
+ // querying between A and B should return snapshotA
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable, timeA +
5)).isEqualTo(snapshotA);
+
+ // querying before any entry should return null
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable,
999L)).isNull();
+ }
+
+ @Test
+ public void snapshotIdAsOfTimeWithDuplicateTimestamps() {
+ long snapshotA = 1L;
+ long snapshotB = 2L;
+
+ long sameTime = 1000L;
+
+ // two entries with the same timestamp — the first one encountered is
returned
+ List<HistoryEntry> history =
+ Arrays.asList(historyEntry(sameTime, snapshotA),
historyEntry(sameTime, snapshotB));
+
+ Table mockTable = mock(Table.class);
+ when(mockTable.history()).thenReturn(history);
+
+ assertThat(SnapshotUtil.nullableSnapshotIdAsOfTime(mockTable,
sameTime)).isEqualTo(snapshotA);
+ }
+
+ private static HistoryEntry historyEntry(long timestampMillis, long
snapshotId) {
+ HistoryEntry entry = mock(HistoryEntry.class);
+ when(entry.timestampMillis()).thenReturn(timestampMillis);
+ when(entry.snapshotId()).thenReturn(snapshotId);
+ return entry;
+ }
}