laskoviymishka commented on code in PR #3260:
URL: https://github.com/apache/iceberg-rust/pull/3260#discussion_r4102928746
##########
crates/iceberg/src/util/snapshot.rs:
##########
@@ -76,10 +77,38 @@ pub fn ancestors_between(
})
}
+/// Resolve the snapshot ID from the latest main-history entry at or before
+/// `timestamp_ms` (milliseconds since the Unix epoch).
+///
+/// Equal timestamps select the first entry. Returns [`ErrorKind::DataInvalid`]
+/// if no matching history exists. The returned snapshot may have expired, so
+/// [`snapshot_by_id`](crate::spec::TableMetadata::snapshot_by_id) can still
return `None`.
+pub fn snapshot_id_as_of_time(table_metadata: &TableMetadataRef, timestamp_ms:
i64) -> Result<i64> {
+ let best = table_metadata
+ .history()
+ .iter()
+ .filter(|entry| entry.timestamp_ms() <= timestamp_ms)
Review Comment:
This is the tie-break restructure from last round — still leaning on the
comment rather than the structure. A one-char `>`→`>=` here silently flips
first-wins to last-wins and nothing but CI would catch it. I'd encode
first-wins in the sort key so misreading the operator can't change semantics:
```rust
table_metadata
.history()
.iter()
.enumerate()
.filter(|(_, entry)| entry.timestamp_ms() <= timestamp_ms)
.max_by_key(|(idx, entry)| (entry.timestamp_ms(), Reverse(*idx)))
.map(|(_, entry)| entry.snapshot_id)
```
`max_by_key` on its own picks the *last* max on ties, so `Reverse(idx)` is
doing the real work; the result is still an `Option`, so it feeds the existing
`.ok_or_else(...)` unchanged.
##########
crates/iceberg/src/util/snapshot.rs:
##########
@@ -76,10 +77,38 @@ pub fn ancestors_between(
})
}
+/// Resolve the snapshot ID from the latest main-history entry at or before
+/// `timestamp_ms` (milliseconds since the Unix epoch).
+///
Review Comment:
Still missing the cross-engine note from last round: ties here are
first-wins (matching Java), but PyIceberg's `snapshot_as_of_timestamp` iterates
history reversed, so it's last-wins. A table with duplicate `timestamp-ms`
entries resolves to a different snapshot across engines — one docstring
sentence calling that out would save someone a confusing debugging session.
##########
crates/iceberg/src/util/snapshot.rs:
##########
@@ -93,6 +122,76 @@ mod tests {
std::sync::Arc::new(fixture.table.metadata().clone())
}
+ type History = [(i64, i64)];
+
+ fn metadata_with_history(history: &History) -> TableMetadataRef {
+ let mut metadata = metadata().as_ref().clone();
+ metadata.snapshot_log = history
+ .iter()
+ .map(|&(timestamp_ms, snapshot_id)| SnapshotLog {
+ timestamp_ms,
+ snapshot_id,
+ })
+ .collect();
+ metadata.into()
+ }
+
+ #[test]
+ fn test_snapshot_id_as_of_time() {
+ let cases: &[(&History, i64, i64)] = &[
+ (&[(1000, S1), (2000, S2)], 1000, S1),
+ (&[(1000, S1), (2000, S2)], 1500, S1),
+ (&[(1000, S1), (2000, S2)], 2000, S2),
+ (&[(1000, S1), (2000, S2)], 3000, S2),
+ // Rollback records when an existing snapshot becomes current
again.
+ (&[(1000, S1), (2000, S2), (3000, S1)], 3500, S1),
+ // The first equal maximum wins; max_by_key would pick S3.
+ (&[(1000, S1), (2000, S2), (2000, S3)], 2000, S2),
+ // Clock skew means the log need not be sorted by timestamp.
+ (&[(1000, S1), (3000, S2), (2500, S3)], 3500, S2),
+ (&[(1000, S1), (3000, S2), (2500, S3)], 2600, S3),
+ (&[(-2000, S1), (-1000, S2)], -1500, S1),
+ (&[(i64::MIN, S1), (0, S2)], i64::MIN, S1),
Review Comment:
Worth flagging honestly: this case actually diverges from Java. `reduce`
seeds on the first filtered element, so querying `i64::MIN` against an
`i64::MIN` entry succeeds and returns S1 — Java seeds `bestTimestamp =
Long.MIN_VALUE` and updates only on strict `>`, so `MIN > MIN` fails and it
returns null (error). No real epoch timestamp hits this, so I wouldn't change
the code — but I'd soften the "matches Java exactly" framing or drop this case,
so a later parity audit doesn't read it as a bug.
##########
crates/iceberg/src/scan/mod.rs:
##########
@@ -1956,6 +1986,329 @@ pub mod tests {
);
}
+ #[test]
+ fn test_table_scan_as_of_time_after_rollback() {
+ for elapsed_ms in [0, 1000] {
+ let table = TableTestFixture::new().table;
+ let mut metadata = table.metadata().clone();
+ let original = metadata.history()[0].snapshot_id;
+ let previous = metadata.history()[1].clone();
+ let rollback_time = previous.timestamp_ms + elapsed_ms;
+ metadata.snapshot_log.push(crate::spec::SnapshotLog {
+ timestamp_ms: rollback_time,
+ snapshot_id: original,
+ });
+ metadata.current_snapshot_id = Some(original);
+ metadata.refs.get_mut(MAIN_BRANCH).unwrap().snapshot_id = original;
+ let table = table.with_metadata(Arc::new(metadata));
+ let scan = table.scan().as_of_time(rollback_time).build().unwrap();
+ // A rollback at the same timestamp keeps the earlier history
entry.
+ let expected = if elapsed_ms == 0 {
+ previous.snapshot_id
+ } else {
+ original
+ };
+ assert_eq!(scan.snapshot().unwrap().snapshot_id(), expected);
+ }
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_rejects_empty_history() {
+ let table = TableTestFixture::new_empty().table;
+ assert!(table.scan().build().unwrap().snapshot().is_none());
+ let err = table.scan().as_of_time(0).build().unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ assert!(err.message().contains("No snapshot history"));
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_rejects_before_history() {
+ let table = TableTestFixture::new().table;
+ let before = table.metadata().history()[0].timestamp_ms - 1;
+ let err = table.scan().as_of_time(before).build().unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ assert!(err.message().contains(&before.to_string()));
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_rejects_all_qualifying_snapshots_expired() {
+ let table = TableTestFixture::new_with_deep_history().table;
+ let mut metadata = table.metadata().clone();
+ let entry = metadata.history()[3].clone();
+ let expired_ids: Vec<_> = metadata
+ .history()
+ .iter()
+ .filter(|log| log.timestamp_ms <= entry.timestamp_ms)
+ .map(|log| log.snapshot_id)
+ .collect();
+ assert_eq!(expired_ids.len(), 4);
+ for snapshot_id in expired_ids {
+ metadata.snapshots.remove(&snapshot_id);
+ }
+ let table = table.with_metadata(Arc::new(metadata));
+ // History resolution succeeds even though every qualifying snapshot
expired.
+ assert_eq!(
+ snapshot_id_as_of_time(&table.metadata_ref(),
entry.timestamp_ms).unwrap(),
+ entry.snapshot_id
+ );
+ assert!(table.scan().build().unwrap().snapshot().is_some());
+ let err = table
+ .scan()
+ .as_of_time(entry.timestamp_ms)
+ .build()
+ .unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ assert_eq!(
+ err.message(),
+ format!("Snapshot with id {} not found", entry.snapshot_id)
+ );
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_conflicts_with_snapshot_id() {
+ for id_first in [true, false] {
+ let table = TableTestFixture::new().table;
+ let entry = &table.metadata().history()[0];
+ let scan = table.scan();
+ let scan = if id_first {
+ scan.snapshot_id(entry.snapshot_id)
+ .as_of_time(entry.timestamp_ms)
+ } else {
+ scan.as_of_time(entry.timestamp_ms)
+ .snapshot_id(entry.snapshot_id)
+ };
+ let err = scan.build().unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ assert_eq!(err.message(), "Cannot combine snapshot_id and
as_of_time");
+ }
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_last_timestamp_wins() {
+ let table = TableTestFixture::new().table;
+ let first = &table.metadata().history()[0];
+ let scan = table
+ .scan()
+ .as_of_time(i64::MAX)
+ .as_of_time(first.timestamp_ms)
+ .build()
+ .unwrap();
+ assert_eq!(scan.snapshot().unwrap().snapshot_id(), first.snapshot_id);
+ // Keep the existing repeated snapshot_id setter behavior as well.
+ let scan = table
+ .scan()
+ .snapshot_id(-1)
+ .snapshot_id(first.snapshot_id)
+ .build()
+ .unwrap();
+ assert_eq!(scan.snapshot().unwrap().snapshot_id(), first.snapshot_id);
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_uses_snapshot_schema() {
+ for select_current in [false, true] {
+ let table = TableTestFixture::new().table;
+ let mut metadata = table.metadata().clone();
+ let entry =
metadata.history()[usize::from(select_current)].clone();
+ // Simulate a schema-only update after this snapshot: schema 0 has
x only,
+ // while the current table schema also contains y and other
columns.
+
Arc::make_mut(metadata.snapshots.get_mut(&entry.snapshot_id).unwrap()).schema_id
=
+ Some(0);
+ let table = table.with_metadata(Arc::new(metadata));
+ let scan = table
+ .scan()
+ .as_of_time(entry.timestamp_ms)
+ .select(["x"])
+ .with_filter(Reference::new("x").greater_than(Datum::long(0)))
+ .build()
+ .unwrap();
+ let context = scan.plan_context.as_ref().unwrap();
+ assert_eq!(context.snapshot_schema.schema_id(), 0);
+ assert!(context.snapshot_bound_predicate.is_some());
+ assert!(
+ table
+ .scan()
+ .as_of_time(entry.timestamp_ms)
+ .select(["y"])
+ .build()
+ .is_err()
+ );
+ assert!(
+ table
+ .scan()
+ .as_of_time(entry.timestamp_ms)
+
.with_filter(Reference::new("y").greater_than(Datum::long(0)))
+ .build()
+ .is_err()
+ );
+ }
+ }
+
+ #[test]
+ fn test_table_scan_as_of_time_schema_compatibility() {
+ let table = TableTestFixture::new().table;
+ let entry = table.metadata().history()[0].clone();
+ // Older snapshots may omit schema-id and use the current-schema
fallback.
+ assert!(
+ table
+ .metadata()
+ .snapshot_by_id(entry.snapshot_id)
+ .unwrap()
+ .schema_id()
+ .is_none()
+ );
+ let scan =
table.scan().as_of_time(entry.timestamp_ms).build().unwrap();
+ assert_eq!(
+ scan.plan_context.as_ref().unwrap().snapshot_schema,
+ *table.metadata().current_schema()
+ );
+
+ let mut metadata = table.metadata().clone();
+
Arc::make_mut(metadata.snapshots.get_mut(&entry.snapshot_id).unwrap()).schema_id
=
+ Some(1234);
+ let table = table.with_metadata(Arc::new(metadata));
+ let err = table
+ .scan()
+ .as_of_time(entry.timestamp_ms)
+ .build()
+ .unwrap_err();
+ assert_eq!(err.kind(), ErrorKind::DataInvalid);
+ assert!(err.message().contains("Schema with id 1234 not found"));
+ }
+
+ #[tokio::test]
+ async fn test_table_scan_as_of_time_reads_historical_rows() {
+ let mut fixture = TableTestFixture::new();
+ fixture.setup_manifest_files().await;
+ let entry = fixture.table.metadata().history()[0].clone();
+
+ // The older snapshot has only x, while the current snapshot has eight
columns.
+ let mut metadata = fixture.table.metadata().clone();
+
Arc::make_mut(metadata.snapshots.get_mut(&entry.snapshot_id).unwrap()).schema_id
= Some(0);
+ fixture.table = fixture.table.with_metadata(Arc::new(metadata));
+ let snapshot = fixture
+ .table
+ .metadata()
+ .snapshot_by_id(entry.snapshot_id)
+ .unwrap();
+ let schema = snapshot.schema(fixture.table.metadata()).unwrap();
+ let arrow_schema =
Arc::new(crate::arrow::schema_to_arrow_schema(&schema).unwrap());
+
+ // Write distinct historical data: x = 100, versus x = 1 in the
current files.
+ let path = format!("{}/historical.parquet", fixture.table_location);
+ let batch = RecordBatch::try_new(arrow_schema.clone(),
vec![Arc::new(Int64Array::from(
+ vec![100],
+ ))])
+ .unwrap();
+ let mut writer =
+ ArrowWriter::try_new(File::create(&path).unwrap(), arrow_schema,
None).unwrap();
+ writer.write(&batch).unwrap();
+ writer.close().unwrap();
+
+ let mut manifest_writer = ManifestWriterBuilder::new(
+ fixture.next_manifest_file(),
+ Some(snapshot.snapshot_id()),
+ schema,
+ fixture
+ .table
+ .metadata()
+ .default_partition_spec()
+ .as_ref()
+ .clone(),
+ )
+ .build_v2_data();
+ manifest_writer
+ .add_entry(
+ ManifestEntry::builder()
+ .status(ManifestStatus::Added)
+ .data_file(
+ DataFileBuilder::default()
+ .partition_spec_id(0)
+ .content(DataContentType::Data)
+ .file_format(DataFileFormat::Parquet)
+
.file_size_in_bytes(fs::metadata(&path).unwrap().len())
+ .file_path(path)
+ .record_count(1)
+
.partition(Struct::from_iter([Some(Literal::long(100))]))
+ .build()
+ .unwrap(),
+ )
+ .build(),
+ )
+ .unwrap();
+ let manifest = manifest_writer.write_manifest_file().await.unwrap();
+ let output = fixture
+ .table
+ .file_io()
+ .new_output(snapshot.manifest_list())
+ .unwrap()
+ .writer()
+ .await
+ .unwrap();
+ let mut list_writer = ManifestListWriter::v2(
+ output,
+ snapshot.snapshot_id(),
+ snapshot.parent_snapshot_id(),
+ snapshot.sequence_number(),
+ );
+ list_writer.add_manifests([manifest].into_iter()).unwrap();
+ list_writer.close().await.unwrap();
+
+ let table = fixture.table;
+ let [by_time, by_id] = [
+ table.scan().as_of_time(entry.timestamp_ms),
+ table.scan().snapshot_id(entry.snapshot_id),
+ ]
+ .map(|scan| {
+
scan.with_filter(Reference::new("x").greater_than_or_equal_to(Datum::long(100)))
+ .build()
+ .unwrap()
+ });
+ let current = table.scan().build().unwrap();
+ let mut planned_files = Vec::new();
+ let mut rows = Vec::new();
+ let mut schemas = Vec::new();
+ for scan in [by_time, by_id, current] {
+ let mut tasks: Vec<_> = scan
+ .plan_files()
+ .await
+ .unwrap()
+ .try_collect()
+ .await
+ .unwrap();
+ tasks.sort_by_key(|task| task.data_file_path().to_string());
+ planned_files.push(tasks);
+ let batches: Vec<_> =
scan.to_arrow().await.unwrap().try_collect().await.unwrap();
+ assert!(!batches.is_empty());
+ schemas.push(batches[0].schema());
+ let mut values: Vec<i64> = batches
+ .iter()
+ .flat_map(|batch| {
+ batch
+ .column_by_name("x")
+ .unwrap()
+ .as_any()
+ .downcast_ref::<Int64Array>()
+ .unwrap()
+ .values()
+ .iter()
+ .copied()
+ })
+ .collect();
+ values.sort_unstable();
+ rows.push(values);
+ }
+ assert!(!planned_files[0].is_empty());
+ assert_eq!(planned_files[0], planned_files[1]);
+ assert_eq!(rows[0], vec![100]);
+ assert_eq!(rows[0], rows[1]);
+ // Two live fixture files contain 1,024 rows each.
+ assert_eq!(rows[2], vec![1; 2048]);
Review Comment:
The `2048` is still decoupled from its source of truth — if
`write_parquet_data_files`' `vec![1; 1024]` ever changes, this assertion
silently goes stale. A `const ROWS_PER_FIXTURE_FILE: usize = 1024;` with
`vec![1; 2 * ROWS_PER_FIXTURE_FILE]` ties them back together.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]