From c226076e243b59788b6449ff6d9a31d61a9195c6 Mon Sep 17 00:00:00 2001 From: Anoop Johnson Date: Mon, 3 Aug 2026 09:14:30 -0700 Subject: [PATCH 1/2] feat: carry first_row_id and data sequence number on FileScanTask Manifest read now inherits a data file's first_row_id (#2922) Add both fields to FileScanTask and populate them in ManifestEntryContext::into_file_scan_task from the manifest entry's data_file().first_row_id() and sequence_number() (the data sequence number, which is what _last_updated_sequence_number maps to, not the file sequence number). This is step 2 of reading the row-lineage metadata columns; positional materialization in the arrow pipeline follows. --- crates/iceberg/public-api.txt | 4 +- crates/iceberg/src/arrow/reader/row_filter.rs | 4 + crates/iceberg/src/scan/context.rs | 2 + crates/iceberg/src/scan/mod.rs | 138 +++++++++++++++++- crates/iceberg/src/scan/task.rs | 15 ++ 5 files changed, 161 insertions(+), 2 deletions(-) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 12f844b362..d88ffa1b97 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -1266,8 +1266,10 @@ pub struct iceberg::scan::FileScanTask pub iceberg::scan::FileScanTask::case_sensitive: bool pub iceberg::scan::FileScanTask::data_file_format: iceberg::spec::DataFileFormat pub iceberg::scan::FileScanTask::data_file_path: alloc::string::String +pub iceberg::scan::FileScanTask::data_sequence_number: core::option::Option pub iceberg::scan::FileScanTask::deletes: alloc::vec::Vec pub iceberg::scan::FileScanTask::file_size_in_bytes: u64 +pub iceberg::scan::FileScanTask::first_row_id: core::option::Option pub iceberg::scan::FileScanTask::key_metadata: core::option::Option> pub iceberg::scan::FileScanTask::length: u64 pub iceberg::scan::FileScanTask::name_mapping: core::option::Option> @@ -1293,7 +1295,7 @@ impl core::fmt::Debug for iceberg::scan::FileScanTask pub fn iceberg::scan::FileScanTask::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result impl core::marker::StructuralPartialEq for iceberg::scan::FileScanTask impl iceberg::scan::FileScanTask -pub fn iceberg::scan::FileScanTask::builder() -> FileScanTaskBuilder<((), (), (), (), (), (), (), (), (), (), (), (), (), (), (), ())> +pub fn iceberg::scan::FileScanTask::builder() -> FileScanTaskBuilder<((), (), (), (), (), (), (), (), (), (), (), (), (), (), (), (), (), ())> impl serde_core::ser::Serialize for iceberg::scan::FileScanTask pub fn iceberg::scan::FileScanTask::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer impl<'de> serde_core::de::Deserialize<'de> for iceberg::scan::FileScanTask diff --git a/crates/iceberg/src/arrow/reader/row_filter.rs b/crates/iceberg/src/arrow/reader/row_filter.rs index d82b9a0aff..2e6b37ba30 100644 --- a/crates/iceberg/src/arrow/reader/row_filter.rs +++ b/crates/iceberg/src/arrow/reader/row_filter.rs @@ -1138,6 +1138,8 @@ mod tests { start: 0, length: 0, record_count: None, + first_row_id: None, + data_sequence_number: None, data_file_path: file_path.clone(), data_file_format: DataFileFormat::Parquet, schema: iceberg_schema.clone(), @@ -1231,6 +1233,8 @@ mod tests { start: 0, length: 0, record_count: None, + first_row_id: None, + data_sequence_number: None, data_file_path: file_path.clone(), data_file_format: DataFileFormat::Parquet, schema: iceberg_schema.clone(), diff --git a/crates/iceberg/src/scan/context.rs b/crates/iceberg/src/scan/context.rs index 176bcba459..67b8cde4d8 100644 --- a/crates/iceberg/src/scan/context.rs +++ b/crates/iceberg/src/scan/context.rs @@ -133,6 +133,8 @@ impl ManifestEntryContext { .with_start(0) .with_length(self.manifest_entry.file_size_in_bytes()) .with_record_count(Some(self.manifest_entry.record_count())) + .with_first_row_id(self.manifest_entry.data_file().first_row_id()) + .with_data_sequence_number(self.manifest_entry.sequence_number()) .with_data_file_path(self.manifest_entry.file_path().to_string()) .with_data_file_format(self.manifest_entry.file_format()) .with_schema(self.snapshot_schema) diff --git a/crates/iceberg/src/scan/mod.rs b/crates/iceberg/src/scan/mod.rs index a1c85c7117..29b0fd3276 100644 --- a/crates/iceberg/src/scan/mod.rs +++ b/crates/iceberg/src/scan/mod.rs @@ -669,7 +669,7 @@ pub mod tests { use crate::scan::FileScanTask; use crate::spec::{ DEFAULT_SCHEMA_NAME_MAPPING, DataContentType, DataFileBuilder, DataFileFormat, Datum, - Literal, MAIN_BRANCH, ManifestEntry, ManifestListWriter, ManifestStatus, + FormatVersion, Literal, MAIN_BRANCH, ManifestEntry, ManifestListWriter, ManifestStatus, ManifestWriterBuilder, NestedField, Operation, PartitionSpec, PrimitiveType, Schema, Snapshot, Struct, StructType, Summary, TableMetadata, TableMetadataBuilder, Type, UnboundPartitionSpec, @@ -996,6 +996,84 @@ pub mod tests { manifest_list_write.close().await.unwrap(); } + /// Writes a v3 data manifest with a manifest-level `first_row_id` of 0, + /// so live entries inherit a per-file `first_row_id` on read. Upgrades the + /// table to v3 first, so the manifest list is read as v3. + pub async fn setup_v3_manifest_files(&mut self) { + let metadata = TableMetadataBuilder::new_from_metadata( + self.table.metadata().clone(), + self.table.metadata_location().map(str::to_string), + ) + .upgrade_format_version(FormatVersion::V3) + .unwrap() + .build() + .unwrap() + .metadata; + self.table = Table::builder() + .metadata(metadata) + .identifier(self.table.identifier().clone()) + .file_io(self.table.file_io().clone()) + .metadata_location(self.table.metadata_location().unwrap().to_string()) + .runtime(test_runtime()) + .build() + .unwrap(); + + let current_snapshot = self.table.metadata().current_snapshot().unwrap(); + let current_schema = current_snapshot.schema(self.table.metadata()).unwrap(); + let current_partition_spec = self.table.metadata().default_partition_spec(); + + let parquet_file_size = self.write_parquet_data_files(); + + let mut writer = ManifestWriterBuilder::new( + self.next_manifest_file(), + Some(current_snapshot.snapshot_id()), + current_schema.clone(), + current_partition_spec.as_ref().clone(), + ) + .build_v3_data(); + writer + .add_entry( + ManifestEntry::builder() + .status(ManifestStatus::Added) + .data_file( + DataFileBuilder::default() + .partition_spec_id(0) + .content(DataContentType::Data) + .file_path(format!("{}/1.parquet", &self.table_location)) + .file_format(DataFileFormat::Parquet) + .file_size_in_bytes(parquet_file_size) + .record_count(1) + .partition(Struct::from_iter([Some(Literal::long(100))])) + .key_metadata(None) + .build() + .unwrap(), + ) + .build(), + ) + .unwrap(); + let data_file_manifest = writer.write_manifest_file().await.unwrap(); + + let manifest_list_writer = self + .table + .file_io() + .new_output(current_snapshot.manifest_list()) + .unwrap() + .writer() + .await + .unwrap(); + let mut manifest_list_write = ManifestListWriter::v3( + manifest_list_writer, + current_snapshot.snapshot_id(), + current_snapshot.parent_snapshot_id(), + current_snapshot.sequence_number(), + Some(0), + ); + manifest_list_write + .add_manifests(vec![data_file_manifest].into_iter()) + .unwrap(); + manifest_list_write.close().await.unwrap(); + } + pub async fn setup_manifest_files_with_partition_evolution(&mut self) { let current_snapshot = self.table.metadata().current_snapshot().unwrap(); let parent_snapshot = current_snapshot @@ -1862,6 +1940,60 @@ pub mod tests { ); } + #[tokio::test] + async fn test_plan_files_carries_row_lineage_into_file_scan_task() { + let mut fixture = TableTestFixture::new(); + fixture.setup_manifest_files().await; + + let mut tasks: Vec<_> = fixture + .table + .scan() + .build() + .unwrap() + .plan_files() + .await + .unwrap() + .try_collect() + .await + .unwrap(); + + tasks.sort_by_key(|task| task.data_file_path.to_string()); + assert_eq!(tasks.len(), 2); + + // The added file inherits the current snapshot's data sequence number, + // the existing file keeps the one it was written with. + assert_eq!(tasks[0].data_sequence_number, Some(1)); + assert_eq!(tasks[1].data_sequence_number, Some(0)); + + // first_row_id is a v3 concept; a v2 manifest carries none. + assert!(tasks.iter().all(|task| task.first_row_id.is_none())); + } + + #[tokio::test] + async fn test_plan_files_carries_inherited_first_row_id() { + let mut fixture = TableTestFixture::new(); + fixture.setup_v3_manifest_files().await; + + let task = fixture + .table + .scan() + .build() + .unwrap() + .plan_files() + .await + .unwrap() + .try_collect::>() + .await + .unwrap() + .into_iter() + .next() + .expect("expected one FileScanTask"); + + // The manifest-level first_row_id (0) is inherited onto the entry on + // read, then carried onto the task. + assert_eq!(task.first_row_id, Some(0)); + } + #[tokio::test] async fn test_filtered_scan_with_dropped_partition_source_column() { let mut fixture = TableTestFixture::new(); @@ -2441,6 +2573,8 @@ pub mod tests { assert_eq!(task.project_field_ids, deserialized.project_field_ids); assert_eq!(task.predicate, deserialized.predicate); assert_eq!(task.schema, deserialized.schema); + assert_eq!(task.first_row_id, deserialized.first_row_id); + assert_eq!(task.data_sequence_number, deserialized.data_sequence_number); }; // without predicate @@ -2462,6 +2596,8 @@ pub mod tests { .with_project_field_ids(vec![1, 2, 3]) .with_schema(schema.clone()) .with_record_count(Some(100)) + .with_first_row_id(Some(1000)) + .with_data_sequence_number(Some(5)) .with_data_file_format(DataFileFormat::Parquet) .with_case_sensitive(false) .build(); diff --git a/crates/iceberg/src/scan/task.rs b/crates/iceberg/src/scan/task.rs index 24a6036870..1b58b9264d 100644 --- a/crates/iceberg/src/scan/task.rs +++ b/crates/iceberg/src/scan/task.rs @@ -67,6 +67,21 @@ pub struct FileScanTask { #[builder(default)] pub record_count: Option, + /// The first row id assigned to the data file. + /// + /// Used to derive the `_row_id` metadata column: for a row without an + /// explicit `_row_id`, it is this value plus the row's ordinal position. + #[serde(skip_serializing_if = "Option::is_none")] + #[builder(default)] + pub first_row_id: Option, + + /// The data sequence number of the data file. + /// + /// Used to derive the `_last_updated_sequence_number` metadata column. + #[serde(skip_serializing_if = "Option::is_none")] + #[builder(default)] + pub data_sequence_number: Option, + /// The data file path corresponding to the task. pub data_file_path: String, From d48a73ab8aa1d01b3315f57434bb04d4432b4b69 Mon Sep 17 00:00:00 2001 From: Anoop Johnson Date: Tue, 4 Aug 2026 08:01:33 -0700 Subject: [PATCH 2/2] test: address review on row-lineage scan task --- crates/iceberg/src/scan/mod.rs | 22 ++++++++++++++++------ crates/iceberg/src/scan/task.rs | 5 ++++- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/crates/iceberg/src/scan/mod.rs b/crates/iceberg/src/scan/mod.rs index 29b0fd3276..561b41696e 100644 --- a/crates/iceberg/src/scan/mod.rs +++ b/crates/iceberg/src/scan/mod.rs @@ -996,7 +996,7 @@ pub mod tests { manifest_list_write.close().await.unwrap(); } - /// Writes a v3 data manifest with a manifest-level `first_row_id` of 0, + /// Writes a v3 data manifest with a manifest-level `first_row_id` of 42, /// so live entries inherit a per-file `first_row_id` on read. Upgrades the /// table to v3 first, so the manifest list is read as v3. pub async fn setup_v3_manifest_files(&mut self) { @@ -1066,7 +1066,7 @@ pub mod tests { current_snapshot.snapshot_id(), current_snapshot.parent_snapshot_id(), current_snapshot.sequence_number(), - Some(0), + Some(42), ); manifest_list_write .add_manifests(vec![data_file_manifest].into_iter()) @@ -1957,12 +1957,20 @@ pub mod tests { .await .unwrap(); - tasks.sort_by_key(|task| task.data_file_path.to_string()); assert_eq!(tasks.len(), 2); + tasks.sort_by_key(|task| task.data_file_path.to_string()); // The added file inherits the current snapshot's data sequence number, // the existing file keeps the one it was written with. + assert_eq!( + tasks[0].data_file_path, + format!("{}/1.parquet", &fixture.table_location) + ); assert_eq!(tasks[0].data_sequence_number, Some(1)); + assert_eq!( + tasks[1].data_file_path, + format!("{}/3.parquet", &fixture.table_location) + ); assert_eq!(tasks[1].data_sequence_number, Some(0)); // first_row_id is a v3 concept; a v2 manifest carries none. @@ -1970,7 +1978,7 @@ pub mod tests { } #[tokio::test] - async fn test_plan_files_carries_inherited_first_row_id() { + async fn test_plan_files_carries_row_lineage_from_v3_manifest() { let mut fixture = TableTestFixture::new(); fixture.setup_v3_manifest_files().await; @@ -1989,9 +1997,11 @@ pub mod tests { .next() .expect("expected one FileScanTask"); - // The manifest-level first_row_id (0) is inherited onto the entry on + // The manifest-level first_row_id (42) is inherited onto the entry on // read, then carried onto the task. - assert_eq!(task.first_row_id, Some(0)); + assert_eq!(task.first_row_id, Some(42)); + // The data sequence number is threaded through the same v3 read path. + assert_eq!(task.data_sequence_number, Some(1)); } #[tokio::test] diff --git a/crates/iceberg/src/scan/task.rs b/crates/iceberg/src/scan/task.rs index 1b58b9264d..e69397b066 100644 --- a/crates/iceberg/src/scan/task.rs +++ b/crates/iceberg/src/scan/task.rs @@ -75,7 +75,10 @@ pub struct FileScanTask { #[builder(default)] pub first_row_id: Option, - /// The data sequence number of the data file. + /// The data sequence number of the file, as opposed to its file sequence + /// number: the sequence number preserved when a file is carried forward + /// across a rewrite. May be null for an existing entry in a malformed + /// manifest that lacks one. /// /// Used to derive the `_last_updated_sequence_number` metadata column. #[serde(skip_serializing_if = "Option::is_none")]