Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions datafusion_iceberg/src/row_lineage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
167 changes: 164 additions & 3 deletions datafusion_iceberg/src/table/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -2173,7 +2174,26 @@ pub async fn write_parquet_data_files(
context: &Arc<TaskContext>,
branch: Option<&str>,
) -> Result<Vec<DataFile>, 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<TaskContext>,
branch: Option<&str>,
) -> Result<Vec<DataFile>, 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.
Expand Down Expand Up @@ -2204,7 +2224,7 @@ pub async fn write_parquet_equality_delete_files(
equality_ids: &[i32],
branch: Option<&str>,
) -> Result<Vec<DataFile>, 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(
Expand All @@ -2217,6 +2237,7 @@ async fn write_parquet_files(
context: &Arc<TaskContext>,
equality_ids: Option<&[i32]>,
branch: Option<&str>,
preserve_row_lineage: bool,
) -> Result<Vec<DataFile>, DataFusionError> {
let object_store = table.object_store();
let metadata = table.metadata();
Expand All @@ -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::<ArrowSchema>::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()
Expand Down Expand Up @@ -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::{
Expand All @@ -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;

Expand Down Expand Up @@ -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<dyn Catalog> = 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::<Int64Array>()
.expect("Int64 lineage column")
.iter()
.collect::<Vec<_>>()
};
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();
Expand Down
1 change: 1 addition & 0 deletions iceberg-rust-spec/src/spec/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
17 changes: 17 additions & 0 deletions iceberg-rust-spec/src/spec/row_lineage.rs
Original file line number Diff line number Diff line change
@@ -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
}
8 changes: 8 additions & 0 deletions iceberg-rust/src/file_format/parquet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down Expand Up @@ -116,6 +117,13 @@ pub fn parquet_to_datafile(
let mut counted_logical_values = HashSet::new();
let mut variant_null_counts = HashMap::<i32, i64>::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
Expand Down
53 changes: 41 additions & 12 deletions iceberg-rust/src/table/manifest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,36 @@ impl FilteredManifestStats {
}
}

fn inherit_first_row_id(
content: manifest_list::Content,
next_row_id: &mut Option<i64>,
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.
///
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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);
Expand All @@ -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;
Expand Down
Loading
Loading