Skip to content
Open
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
165 changes: 133 additions & 32 deletions crates/datafusion/src/physical_plan/commit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ use datafusion::arrow::datatypes::{
DataType, Field, Schema as ArrowSchema, SchemaRef as ArrowSchemaRef,
};
use datafusion::common::{
internal_datafusion_err, internal_err, tree_node::TreeNodeRecursion,
exec_err, internal_datafusion_err, internal_err, tree_node::TreeNodeRecursion,
};
use datafusion::error::{DataFusionError, Result};
use datafusion::execution::{SendableRecordBatchStream, TaskContext};
Expand All @@ -44,22 +44,30 @@ use crate::to_datafusion_error;

/// IcebergCommitExec is responsible for collecting the files written and use
/// [`Transaction::fast_append`] to commit the data files written.
///
/// Its input produces the data files to commit, in the form
/// [`IcebergWriteExec`](super::IcebergWriteExec) outputs them, in a single
/// partition. Its output is one row holding the number of rows committed.
#[derive(Debug)]
pub(crate) struct IcebergCommitExec {
pub struct IcebergCommitExec {
table: Table,
catalog: Arc<dyn Catalog>,
input: Arc<dyn ExecutionPlan>,
schema: ArrowSchemaRef,
count_schema: ArrowSchemaRef,
plan_properties: Arc<PlanProperties>,
}

impl IcebergCommitExec {
/// Commits the data files `input` produces to `table` through `catalog`.
///
/// `input` must have a single partition, such as a
/// [`CoalescePartitionsExec`](datafusion::physical_plan::coalesce_partitions::CoalescePartitionsExec)
/// over an [`IcebergWriteExec`](super::IcebergWriteExec); executing the node
/// fails otherwise.
pub fn new(
table: Table,
catalog: Arc<dyn Catalog>,
input: Arc<dyn ExecutionPlan>,
schema: ArrowSchemaRef,
) -> Self {
let count_schema = Self::make_count_schema();

Expand All @@ -69,12 +77,21 @@ impl IcebergCommitExec {
table,
catalog,
input,
schema,
count_schema,
plan_properties,
}
}

/// The catalog this node commits through.
pub fn catalog(&self) -> &Arc<dyn Catalog> {
&self.catalog
}

/// The table this node commits to, as loaded when the node was planned.
pub fn table(&self) -> &Table {
&self.table
}

// Compute the plan properties for this execution plan
fn compute_properties(schema: ArrowSchemaRef) -> Arc<PlanProperties> {
Arc::new(PlanProperties::new(
Expand Down Expand Up @@ -115,11 +132,22 @@ impl DisplayAs for IcebergCommitExec {
write!(f, "IcebergCommitExec: table={}", self.table.identifier())
}
DisplayFormatType::Verbose => {
// The fields on one line, as `Schema`'s own display spans
// several.
let fields = self
.table
.metadata()
.current_schema()
.as_struct()
.fields()
.iter()
.map(|field| field.to_string().trim_end().to_string())
.collect::<Vec<_>>()
.join(", ");
write!(
f,
"IcebergCommitExec: table={}, schema={:?}",
self.table.identifier(),
self.schema
"IcebergCommitExec: table={}, schema=[{fields}]",
self.table.identifier()
)
}
DisplayFormatType::TreeRender => {
Expand Down Expand Up @@ -177,7 +205,6 @@ impl ExecutionPlan for IcebergCommitExec {
self.table.clone(),
self.catalog.clone(),
children[0].clone(),
self.schema.clone(),
)))
}

Expand All @@ -193,6 +220,16 @@ impl ExecutionPlan for IcebergCommitExec {
);
}

// Only partition 0 of the input is read below, so the files of any
// other partition would silently go uncommitted.
let input_partitions = self.input.properties().partitioning.partition_count();
if input_partitions != 1 {
return exec_err!(
"IcebergCommitExec requires an input with one partition, but it has \
{input_partitions}; coalesce the input first"
);
}

let table = self.table.clone();
let input_plan = self.input.clone();

Expand Down Expand Up @@ -297,6 +334,7 @@ mod tests {
use datafusion::physical_plan::common::collect;
use datafusion::physical_plan::execution_plan::Boundedness;
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::union::UnionExec;
use datafusion::physical_plan::{
DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties,
};
Expand Down Expand Up @@ -487,19 +525,8 @@ mod tests {
let input_exec =
Arc::new(MockWriteExec::new(vec![data_file1_json, data_file2_json]));

// Create the IcebergCommitExec
let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
DATA_FILES_COL_NAME,
DataType::Utf8,
false,
)]));

let commit_exec = IcebergCommitExec::new(
table.clone(),
catalog.clone(),
input_exec,
arrow_schema,
);
let commit_exec =
IcebergCommitExec::new(table.clone(), catalog.clone(), input_exec);

// Verify Execution Plan schema matches the count schema
assert_eq!(commit_exec.schema(), IcebergCommitExec::make_count_schema());
Expand Down Expand Up @@ -559,6 +586,76 @@ mod tests {
Ok(())
}

