From 7f2e57f84958310e10002a660ff112b7482b4b26 Mon Sep 17 00:00:00 2001 From: osipovartem Date: Tue, 15 Sep 2026 15:29:57 +0300 Subject: [PATCH 1/2] Add incremental changed data file scans --- datafusion_iceberg/src/table/mod.rs | 230 +++++++++++++++++++++++++--- 1 file changed, 212 insertions(+), 18 deletions(-) diff --git a/datafusion_iceberg/src/table/mod.rs b/datafusion_iceberg/src/table/mod.rs index 8302200a..c331fe23 100644 --- a/datafusion_iceberg/src/table/mod.rs +++ b/datafusion_iceberg/src/table/mod.rs @@ -127,12 +127,20 @@ static MANIFEST_FILE_PATH_COLUMN: &str = "__manifest_file_path"; static DATA_FILE_ROW_POSITION_COLUMN: &str = "__iceberg_file_row_position"; static DATA_FILE_SEQUENCE_NUMBER_COLUMN: &str = "__iceberg_data_sequence_number"; static DELETE_FILE_SEQUENCE_NUMBER_COLUMN: &str = "__iceberg_delete_sequence_number"; +pub const CHANGE_FILE_STATUS_COLUMN: &str = "__iceberg_change_file_status"; static POSITION_DELETE_FILE_PATH_COLUMN: &str = "file_path"; static POSITION_DELETE_POS_COLUMN: &str = "pos"; const POSITION_DELETE_FILE_PATH_FIELD_ID: i32 = i32::MAX - 101; const POSITION_DELETE_POS_FIELD_ID: i32 = i32::MAX - 102; type PhysicalProjection = Vec<(Arc, String)>; +#[derive(Clone, Copy, Debug, Default)] +struct PartitionedFileMetadataColumns { + first_row_id: bool, + sequence_number: bool, + change_status: bool, +} + fn row_lineage_field(name: &str, field_id: i32) -> Field { Field::new(name, DataType::Int64, true).with_metadata(HashMap::from([( PARQUET_FIELD_ID_META_KEY.to_owned(), @@ -232,6 +240,9 @@ pub struct DataFusionTableConfig { /// Expose the Iceberg v3 `_last_updated_sequence_number` metadata column. #[builder(default)] enable_last_updated_sequence_number_column: bool, + /// Scan only data files added or deleted by the selected snapshot range. + #[builder(default)] + scan_changed_data_files: bool, } impl DataFusionTable { @@ -298,6 +309,13 @@ impl DataFusionTable { LAST_UPDATED_SEQUENCE_NUMBER_FIELD_ID, )); } + if config + .as_ref() + .map(|x| x.scan_changed_data_files) + .unwrap_or_default() + { + builder.push(Field::new(CHANGE_FILE_STATUS_COLUMN, DataType::Utf8, false)); + } Arc::new(builder.finish()) } Tabular::View(view) => { @@ -593,6 +611,10 @@ async fn table_scan( let enable_row_lineage = enable_row_id_column || enable_last_updated_sequence_number_column; + let scan_changed_data_files = config + .map(|x| x.scan_changed_data_files) + .unwrap_or_default(); + let partition_fields = &snapshot_range .1 .and_then(|snapshot_id| table.metadata().partition_fields(snapshot_id).ok()) @@ -603,6 +625,13 @@ async fn table_scan( .map(|x| x.and_then(|y| table.metadata().sequence_number(y))) .collect_tuple::<(Option, Option)>() .unwrap(); + let manifest_entry_sequence_range = if scan_changed_data_files { + // The enclosing manifest is already constrained to the snapshot range. Deleted + // entries retain the data file's original sequence number and must not be removed. + (None, None) + } else { + sequence_number_range + }; // If there is a filter expression the manifests to read are pruned based on the pruning statistics available in the manifest_list file. // Row lineage values may be synthesized from manifest metadata, so they must not be @@ -712,7 +741,11 @@ async fn table_scan( pruning_predicate.prune(&PruneManifests::new(partition_fields, &manifests))?; table - .datafiles(&manifests, Some(manifests_to_prune), sequence_number_range) + .datafiles( + &manifests, + Some(manifests_to_prune), + manifest_entry_sequence_range, + ) .await .map_err(DataFusionIcebergError::from)? .try_collect() @@ -720,7 +753,7 @@ async fn table_scan( .map_err(DataFusionIcebergError::from)? } else { table - .datafiles(&manifests, None, sequence_number_range) + .datafiles(&manifests, None, manifest_entry_sequence_range) .await .map_err(DataFusionIcebergError::from)? .try_collect() @@ -762,7 +795,7 @@ async fn table_scan( .await .map_err(DataFusionIcebergError::from)?; let data_files: Vec<_> = table - .datafiles(&manifests, None, sequence_number_range) + .datafiles(&manifests, None, manifest_entry_sequence_range) .await .map_err(DataFusionIcebergError::from)? .try_collect() @@ -794,7 +827,7 @@ async fn table_scan( let mut equality_deletes = Vec::new(); let mut position_deletes = Vec::new(); content_file_iter - .filter(|manifest| *manifest.1.status() != Status::Deleted) + .filter(|manifest| scan_manifest_entry(&manifest.1, scan_changed_data_files)) .for_each(|manifest| match manifest.1.data_file().content() { Content::Data => data_files.push(manifest), Content::EqualityDeletes => equality_deletes.push(manifest), @@ -823,7 +856,7 @@ async fn table_scan( } } else { content_file_iter.for_each(|manifest| { - if *manifest.1.status() != Status::Deleted { + if scan_manifest_entry(&manifest.1, scan_changed_data_files) { match manifest.1.data_file().content() { Content::Data => { data_file_groups @@ -879,6 +912,12 @@ async fn table_scan( .column_statistics .push(ColumnStatistics::new_unknown()); } + if scan_changed_data_files { + table_partition_cols.push(Field::new(CHANGE_FILE_STATUS_COLUMN, DataType::Utf8, false)); + statistics + .column_statistics + .push(ColumnStatistics::new_unknown()); + } let mut table_schema_builder = TableSchema::builder(file_schema.clone()) .with_table_partition_cols( @@ -1014,8 +1053,12 @@ async fn table_scan( &data_manifest.1, last_updated_ms, include_data_file_path_column, - enable_row_lineage, - has_position_deletes || enable_last_updated_sequence_number_column, + PartitionedFileMetadataColumns { + first_row_id: enable_row_lineage, + sequence_number: has_position_deletes + || enable_last_updated_sequence_number_column, + change_status: scan_changed_data_files, + }, manifest_path, ) .unwrap(); @@ -1071,8 +1114,7 @@ async fn table_scan( &delete_manifest.1, last_updated_ms, enable_data_file_path_column, - false, - false, + PartitionedFileMetadataColumns::default(), manifest_path, )?; @@ -1170,8 +1212,12 @@ async fn table_scan( &x.1, last_updated_ms, include_data_file_path_column, - enable_row_lineage, - has_position_deletes || enable_last_updated_sequence_number_column, + PartitionedFileMetadataColumns { + first_row_id: enable_row_lineage, + sequence_number: has_position_deletes + || enable_last_updated_sequence_number_column, + change_status: scan_changed_data_files, + }, manifest_path, ) }) @@ -1246,8 +1292,12 @@ async fn table_scan( &entry, last_updated_ms, include_data_file_path_column, - enable_row_lineage, - has_position_deletes || enable_last_updated_sequence_number_column, + PartitionedFileMetadataColumns { + first_row_id: enable_row_lineage, + sequence_number: has_position_deletes + || enable_last_updated_sequence_number_column, + change_status: scan_changed_data_files, + }, manifest_path, )?; if is_attested { @@ -1327,6 +1377,55 @@ async fn table_scan( } } +fn scan_manifest_entry(entry: &ManifestEntry, scan_changed_data_files: bool) -> bool { + scan_manifest_status( + *entry.status(), + entry.data_file().content(), + scan_changed_data_files, + ) +} + +fn scan_manifest_status(status: Status, content: &Content, scan_changed_data_files: bool) -> bool { + if scan_changed_data_files { + content == &Content::Data && matches!(status, Status::Added | Status::Deleted) + } else { + status != Status::Deleted + } +} + +#[cfg(test)] +mod changed_scan_tests { + use super::{scan_manifest_status, Content, Status}; + + #[test] + fn changed_scan_includes_only_added_and_deleted_data_files() { + assert!(scan_manifest_status(Status::Added, &Content::Data, true)); + assert!(scan_manifest_status(Status::Deleted, &Content::Data, true)); + assert!(!scan_manifest_status( + Status::Existing, + &Content::Data, + true + )); + assert!(!scan_manifest_status( + Status::Added, + &Content::PositionDeletes, + true + )); + + assert!(scan_manifest_status(Status::Added, &Content::Data, false)); + assert!(scan_manifest_status( + Status::Existing, + &Content::Data, + false + )); + assert!(!scan_manifest_status( + Status::Deleted, + &Content::Data, + false + )); + } +} + fn row_lineage_projection( output_schema: &SchemaRef, scan_schema: &SchemaRef, @@ -1595,8 +1694,7 @@ fn generate_partitioned_file( manifest: &ManifestEntry, last_updated_ms: i64, enable_data_file_path: bool, - include_first_row_id: bool, - include_sequence_number: bool, + metadata_columns: PartitionedFileMetadataColumns, manifest_file_path: Option, ) -> Result { let manifest_statistics = manifest_statistics(schema, manifest); @@ -1621,11 +1719,11 @@ fn generate_partitioned_file( partition_values.push(ScalarValue::Utf8(Some(manifest_file_path))); } - if include_first_row_id { + if metadata_columns.first_row_id { partition_values.push(ScalarValue::Int64(*manifest.data_file().first_row_id())); } - if include_sequence_number { + if metadata_columns.sequence_number { let sequence_number = manifest .sequence_number() .as_ref() @@ -1639,6 +1737,15 @@ fn generate_partitioned_file( partition_values.push(ScalarValue::Int64(Some(sequence_number))); } + if metadata_columns.change_status { + let status = match manifest.status() { + Status::Added => "ADDED", + Status::Existing => "EXISTING", + Status::Deleted => "DELETED", + }; + partition_values.push(ScalarValue::Utf8(Some(status.to_owned()))); + } + let object_meta = ObjectMeta { location: util::strip_prefix(manifest.data_file().file_path()).into(), size: *manifest.data_file().file_size_in_bytes() as u64, @@ -2168,7 +2275,9 @@ struct PartitionDeleteFileIndex { mod tests { use datafusion::{ - arrow::array::Int64Array, execution::object_store::ObjectStoreUrl, prelude::SessionContext, + arrow::array::{Int64Array, StringArray}, + execution::object_store::ObjectStoreUrl, + prelude::SessionContext, scalar::ScalarValue, }; use iceberg_rust::{ @@ -2328,6 +2437,13 @@ mod tests { }; table.clone() }; + let first_snapshot_id = updated_table + .metadata() + .current_snapshot(None) + .expect("read current snapshot") + .expect("first snapshot") + .snapshot_id() + .to_owned(); let config = super::DataFusionTableConfigBuilder::default() .enable_data_file_path_column(false) .enable_data_file_row_position_column(false) @@ -2371,6 +2487,84 @@ mod tests { assert_eq!(values(0), vec![Some(20), Some(30)]); assert_eq!(values(1), vec![Some(1), Some(2)]); assert_eq!(values(2), vec![Some(1), Some(1)]); + + ctx.deregister_table("lineage_numbers") + .expect("deregister lineage table"); + ctx.register_table("lineage_numbers", writable.clone()) + .expect("restore writable table"); + ctx.sql("INSERT INTO lineage_numbers (id) VALUES (40), (50)") + .await + .expect("plan second v3 append") + .collect() + .await + .expect("append more v3 rows"); + let (appended_table, second_snapshot_id) = { + let tabular = writable.tabular.read().expect("read appended table"); + let Tabular::Table(table) = &*tabular else { + panic!("expected an Iceberg table"); + }; + let snapshot_id = table + .metadata() + .current_snapshot(None) + .expect("read second snapshot") + .expect("second snapshot") + .snapshot_id() + .to_owned(); + (table.clone(), snapshot_id) + }; + let changes_config = super::DataFusionTableConfigBuilder::default() + .enable_data_file_path_column(false) + .enable_data_file_row_position_column(false) + .enable_manifest_file_path_column(false) + .enable_row_id_column(true) + .enable_last_updated_sequence_number_column(true) + .scan_changed_data_files(true) + .build() + .expect("build change file scan config"); + let changes = Arc::new(DataFusionTable::new_with_config( + Tabular::Table(appended_table), + Some(first_snapshot_id), + Some(second_snapshot_id), + None, + Some(changes_config), + )); + ctx.deregister_table("lineage_numbers") + .expect("deregister writable table"); + ctx.register_table("lineage_numbers", changes) + .expect("register changed-file table"); + + let batches = ctx + .sql( + "SELECT id, _row_id, _last_updated_sequence_number, \ + __iceberg_change_file_status \ + FROM lineage_numbers ORDER BY __iceberg_change_file_status, id", + ) + .await + .expect("plan changed-file scan") + .collect() + .await + .expect("scan changed files"); + let batch = batches.first().expect("changed-file result batch"); + let int_values = |column: usize| { + batch + .column(column) + .as_any() + .downcast_ref::() + .expect("Int64 change result") + .iter() + .collect::>() + }; + let statuses = batch + .column(3) + .as_any() + .downcast_ref::() + .expect("Utf8 change status") + .iter() + .collect::>(); + assert_eq!(int_values(0), vec![Some(40), Some(50)]); + assert_eq!(int_values(1), vec![Some(3), Some(4)]); + assert_eq!(int_values(2), vec![Some(2), Some(2)]); + assert_eq!(statuses, vec![Some("ADDED"), Some("ADDED")]); } #[tokio::test] From 4cdb55a8bff2a419f9c98f496ded7f572228319f Mon Sep 17 00:00:00 2001 From: osipovartem Date: Tue, 15 Sep 2026 16:14:16 +0300 Subject: [PATCH 2/2] Preserve changes across append snapshot ranges --- datafusion_iceberg/src/table/mod.rs | 100 ++++++++++++++++++++++++++-- iceberg-rust/src/table/mod.rs | 4 ++ 2 files changed, 97 insertions(+), 7 deletions(-) diff --git a/datafusion_iceberg/src/table/mod.rs b/datafusion_iceberg/src/table/mod.rs index c331fe23..0504be88 100644 --- a/datafusion_iceberg/src/table/mod.rs +++ b/datafusion_iceberg/src/table/mod.rs @@ -632,6 +632,9 @@ async fn table_scan( } else { sequence_number_range }; + let change_snapshot_ids = scan_changed_data_files + .then(|| snapshot_ids_in_range(table.metadata(), snapshot_range.0, snapshot_range.1)) + .transpose()?; // If there is a filter expression the manifests to read are pruned based on the pruning statistics available in the manifest_list file. // Row lineage values may be synthesized from manifest metadata, so they must not be @@ -760,6 +763,7 @@ async fn table_scan( .await .map_err(DataFusionIcebergError::from)? }; + let data_files = prepare_changed_data_files(data_files, change_snapshot_ids.as_ref()); let pruning_predicate = PruningPredicateBuilder::new() .with_file_schema(arrow_schema.clone()) @@ -801,6 +805,7 @@ async fn table_scan( .try_collect() .await .map_err(DataFusionIcebergError::from)?; + let data_files = prepare_changed_data_files(data_files, change_snapshot_ids.as_ref()); let mut statistics = statistics_from_datafiles(&schema, &data_files); for _ in 0..usize::from(enable_row_id_column) @@ -1385,6 +1390,67 @@ fn scan_manifest_entry(entry: &ManifestEntry, scan_changed_data_files: bool) -> ) } +fn snapshot_ids_in_range( + metadata: &TableMetadata, + start_snapshot_id: Option, + end_snapshot_id: Option, +) -> Result, DataFusionError> { + let mut current = match end_snapshot_id { + Some(snapshot_id) => Some(snapshot_id), + None => metadata + .current_snapshot(None) + .map_err(DataFusionIcebergError::from)? + .map(|snapshot| *snapshot.snapshot_id()), + } + .ok_or_else(|| DataFusionError::Plan("Iceberg table has no snapshot".to_owned()))?; + let mut snapshots = HashSet::new(); + loop { + if Some(current) == start_snapshot_id { + return Ok(snapshots); + } + snapshots.insert(current); + let snapshot = metadata.snapshots.get(¤t).ok_or_else(|| { + DataFusionError::Plan(format!( + "Iceberg snapshot {current} is missing from metadata" + )) + })?; + match snapshot.parent_snapshot_id() { + Some(parent) => current = *parent, + None if start_snapshot_id.is_none() => return Ok(snapshots), + None => { + return plan_err!( + "Iceberg scan start snapshot is not an ancestor of the end snapshot" + ); + } + } + } +} + +fn prepare_changed_data_files( + data_files: Vec<(ManifestPath, ManifestEntry)>, + change_snapshot_ids: Option<&HashSet>, +) -> Vec<(ManifestPath, ManifestEntry)> { + let Some(change_snapshot_ids) = change_snapshot_ids else { + return data_files; + }; + data_files + .into_iter() + .filter_map(|(path, mut entry)| { + let is_changed_data_file = entry.data_file().content() == &Content::Data + && entry + .snapshot_id() + .is_some_and(|snapshot_id| change_snapshot_ids.contains(&snapshot_id)); + if !is_changed_data_file { + return None; + } + if *entry.status() == Status::Existing { + *entry.status_mut() = Status::Added; + } + Some((path, entry)) + }) + .collect() +} + fn scan_manifest_status(status: Status, content: &Content, scan_changed_data_files: bool) -> bool { if scan_changed_data_files { content == &Content::Data && matches!(status, Status::Added | Status::Deleted) @@ -2498,20 +2564,40 @@ mod tests { .collect() .await .expect("append more v3 rows"); - let (appended_table, second_snapshot_id) = { + let second_snapshot_id = { let tabular = writable.tabular.read().expect("read appended table"); let Tabular::Table(table) = &*tabular else { panic!("expected an Iceberg table"); }; - let snapshot_id = table + table .metadata() .current_snapshot(None) .expect("read second snapshot") .expect("second snapshot") .snapshot_id() + .to_owned() + }; + ctx.sql("INSERT INTO lineage_numbers (id) VALUES (60)") + .await + .expect("plan third v3 append") + .collect() + .await + .expect("append final v3 row"); + let (appended_table, third_snapshot_id) = { + let tabular = writable.tabular.read().expect("read final table"); + let Tabular::Table(table) = &*tabular else { + panic!("expected an Iceberg table"); + }; + let snapshot_id = table + .metadata() + .current_snapshot(None) + .expect("read third snapshot") + .expect("third snapshot") + .snapshot_id() .to_owned(); (table.clone(), snapshot_id) }; + assert_ne!(second_snapshot_id, third_snapshot_id); let changes_config = super::DataFusionTableConfigBuilder::default() .enable_data_file_path_column(false) .enable_data_file_row_position_column(false) @@ -2524,7 +2610,7 @@ mod tests { let changes = Arc::new(DataFusionTable::new_with_config( Tabular::Table(appended_table), Some(first_snapshot_id), - Some(second_snapshot_id), + Some(third_snapshot_id), None, Some(changes_config), )); @@ -2561,10 +2647,10 @@ mod tests { .expect("Utf8 change status") .iter() .collect::>(); - assert_eq!(int_values(0), vec![Some(40), Some(50)]); - assert_eq!(int_values(1), vec![Some(3), Some(4)]); - assert_eq!(int_values(2), vec![Some(2), Some(2)]); - assert_eq!(statuses, vec![Some("ADDED"), Some("ADDED")]); + assert_eq!(int_values(0), vec![Some(40), Some(50), Some(60)]); + assert_eq!(int_values(1), vec![Some(3), Some(4), Some(5)]); + assert_eq!(int_values(2), vec![Some(2), Some(2), Some(3)]); + assert_eq!(statuses, vec![Some("ADDED"), Some("ADDED"), Some("ADDED")]); } #[tokio::test] diff --git a/iceberg-rust/src/table/mod.rs b/iceberg-rust/src/table/mod.rs index d50164eb..cc3a59b7 100644 --- a/iceberg-rust/src/table/mod.rs +++ b/iceberg-rust/src/table/mod.rs @@ -327,6 +327,7 @@ async fn datafiles( let object_store = object_store.clone(); let manifest_path = file.manifest_path.clone(); let manifest_sequence_number = file.sequence_number; + let manifest_snapshot_id = file.added_snapshot_id; let manifest_first_row_id = file.first_row_id; let manifest_content = file.content; async move { @@ -358,6 +359,9 @@ async fn datafiles( entries .into_iter() .filter_map(|mut x| { + if x.snapshot_id().is_none() { + *x.snapshot_id_mut() = Some(manifest_snapshot_id); + } let sequence_number = if let Some(sequence_number) = x.sequence_number() { *sequence_number