From d4c0d9032ec02416a20311b89da3e38dac651648 Mon Sep 17 00:00:00 2001 From: osipovartem Date: Tue, 22 Sep 2026 19:31:36 +0300 Subject: [PATCH 1/2] Write Iceberg v3 row lineage columns --- datafusion_iceberg/src/row_lineage.rs | 8 +- datafusion_iceberg/src/table/mod.rs | 167 +++++++++++++++++++++- iceberg-rust-spec/src/spec/mod.rs | 1 + iceberg-rust-spec/src/spec/row_lineage.rs | 17 +++ iceberg-rust/src/file_format/parquet.rs | 8 ++ 5 files changed, 194 insertions(+), 7 deletions(-) create mode 100644 iceberg-rust-spec/src/spec/row_lineage.rs diff --git a/datafusion_iceberg/src/row_lineage.rs b/datafusion_iceberg/src/row_lineage.rs index b92cfa1e..eb8b998e 100644 --- a/datafusion_iceberg/src/row_lineage.rs +++ b/datafusion_iceberg/src/row_lineage.rs @@ -10,15 +10,15 @@ use datafusion::common::{exec_err, Result}; use datafusion::physical_plan::PhysicalExpr; use datafusion_expr::ColumnarValue; use iceberg_rust::spec::arrow::schema::PARQUET_FIELD_ID_META_KEY; +pub use iceberg_rust::spec::row_lineage::{ + LAST_UPDATED_SEQUENCE_NUMBER_COLUMN_NAME as LAST_UPDATED_SEQUENCE_NUMBER_COLUMN, + LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID, ROW_ID_COLUMN_NAME as ROW_ID_COLUMN, ROW_ID_FIELD_ID, +}; -pub(crate) const ROW_ID_COLUMN: &str = "_row_id"; -pub(crate) const LAST_UPDATED_SEQUENCE_NUMBER_COLUMN: &str = "_last_updated_sequence_number"; pub(crate) const PHYSICAL_ROW_ID_COLUMN: &str = "__iceberg_physical_row_id"; pub(crate) const PHYSICAL_LAST_UPDATED_SEQUENCE_NUMBER_COLUMN: &str = "__iceberg_physical_last_updated_sequence_number"; pub(crate) const FIRST_ROW_ID_COLUMN: &str = "__iceberg_first_row_id"; -pub(crate) const ROW_ID_FIELD_ID: i32 = i32::MAX - 107; -pub(crate) const LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID: i32 = i32::MAX - 108; #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] pub(crate) enum RowLineageKind { diff --git a/datafusion_iceberg/src/table/mod.rs b/datafusion_iceberg/src/table/mod.rs index 6144d2e7..e46c7144 100644 --- a/datafusion_iceberg/src/table/mod.rs +++ b/datafusion_iceberg/src/table/mod.rs @@ -108,6 +108,7 @@ use iceberg_rust::spec::{ partition::Transform, schema::Schema, sort::{NullOrder, SortDirection, SortOrder}, + table_metadata::FormatVersion, view_metadata::ViewRepresentation, }; use iceberg_rust::{ @@ -2173,7 +2174,26 @@ pub async fn write_parquet_data_files( context: &Arc, branch: Option<&str>, ) -> Result, DataFusionError> { - write_parquet_files(table, batches, context, None, branch).await + write_parquet_files(table, batches, context, None, branch, false).await +} + +/// Writes copy-on-write data while preserving Iceberg v3 row lineage. +/// +/// The input must contain the table columns followed by nullable `_row_id` and +/// `_last_updated_sequence_number` Int64 columns. Existing rows carry their +/// current values; modified rows carry their row ID and a null last-update +/// value; new rows carry nulls for both. Iceberg inheritance fills null values +/// from the committed file metadata without a post-commit rewrite. +pub async fn write_parquet_data_files_with_row_lineage( + table: &Table, + batches: SendableRecordBatchStream, + context: &Arc, + branch: Option<&str>, +) -> Result, DataFusionError> { + if table.metadata().format_version != FormatVersion::V3 { + return plan_err!("Iceberg row-lineage writes require a v3 table"); + } + write_parquet_files(table, batches, context, None, branch, true).await } /// Writes record batches as Parquet equality delete files to an Iceberg table. @@ -2204,7 +2224,7 @@ pub async fn write_parquet_equality_delete_files( equality_ids: &[i32], branch: Option<&str>, ) -> Result, DataFusionError> { - write_parquet_files(table, batches, context, Some(equality_ids), branch).await + write_parquet_files(table, batches, context, Some(equality_ids), branch, false).await } #[instrument(name = "datafusion_iceberg::write_parquet_files", level = "debug", skip(table, batches, context), fields( @@ -2217,6 +2237,7 @@ async fn write_parquet_files( context: &Arc, equality_ids: Option<&[i32]>, branch: Option<&str>, + preserve_row_lineage: bool, ) -> Result, DataFusionError> { let object_store = table.object_store(); let metadata = table.metadata(); @@ -2225,9 +2246,26 @@ async fn write_parquet_files( let schema = table .current_schema() .map_err(DataFusionIcebergError::from)?; - let arrow_schema = Arc::new( + let table_arrow_schema = Arc::new( TryInto::::try_into(schema.fields()).map_err(DataFusionIcebergError::from)?, ); + let arrow_schema = if preserve_row_lineage { + let mut builder = SchemaBuilder::from(table_arrow_schema.as_ref().clone()); + builder.push(row_lineage_field(ROW_ID_COLUMN, ROW_ID_FIELD_ID)); + builder.push(row_lineage_field( + LAST_UPDATED_SEQUENCE_NUMBER_COLUMN, + LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID, + )); + let write_schema = Arc::new(builder.finish()); + if !write_schema.equivalent_names_and_types(&batches.schema()) { + return plan_err!( + "Iceberg row-lineage write input must contain table columns followed by nullable {ROW_ID_COLUMN} and {LAST_UPDATED_SEQUENCE_NUMBER_COLUMN} Int64 columns" + ); + } + write_schema + } else { + table_arrow_schema + }; let partition_fields = metadata .current_partition_fields() @@ -2480,21 +2518,29 @@ mod tests { use datafusion::{ arrow::{ array::{BooleanArray, Int64Array, StringArray}, + datatypes::{DataType, Field, Schema as ArrowSchema}, record_batch::RecordBatch, }, execution::object_store::ObjectStoreUrl, + physical_plan::stream::RecordBatchStreamAdapter, prelude::SessionContext, scalar::ScalarValue, }; + use futures::stream; use iceberg_rust::{ catalog::tabular::Tabular, object_store::ObjectStoreBuilder, spec::{ namespace::Namespace, partition::{PartitionField, Transform}, + row_lineage::{ + LAST_UPDATED_SEQUENCE_NUMBER_COLUMN_NAME, LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID, + ROW_ID_COLUMN_NAME, ROW_ID_FIELD_ID, + }, schema::Schema, snapshot::SnapshotBuilder, types::{PrimitiveType, StructField, Type}, + util, }, }; use iceberg_rust::{ @@ -2507,6 +2553,8 @@ mod tests { view::View, }; use iceberg_sql_catalog::SqlCatalog; + use object_store::ObjectStoreExt; + use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; use std::sync::Arc; @@ -2839,6 +2887,119 @@ mod tests { assert_eq!(int_values(4), vec![Some(2), Some(2), Some(3)]); } + #[tokio::test] + async fn test_v3_row_lineage_writer_preserves_physical_values() { + let object_store = ObjectStoreBuilder::memory(); + let catalog: Arc = Arc::new( + SqlCatalog::new("sqlite://", "test", object_store) + .await + .expect("create catalog"), + ); + let schema = Schema::builder() + .with_struct_field(StructField { + id: 1, + name: "id".to_owned(), + required: true, + field_type: Type::Primitive(PrimitiveType::Long), + doc: None, + initial_default: None, + write_default: None, + }) + .build() + .expect("build schema"); + let table = Table::builder() + .with_name("lineage_writer") + .with_location("memory:///test/lineage_writer") + .with_schema(schema) + .with_property(("format-version".to_owned(), "3".to_owned())) + .build(&["test".to_owned()], catalog) + .await + .expect("create v3 table"); + + let write_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int64, false), + super::row_lineage_field(ROW_ID_COLUMN_NAME, ROW_ID_FIELD_ID), + super::row_lineage_field( + LAST_UPDATED_SEQUENCE_NUMBER_COLUMN_NAME, + LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID, + ), + ])); + let batch = RecordBatch::try_new( + write_schema.clone(), + vec![ + Arc::new(Int64Array::from(vec![10, 20, 30])), + Arc::new(Int64Array::from(vec![Some(7), Some(8), None])), + Arc::new(Int64Array::from(vec![Some(1), None, None])), + ], + ) + .expect("build lineage batch"); + let batches = Box::pin(RecordBatchStreamAdapter::new( + write_schema, + stream::iter([Ok(batch)]), + )); + let ctx = SessionContext::new(); + let files = super::write_parquet_data_files_with_row_lineage( + &table, + batches, + &ctx.task_ctx(), + None, + ) + .await + .expect("write lineage parquet"); + assert_eq!(files.len(), 1); + let file = &files[0]; + assert_eq!(*file.record_count(), 3); + for metrics in [ + file.column_sizes().as_ref(), + file.value_counts().as_ref(), + file.null_value_counts().as_ref(), + ] { + let metrics = metrics.expect("data file metrics"); + assert!(!metrics.contains_key(&ROW_ID_FIELD_ID)); + assert!(!metrics.contains_key(&LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID)); + } + + let path = util::strip_prefix(file.file_path()).into(); + let bytes = table + .object_store() + .get(&path) + .await + .expect("read lineage parquet") + .bytes() + .await + .expect("load lineage parquet bytes"); + let builder = + ParquetRecordBatchReaderBuilder::try_new(bytes).expect("open lineage parquet reader"); + let columns = builder.metadata().file_metadata().schema_descr().columns(); + assert_eq!( + columns[1].self_type().get_basic_info().id(), + ROW_ID_FIELD_ID + ); + assert_eq!( + columns[2].self_type().get_basic_info().id(), + LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID + ); + let output = builder + .build() + .expect("build lineage parquet reader") + .next() + .transpose() + .expect("read lineage parquet batch") + .expect("lineage parquet batch"); + let values = |column: usize| { + output + .column(column) + .as_any() + .downcast_ref::() + .expect("Int64 lineage column") + .iter() + .collect::>() + }; + assert_eq!(values(0), vec![Some(10), Some(20), Some(30)]); + assert_eq!(values(1), vec![Some(7), Some(8), None]); + assert_eq!(values(2), vec![Some(1), None, None]); + } + #[tokio::test] async fn test_changed_file_scan_preserves_intermediate_snapshot_manifests() { let object_store = ObjectStoreBuilder::memory(); diff --git a/iceberg-rust-spec/src/spec/mod.rs b/iceberg-rust-spec/src/spec/mod.rs index f809f0ce..9791ceda 100644 --- a/iceberg-rust-spec/src/spec/mod.rs +++ b/iceberg-rust-spec/src/spec/mod.rs @@ -23,6 +23,7 @@ pub mod materialized_view_metadata; pub mod namespace; pub mod partition; pub mod puffin; +pub mod row_lineage; pub mod schema; pub mod snapshot; pub mod sort; diff --git a/iceberg-rust-spec/src/spec/row_lineage.rs b/iceberg-rust-spec/src/spec/row_lineage.rs new file mode 100644 index 00000000..d3d71544 --- /dev/null +++ b/iceberg-rust-spec/src/spec/row_lineage.rs @@ -0,0 +1,17 @@ +//! Iceberg v3 row-lineage metadata columns. + +/// Logical metadata column carrying a row's stable table-wide identifier. +pub const ROW_ID_COLUMN_NAME: &str = "_row_id"; +/// Logical metadata column carrying the sequence number of a row's last update. +pub const LAST_UPDATED_SEQUENCE_NUMBER_COLUMN_NAME: &str = "_last_updated_sequence_number"; + +/// Reserved Iceberg field ID for [`ROW_ID_COLUMN_NAME`]. +pub const ROW_ID_FIELD_ID: i32 = i32::MAX - 107; +/// Reserved Iceberg field ID for [`LAST_UPDATED_SEQUENCE_NUMBER_COLUMN_NAME`]. +pub const LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID: i32 = i32::MAX - 108; + +/// Returns whether `field_id` identifies an Iceberg row-lineage metadata column. +#[must_use] +pub const fn is_row_lineage_field_id(field_id: i32) -> bool { + field_id == ROW_ID_FIELD_ID || field_id == LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID +} diff --git a/iceberg-rust/src/file_format/parquet.rs b/iceberg-rust/src/file_format/parquet.rs index 16780572..09747dc2 100644 --- a/iceberg-rust/src/file_format/parquet.rs +++ b/iceberg-rust/src/file_format/parquet.rs @@ -12,6 +12,7 @@ use iceberg_rust_spec::{ spec::{ manifest::{AvroMap, Content, DataFile, FileFormat}, partition::PartitionField, + row_lineage::is_row_lineage_field_id, schema::Schema, types::{PrimitiveType, Type}, values::{PhysicalTypeHint, Struct, Value}, @@ -116,6 +117,13 @@ pub fn parquet_to_datafile( let mut counted_logical_values = HashSet::new(); let mut variant_null_counts = HashMap::::new(); for column in row_group.columns() { + let parquet_field = column.column_descr().self_type().get_basic_info(); + if parquet_field.has_id() && is_row_lineage_field_id(parquet_field.id()) { + // Row-lineage columns are reserved metadata, not table fields. Their + // values are still written to Parquet, but v1-v3 DataFile metrics do + // not have table-schema entries through which to describe them. + continue; + } let column_name = column.column_descr().name(); let column_path = column.column_path().parts().join("."); let id = schema From 3c91ba996a78a119800703f289344f651fc73b9a Mon Sep 17 00:00:00 2001 From: osipovartem Date: Tue, 22 Sep 2026 19:58:03 +0300 Subject: [PATCH 2/2] Preserve row lineage during v3 overwrites --- iceberg-rust/src/table/manifest.rs | 53 +++++++-- iceberg-rust/src/table/manifest_list.rs | 56 +++++++-- .../src/table/transaction/operation.rs | 109 ++++++++++++++++-- 3 files changed, 188 insertions(+), 30 deletions(-) diff --git a/iceberg-rust/src/table/manifest.rs b/iceberg-rust/src/table/manifest.rs index f4aad012..96e3dbb1 100644 --- a/iceberg-rust/src/table/manifest.rs +++ b/iceberg-rust/src/table/manifest.rs @@ -180,6 +180,36 @@ impl FilteredManifestStats { } } +fn inherit_first_row_id( + content: manifest_list::Content, + next_row_id: &mut Option, + data_file: &mut iceberg_rust_spec::manifest::DataFile, +) -> Result<(), Error> { + if content != manifest_list::Content::Data || data_file.first_row_id().is_some() { + return Ok(()); + } + + let Some(first_row_id) = *next_row_id else { + return Ok(()); + }; + let record_count = *data_file.record_count(); + if first_row_id < 0 || record_count < 0 { + return Err(Error::InvalidFormat(format!( + "Invalid row lineage range for data file {}: first_row_id={first_row_id}, record_count={record_count}", + data_file.file_path() + ))); + } + + *data_file.first_row_id_mut() = Some(first_row_id); + *next_row_id = Some(first_row_id.checked_add(record_count).ok_or_else(|| { + Error::InvalidFormat(format!( + "Row ID overflow while rewriting data file {}", + data_file.file_path() + )) + })?); + Ok(()) +} + impl<'schema, 'metadata> ManifestWriter<'schema, 'metadata> { /// Creates a new ManifestWriter for writing manifest entries to a new manifest file. /// @@ -465,6 +495,7 @@ impl<'schema, 'metadata> ManifestWriter<'schema, 'metadata> { let inherited_snapshot_id = manifest.added_snapshot_id; let current_schema = table_metadata.current_schema()?; let manifest_reader = ManifestReader::new(bytes)?; + let mut next_inherited_row_id = manifest.first_row_id; let mut writer = AvroWriter::new(schema, Vec::new()); let mut filtered_stats = FilteredManifestStats::default(); @@ -523,17 +554,16 @@ impl<'schema, 'metadata> ManifestWriter<'schema, 'metadata> { }, )?; - writer.extend(manifest_reader.filter_map(|entry| { - let mut entry = entry - .map_err(|err| { - apache_avro::Error::new(apache_avro::error::Details::DeserializeValue( - err.to_string(), - )) - }) - .unwrap(); + for entry in manifest_reader { + let mut entry = entry?; if *entry.status() == Status::Deleted { - return None; + continue; } + inherit_first_row_id( + manifest.content, + &mut next_inherited_row_id, + entry.data_file_mut(), + )?; entry .data_file_mut() .promote_bounds_to_schema(current_schema); @@ -552,14 +582,13 @@ impl<'schema, 'metadata> ManifestWriter<'schema, 'metadata> { filtered_stats.removed_data_files += 1; *entry.status_mut() = Status::Deleted; filtered_stats.filtered_entries.push(entry); - None } else { existing_files += 1; existing_rows += entry.data_file().record_count(); *entry.status_mut() = Status::Existing; - Some(to_value(entry).unwrap()) + writer.append_ser(entry)?; } - }))?; + } manifest.sequence_number = table_metadata.last_sequence_number + 1; manifest.added_snapshot_id = snapshot_id; diff --git a/iceberg-rust/src/table/manifest_list.rs b/iceberg-rust/src/table/manifest_list.rs index 500b044b..69a2029a 100644 --- a/iceberg-rust/src/table/manifest_list.rs +++ b/iceberg-rust/src/table/manifest_list.rs @@ -590,6 +590,46 @@ impl<'schema, 'metadata> ManifestListWriter<'schema, 'metadata> { let mut writer = AvroWriter::new(schema, Vec::new()); + // Rewriting an unaffected v3 manifest is both unnecessary I/O and a + // row-lineage hazard. Preserve it by path; affected manifests are + // rewritten below after their inherited row IDs are materialized. + if table_metadata.format_version == FormatVersion::V3 { + let mut manifests = Vec::new(); + let mut file_count_all_entries = 0usize; + for manifest in manifest_list_reader { + let manifest = manifest?; + let file_count = manifest + .added_files_count + .unwrap_or(0) + .checked_add(manifest.existing_files_count.unwrap_or(0)) + .ok_or_else(|| Error::InvalidFormat("manifest file count".to_string()))?; + file_count_all_entries = file_count_all_entries + .checked_add(file_count.try_into()?) + .ok_or_else(|| Error::InvalidFormat("manifest file count".to_string()))?; + + if manifests_to_overwrite.contains(&manifest.manifest_path) { + manifests.push(manifest); + } else { + writer.append_ser(manifest)?; + } + } + + return Ok(( + Self { + table_metadata, + writer, + selected_data_manifest: None, + selected_delete_manifest: None, + bounding_partition_values, + n_existing_files: file_count_all_entries, + commit_uuid, + manifest_count: 0, + next_row_id: Some(table_metadata.next_row_id), + }, + manifests, + )); + } + let OverwriteManifest { manifest, file_count_all_entries, @@ -1412,17 +1452,17 @@ impl<'schema, 'metadata> ManifestListWriter<'schema, 'metadata> { if manifest.content == Content::Data { if let Some(next_row_id) = self.next_row_id.as_mut() { let added_rows = manifest.added_rows_count.unwrap_or(0); - if added_rows < 0 { + let existing_rows = manifest.existing_rows_count.unwrap_or(0); + if added_rows < 0 || existing_rows < 0 { return Err(Error::InvalidFormat( - "manifest added row count must be non-negative".to_string(), + "manifest row counts must be non-negative".to_string(), )); } - if added_rows > 0 { - manifest.first_row_id = Some(*next_row_id); - *next_row_id = next_row_id - .checked_add(added_rows) - .ok_or_else(|| Error::InvalidFormat("next row id overflow".to_string()))?; - } + manifest.first_row_id = Some(*next_row_id); + *next_row_id = next_row_id + .checked_add(existing_rows) + .and_then(|value| value.checked_add(added_rows)) + .ok_or_else(|| Error::InvalidFormat("next row id overflow".to_string()))?; } } self.writer.append_ser(manifest)?; diff --git a/iceberg-rust/src/table/transaction/operation.rs b/iceberg-rust/src/table/transaction/operation.rs index 5a890ee0..33166f41 100644 --- a/iceberg-rust/src/table/transaction/operation.rs +++ b/iceberg-rust/src/table/transaction/operation.rs @@ -130,10 +130,7 @@ impl Operation { object_store: Arc, ) -> Result<(Option, Vec), Error> { if table_metadata.format_version == FormatVersion::V3 - && matches!( - &self, - Operation::Replace { .. } | Operation::Overwrite { .. } - ) + && matches!(&self, Operation::Replace { .. }) { return Err(Error::NotSupported( "Iceberg v3 manifest rewrites with row lineage".to_string(), @@ -731,12 +728,9 @@ impl Operation { .build() .map_err(Error::from) }); - let selected_manifest_location = manifest_list_writer + let files_to_filter = manifest_list_writer .selected_data_manifest() - .map(|x| x.manifest_path.clone()) - .ok_or(Error::NotFound("Selected manifest".to_owned()))?; - let files_to_filter = files_to_overwrite - .get(&selected_manifest_location) + .and_then(|manifest| files_to_overwrite.get(&manifest.manifest_path)) .map(|filter_files| filter_files.iter().cloned().collect::>()); let selected_filter_stats = if n_splits == 0 { @@ -795,7 +789,7 @@ impl Operation { } } - let (new_manifest_list_location, _) = manifest_list_writer + let (new_manifest_list_location, next_row_id) = manifest_list_writer .finish(snapshot_id, object_store) .await?; @@ -820,6 +814,7 @@ impl Operation { }) .with_schema_id(*table_metadata.current_schema()?.schema_id()); snapshot_builder.with_parent_snapshot_id(*old_snapshot.snapshot_id()); + apply_v3_row_lineage(&mut snapshot_builder, table_metadata, next_row_id)?; let snapshot = snapshot_builder.build()?; Ok(( @@ -1262,6 +1257,7 @@ pub fn compute_n_splits( #[cfg(test)] mod tests { use super::*; + use crate::table::ManifestReader; use futures::executor::block_on; use iceberg_rust_spec::manifest::FileFormat; use iceberg_rust_spec::spec::schema::SchemaBuilder; @@ -1442,6 +1438,99 @@ mod tests { assert_eq!(metadata.next_row_id, 22); } + #[tokio::test] + async fn v3_overwrite_preserves_existing_row_ids_and_reserves_a_new_range() { + let mut metadata = sample_metadata(&[], None, &[]); + metadata.format_version = FormatVersion::V3; + metadata.next_row_id = 10; + let store = Arc::new(InMemory::new()); + let old_path = "s3://tests/table/data/old.parquet"; + + let append = Operation::Append { + branch: None, + data_files: vec![data_file(old_path, 5)], + delete_files: Vec::new(), + additional_summary: None, + }; + let (_, append_updates) = append.execute(&metadata, store.clone()).await.unwrap(); + crate::catalog::commit::apply_table_updates(&mut metadata, append_updates).unwrap(); + assert_eq!(metadata.next_row_id, 15); + + let current_snapshot = metadata.current_snapshot(None).unwrap().unwrap(); + let manifest_list_bytes = store + .get(&strip_prefix(current_snapshot.manifest_list()).into()) + .await + .unwrap() + .bytes() + .await + .unwrap(); + let original_manifest = ManifestListReader::new(&manifest_list_bytes[..], &metadata) + .unwrap() + .next() + .unwrap() + .unwrap(); + + let mut replacement = data_file("s3://tests/table/data/replacement.parquet", 5); + *replacement.first_row_id_mut() = Some(10); + let overwrite = Operation::Overwrite { + branch: None, + data_files: vec![replacement], + files_to_overwrite: HashMap::from([( + original_manifest.manifest_path.clone(), + vec![old_path.to_string()], + )]), + additional_summary: None, + }; + let (_, overwrite_updates) = overwrite.execute(&metadata, store.clone()).await.unwrap(); + let overwrite_snapshot = overwrite_updates + .iter() + .find_map(|update| match update { + TableUpdate::AddSnapshot { snapshot } => Some(snapshot), + _ => None, + }) + .unwrap(); + + assert_eq!(*overwrite_snapshot.first_row_id(), Some(15)); + assert_eq!(*overwrite_snapshot.added_rows(), Some(5)); + + let manifest_list_bytes = store + .get(&strip_prefix(overwrite_snapshot.manifest_list()).into()) + .await + .unwrap() + .bytes() + .await + .unwrap(); + let manifests = ManifestListReader::new(&manifest_list_bytes[..], &metadata) + .unwrap() + .collect::, _>>() + .unwrap(); + let mut row_ids_by_path = HashMap::new(); + for manifest in manifests { + let bytes = store + .get(&strip_prefix(&manifest.manifest_path).into()) + .await + .unwrap() + .bytes() + .await + .unwrap(); + for entry in ManifestReader::new(&bytes[..]).unwrap() { + let entry = entry.unwrap(); + row_ids_by_path.insert( + entry.data_file().file_path().clone(), + *entry.data_file().first_row_id(), + ); + } + } + assert_eq!(row_ids_by_path.get(old_path), Some(&Some(10))); + assert_eq!( + row_ids_by_path.get("s3://tests/table/data/replacement.parquet"), + Some(&Some(10)) + ); + + crate::catalog::commit::apply_table_updates(&mut metadata, overwrite_updates).unwrap(); + assert_eq!(metadata.next_row_id, 20); + } + #[tokio::test] async fn v3_manifest_rewrites_remain_rejected() { let mut metadata = sample_metadata(&[], None, &[]);