/// The commit reads a single input partition, so an input with more is
/// refused rather than committing the files of its first partition alone.
#[tokio::test]
async fn test_iceberg_commit_exec_rejects_multiple_input_partitions()
-> Result<(), Box<dyn std::error::Error>> {
let catalog = Arc::new(
MemoryCatalogBuilder::default()
.load(
"memory",
HashMap::from([(
MEMORY_CATALOG_WAREHOUSE.to_string(),
"memory://root".to_string(),
)]),
)
.await?,
);
let namespace = NamespaceIdent::new("test_namespace".to_string());
catalog.create_namespace(&namespace, HashMap::new()).await?;
let schema = Schema::builder()
.with_fields(vec![
NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int))
.into(),
])
.build()?;
let table_creation = TableCreation::builder()
.name("test_table".to_string())
.schema(schema)
.location("memory://root/test_table".to_string())
.build();
let table = catalog.create_table(&namespace, table_creation).await?;

// One data file in each of two input partitions.
let partition_type = table.metadata().default_partition_type().clone();
let mut partitions: Vec<Arc<dyn ExecutionPlan>> = vec![];
for path in ["path/to/file1.parquet", "path/to/file2.parquet"] {
let data_file = DataFileBuilder::default()
.content(DataContentType::Data)
.file_path(path.to_string())
.file_format(DataFileFormat::Parquet)
.file_size_in_bytes(1024)
.record_count(100)
.partition_spec_id(table.metadata().default_partition_spec_id())
.partition(Struct::empty())
.build()?;
let json = iceberg::spec::serialize_data_file_to_json(
data_file,
&partition_type,
table.metadata().format_version(),
)?;
partitions.push(Arc::new(MockWriteExec::new(vec![json])));
}
let input = UnionExec::try_new(partitions)?;
assert_eq!(input.properties().partitioning.partition_count(), 2);
let commit_exec = IcebergCommitExec::new(table.clone(), catalog.clone(), input);
let err = match commit_exec.execute(0, Arc::new(TaskContext::default())) {
Ok(_) => panic!("a commit over two input partitions must not execute"),
Err(err) => err.to_string(),
};
assert_eq!(
err,
"Execution error: IcebergCommitExec requires an input with one \
partition, but it has 2; coalesce the input first"
);

let table = catalog.load_table(table.identifier()).await?;
assert!(table.metadata().current_snapshot().is_none());

Ok(())
}

#[tokio::test]
async fn test_iceberg_commit_exec_empty_insert()
-> Result<(), Box<dyn std::error::Error>> {
Expand Down Expand Up @@ -598,16 +695,20 @@ mod tests {

// Mock write plan produces no data files
let input_exec = Arc::new(MockWriteExec::new(vec![]));
let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
DATA_FILES_COL_NAME,
DataType::Utf8,
false,
)]));
let commit_exec = IcebergCommitExec::new(
table.clone(),
catalog.clone(),
input_exec,
arrow_schema,
let commit_exec =
IcebergCommitExec::new(table.clone(), catalog.clone(), input_exec);

// The verbose display shows the table's schema on one line.
struct Verbose<'a>(&'a IcebergCommitExec);
impl fmt::Display for Verbose<'_> {
fn fmt(&self, f: &mut Formatter) -> fmt::Result {
self.0.fmt_as(DisplayFormatType::Verbose, f)
}
}
assert_eq!(
Verbose(&commit_exec).to_string(),
"IcebergCommitExec: table=test_empty_insert.empty_insert_table, \
schema=[1: id: required int]"
);

let task_ctx = Arc::new(TaskContext::default());
Expand Down
7 changes: 7 additions & 0 deletions crates/datafusion/src/physical_plan/metadata_scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,15 @@ use futures::TryStreamExt;

use crate::metadata_table::IcebergMetadataTableProvider;

/// Scans an Iceberg metadata table, such as `$snapshots`.
#[derive(Debug)]
pub struct IcebergMetadataScan {
provider: IcebergMetadataTableProvider,
properties: Arc<PlanProperties>,
}

impl IcebergMetadataScan {
/// Creates a scan of the metadata table `provider` reads.
pub fn new(provider: IcebergMetadataTableProvider) -> Self {
let properties = Arc::new(PlanProperties::new(
EquivalenceProperties::new(provider.schema()),
Expand All @@ -46,6 +48,11 @@ impl IcebergMetadataScan {
properties,
}
}

/// The provider this node scans.
pub fn provider(&self) -> &IcebergMetadataTableProvider {
&self.provider
}
}

impl DisplayAs for IcebergMetadataScan {
Expand Down
3 changes: 3 additions & 0 deletions crates/datafusion/src/physical_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,9 @@ pub(crate) mod write;

pub(crate) const DATA_FILES_COL_NAME: &str = "data_files";

pub use commit::IcebergCommitExec;
pub use expr_to_predicate::convert_filters_to_predicate;
pub use metadata_scan::IcebergMetadataScan;
pub use project::{PartitionExpr, project_with_partition};
pub use scan::IcebergTableScan;
pub use write::IcebergWriteExec;
Loading
Loading