mbutrovich commented on code in PR #55: URL: https://github.com/apache/datafusion-iceberg/pull/55#discussion_r4235376704
########## crates/datafusion/src/config.rs: ########## @@ -0,0 +1,66 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Session options for Iceberg tables, set with `SET iceberg.<group>.<option>`. + +use datafusion::common::config::ConfigExtension; +use datafusion::common::extensions_options; + +extensions_options! { + /// Session options for reading Iceberg tables through DataFusion. + /// + /// `SET iceberg.<group>.<option> = <value>` only works once this extension + /// is registered on the session, with + /// [`SessionConfig::with_option_extension`](datafusion::prelude::SessionConfig::with_option_extension). + /// Without it, every option keeps its default. DataFusion lists these + /// options in `information_schema.df_settings` by their last name alone, + /// such as `max_merge_files`, and `SHOW` and `RESET` do not support them. + pub struct IcebergDataFusionConfig { + /// Options for planning table scans. + pub planning: IcebergPlanningConfig, default = IcebergPlanningConfig::default() + } +} + +extensions_options! { + /// Options for planning Iceberg table scans, under `iceberg.planning`. + pub struct IcebergPlanningConfig { + /// When true, a scan lists its data files while it is planned and, if + /// every file records the same sort order, reports that order to + /// DataFusion so that it can skip sorting the scan's output. The scan + /// then merges the files in that order, with every file open at once. + /// + /// Only the leading sort fields that are identity transforms of + /// projected top-level columns are reported, stopping at the first + /// floating point or UUID field, whose orders in Iceberg writers and in + /// DataFusion can differ. + /// + /// Iceberg sorts ascending fields with nulls first by default, while + /// DataFusion's `ORDER BY c` puts them last. On a nullable sort column, + /// only `ORDER BY c NULLS FIRST`, or a session with + /// `datafusion.sql_parser.default_null_ordering = 'nulls_min'`, can + /// skip its sort. Other queries sort the merged output again. Review Comment: > I'd rather look at DataFusion's sort pushdown as a follow-up than grow this PR if that works for you. Could you file an issue for the follow-up and link it from the PR description? The issue should cover more than sort pushdown. DataFusion 55's `PushdownSort` rule offers a `SortExec`'s order to the plan below it through `ExecutionPlan::try_pushdown_sort`, and drops the sort when the input answers `Exact`. That would cover `ORDER BY`, where the plan is `SortExec` over `CooperativeExec` over the scan. It wouldn't cover the sort-merge join that `test_sorted_scan_removes_sort_from_sort_merge_join` checks. With no ordering reported, each join input plans as `SortExec(preserve_partitioning)` over `RepartitionExec: partitioning=Hash([id@0], 4)`, and [`RepartitionExec::try_pushdown_sort`](https://github.com/apache/datafusion/blob/7d3835c71f30cbd3c3ae4041732267f1f453097a/datafusion/physical-plan/src/repartition/mod.rs#L1907-L1915) returns `Unsupported` when the repartition doesn't keep its input's order. A design that merges only on pushdown would lose the join case, so the issue should re cord that. ########## crates/datafusion/src/config.rs: ########## @@ -0,0 +1,66 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Session options for Iceberg tables, set with `SET iceberg.<group>.<option>`. + +use datafusion::common::config::ConfigExtension; +use datafusion::common::extensions_options; + +extensions_options! { + /// Session options for reading Iceberg tables through DataFusion. + /// + /// `SET iceberg.<group>.<option> = <value>` only works once this extension + /// is registered on the session, with + /// [`SessionConfig::with_option_extension`](datafusion::prelude::SessionConfig::with_option_extension). + /// Without it, every option keeps its default. DataFusion lists these + /// options in `information_schema.df_settings` by their last name alone, + /// such as `max_merge_files`, and `SHOW` and `RESET` do not support them. + pub struct IcebergDataFusionConfig { + /// Options for planning table scans. + pub planning: IcebergPlanningConfig, default = IcebergPlanningConfig::default() + } +} + +extensions_options! { + /// Options for planning Iceberg table scans, under `iceberg.planning`. + pub struct IcebergPlanningConfig { + /// When true, a scan lists its data files while it is planned and, if + /// every file records the same sort order, reports that order to + /// DataFusion so that it can skip sorting the scan's output. The scan + /// then merges the files in that order, with every file open at once. + /// + /// Only the leading sort fields that are identity transforms of + /// projected top-level columns are reported, stopping at the first + /// floating point or UUID field, whose orders in Iceberg writers and in + /// DataFusion can differ. + /// + /// Iceberg sorts ascending fields with nulls first by default, while + /// DataFusion's `ORDER BY c` puts them last. On a nullable sort column, + /// only `ORDER BY c NULLS FIRST`, or a session with + /// `datafusion.sql_parser.default_null_ordering = 'nulls_min'`, can + /// skip its sort. Other queries sort the merged output again. + pub preserve_data_ordering: bool, default = false + /// The most data files a scan merges to preserve their sort order. + /// A scan of more files reports no order, and DataFusion sorts its + /// output instead, which can spill to disk where the merge cannot. + pub max_merge_files: usize, default = 64 Review Comment: > Bounding it instead by how many files' key ranges overlap would need per-file column bounds on `FileScanTask`, which is follow-up work. Could you file an issue for this one too? It needs an iceberg-rust change first, so the issue can track both sides. ########## crates/datafusion/src/physical_plan/scan.rs: ########## @@ -323,12 +540,298 @@ async fn get_batch_stream( if let Some(pred) = predicates { scan_builder = scan_builder.with_filter(pred); } - let table_scan = scan_builder.build().map_err(to_datafusion_error)?; + scan_builder.build().map_err(to_datafusion_error) +} - let stream = table_scan - .to_arrow() - .await - .map_err(to_datafusion_error)? - .map_err(to_datafusion_error); - Ok(Box::pin(stream)) +/// The order of the columns of `output` that every one of `tasks` is sorted +/// in, or `None` if they record different sort orders, one records none, or +/// the schema they are read with does not match `output`. +fn shared_ordering(tasks: &[FileScanTask], output: &ArrowSchema) -> Option<LexOrdering> { + let first = tasks.first()?; + let order = first.sort_order()?; + if !tasks.iter().all(|task| { + task.sort_order_id() == first.sort_order_id() && task.sort_order().is_some() + }) { + return None; + } + if !reads_as_declared(first.schema(), output) { + return None; + } + lex_ordering(order, first.schema(), output) +} + +/// Whether reading `schema`, the snapshot schema the tasks are read with, +/// yields every column of `output` with the type it declares, and no nulls +/// where it declares none. A provider's schema can be older than the +/// snapshot's, and DataFusion ignores the null order of a non-nullable +/// column, while the merge requires its input to match `output`. +fn reads_as_declared(schema: &Schema, output: &ArrowSchema) -> bool { + let Ok(read) = schema_to_arrow_schema(schema) else { + return false; + }; + output.fields().iter().all(|field| { + read.field_with_name(field.name()).is_ok_and(|read| { + read.data_type() == field.data_type() + && (field.is_nullable() || !read.is_nullable()) + }) + }) +} + +/// The longest leading part of `order` that DataFusion can merge and report as +/// an order of the columns of `output`, or `None` if that is empty. `schema` +/// is the Iceberg schema the scan reads. +fn lex_ordering( + order: &SortOrder, + schema: &Schema, + output: &ArrowSchema, +) -> Option<LexOrdering> { + let sort_exprs = order.fields.iter().map_while(|field| { + if field.transform != Transform::Identity { + return None; + } + let source = schema.as_struct().field_by_id(field.source_id)?; + if !orders_like_writers(&source.field_type) { + return None; + } + let index = output.index_of(&source.name).ok()?; + Some(PhysicalSortExpr::new( + Arc::new(Column::new(&source.name, index)), + SortOptions { + descending: field.direction == SortDirection::Descending, + nulls_first: field.null_order == NullOrder::First, + }, + )) + }); + LexOrdering::new(sort_exprs) +} + +/// Whether DataFusion orders values of `field_type` as the engines that write +/// sorted Iceberg files do. It orders -0.0 before 0.0 and negative NaN before +/// every number, where Spark treats the zeros as equal and NaN as the largest +/// value, and the Iceberg spec leaves the order of UUIDs undefined +/// (apache/iceberg#14216). Review Comment: As I read the spec's [sort order section](https://github.com/apache/iceberg/blob/79bc69a51da6d95ccd6da6d46ef156c74e30b8cf/format/spec.md?plain=1#L654), `-NaN < -Infinity < -value < -0 < 0 < value < Infinity < NaN` is the order writers should produce, and Arrow sorts floats in that order. So the files this guards against come from writers like Spark that don't follow the spec here. Should the comment say so? As written, it reads as though DataFusion's order is the one that differs. Someone deciding later whether to lift the guard would need to know that it depends on the writers. ```suggestion /// Whether DataFusion orders values of `field_type` as the engines that write /// sorted Iceberg files do. For floating point, DataFusion uses the order the /// spec defines (<https://iceberg.apache.org/spec/#sort-orders>), with -0.0 /// before 0.0 and negative NaN first. Spark treats the zeros as equal and /// every NaN as the largest value, so files it writes may not be sorted in the /// spec's order. The spec leaves the order of UUIDs undefined /// (apache/iceberg#14216). ``` ########## crates/datafusion/tests/integration_datafusion_test.rs: ########## @@ -1384,3 +1409,750 @@ async fn test_scan_after_schema_evolution_reads_provider_columns() Ok(()) } + +/// The id of the sort order of the tables `create_sorted_table` creates. +const SORTED: Option<i32> = Some(1); + +/// Creates table `name`, sorted by `id` ascending, with columns `id`, `data`, +/// a struct `info` and a list `tags`, whose values are derived from `id`. Each +/// entry of `files` becomes one data file of the given ids, in the given order, +/// which records the given sort order id. +async fn create_sorted_table( + catalog: &Arc<dyn Catalog>, + namespace: &NamespaceIdent, + name: &str, + files: &[(Option<i32>, &[i32])], +) -> Result<Table, Box<dyn Error>> { + let sort_order = SortOrder::builder() + .with_sort_field( + SortField::builder() + .source_id(1) + .transform(Transform::Identity) + .direction(SortDirection::Ascending) + .null_order(NullOrder::First) + .build(), + ) + .build_unbound()?; + let creation = TableCreation::builder() + .location(temp_path()) + .name(name.to_string()) + .schema( + Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)) + .into(), + NestedField::required( + 2, + "data", + Type::Primitive(PrimitiveType::String), + ) + .into(), + NestedField::optional( + 3, + "info", + Type::Struct(StructType::new(vec![ + NestedField::optional( + 4, + "tag", + Type::Primitive(PrimitiveType::String), + ) + .into(), + ])), + ) + .into(), + NestedField::optional( + 5, + "tags", + Type::List(ListType::new( + NestedField::list_element( + 6, + Type::Primitive(PrimitiveType::Int), + false, + ) + .into(), + )), + ) + .into(), + ]) + .build()?, + ) + .sort_order(sort_order) + .build(); + let table = catalog.create_table(namespace, creation).await?; + assert_eq!( + Some(table.metadata().default_sort_order_id() as i32), + SORTED + ); + let files: Vec<_> = files + .iter() + .map(|(sort_order_id, ids)| { + (*sort_order_id, ids.iter().copied().map(Some).collect()) + }) + .collect(); + append_sorted_files(catalog, &table, &files).await +} + +/// Commits one data file of `table` per entry of `files`, with the given ids, +/// in the given order, recording the given sort order id. The other columns' +/// values are derived from `id`. +async fn append_sorted_files( + catalog: &Arc<dyn Catalog>, + table: &Table, + files: &[(Option<i32>, Vec<Option<i32>>)], +) -> Result<Table, Box<dyn Error>> { + if files.is_empty() { + return Ok(table.clone()); + } + let arrow_schema = + Arc::new(schema_to_arrow_schema(table.metadata().current_schema())?); + let (DataType::Struct(info_fields), DataType::List(tags_field)) = ( + arrow_schema.field(2).data_type().clone(), + arrow_schema.field(3).data_type().clone(), + ) else { + unreachable!("info is a struct and tags a list"); + }; + let data_dir = format!("{}/data", table.metadata().location()); + std::fs::create_dir_all(&data_dir)?; + let snapshots = table.metadata().snapshots().count(); + let mut data_files = Vec::new(); + for (i, (sort_order_id, ids)) in files.iter().enumerate() { + let path = format!("{data_dir}/{snapshots}-{i}.parquet"); + let label = |prefix: &str, id: &Option<i32>| match id { + Some(id) => format!("{prefix} {id}"), + None => format!("{prefix} null"), + }; + let batch = RecordBatch::try_new( + arrow_schema.clone(), + vec![ + Arc::new(Int32Array::from(ids.clone())), + Arc::new(StringArray::from_iter_values( + ids.iter().map(|id| label("row", id)), + )), + Arc::new(StructArray::new( + info_fields.clone(), + vec![Arc::new(StringArray::from_iter_values( + ids.iter().map(|id| label("tag", id)), + ))], + None, + )), + Arc::new(ListArray::new( + tags_field.clone(), + OffsetBuffer::from_lengths(ids.iter().map(|_| 2)), + Arc::new(Int32Array::from_iter_values(ids.iter().flat_map(|id| { + let id = id.unwrap_or_default(); + [id, -id] + }))), + None, + )), + ], + )?; + let mut writer = ArrowWriter::try_new( + std::fs::File::create(&path)?, + arrow_schema.clone(), + None, + )?; + writer.write(&batch)?; + writer.close()?; + + let mut builder = DataFileBuilder::default(); + builder + .content(DataContentType::Data) + .file_path(path.clone()) + .file_format(DataFileFormat::Parquet) + .file_size_in_bytes(std::fs::metadata(&path)?.len()) + .record_count(ids.len() as u64) + .partition_spec_id(0) + .partition(Struct::empty()); + if let Some(sort_order_id) = sort_order_id { + builder.sort_order_id(*sort_order_id); + } + data_files.push(builder.build()?); + } + let tx = Transaction::new(table); + let tx = tx.fast_append().add_data_files(data_files).apply(tx)?; + Ok(tx.commit(catalog.as_ref()).await?) +} + +/// Stands in for another engine committing, without writing data, the current +/// schema of `table` with the field of the same id as `field` replaced by it. +/// iceberg-rust has no action that makes a column optional or promotes its +/// type, so this rewrites the table's current metadata file in place. +async fn replace_field_externally( + table: &Table, + field: NestedField, +) -> Result<(), Box<dyn Error>> { + let current = table.metadata().current_schema(); + let schema = Schema::builder() + .with_fields(current.as_struct().fields().iter().map(|current| { + if current.id == field.id { + Arc::new(field.clone()) + } else { + current.clone() + } + })) + .build()?; + let metadata = table + .metadata() + .clone() + .into_builder(None) + .add_current_schema(schema)? + .build()? + .metadata; + let location = MetadataLocation::from_str(table.metadata_location_result()?)?; + metadata.write_to(table.file_io(), &location).await?; + Ok(()) +} + +/// The values, or nulls, of the `id` column, which must be the first, of +/// `batches`. +fn nullable_ids(batches: &[RecordBatch]) -> Vec<Option<i32>> { + batches + .iter() + .flat_map(|batch| batch.column(0).as_primitive::<Int32Type>().iter()) + .collect() +} + +/// A session that registers the Iceberg options, with +/// `iceberg.planning.preserve_data_ordering` set to `preserve`. +async fn ordering_session( + preserve: bool, + target_partitions: usize, +) -> Result<SessionContext, Box<dyn Error>> { + let config = SessionConfig::new() + .with_target_partitions(target_partitions) + .with_option_extension(IcebergDataFusionConfig::default()); + let ctx = SessionContext::new_with_config(config); + ctx.sql(&format!( + "SET iceberg.planning.preserve_data_ordering = {preserve}" + )) + .await? + .collect() + .await?; + Ok(ctx) +} + +/// Plans `table` as a full scan in `ctx`. +async fn plan_sorted_scan( + table: &Table, + ctx: &SessionContext, +) -> Result<Arc<dyn ExecutionPlan>, Box<dyn Error>> { + let provider = IcebergStaticTableProvider::try_new_from_table(table.clone()).await?; + Ok(provider.scan(&ctx.state(), None, &[], None).await?) +} + +/// The values of the `id` column, which must be the first, of `batches`. +fn ids(batches: &[RecordBatch]) -> Vec<i32> { + batches + .iter() + .flat_map(|batch| { + batch + .column(0) + .as_primitive::<Int32Type>() + .values() + .to_vec() + }) + .collect() +} + +async fn sorted_namespace() -> Result<(Arc<dyn Catalog>, NamespaceIdent), Box<dyn Error>> +{ + let catalog = get_iceberg_catalog().await; + let namespace = NamespaceIdent::new("sorted".to_string()); + set_test_namespace(&catalog, &namespace).await?; + Ok((Arc::new(catalog), namespace)) +} + +#[tokio::test] +async fn test_sorted_scan_merges_files_with_overlapping_ranges() +-> Result<(), Box<dyn Error>> { + let (catalog, namespace) = sorted_namespace().await?; + let table = create_sorted_table( + &catalog, + &namespace, + "t", + &[ + (SORTED, &[1, 4, 7, 7]), + (SORTED, &[2, 4, 8]), + (SORTED, &[3, 5, 9]), + ], Review Comment: Each file in the sorted-scan tests holds at most four rows, so every per-file stream yields a single batch, and the merge never advances a stream past its first batch. What do you think about making one of these files longer than a reader batch, or writing it with several row groups, so that a test covers per-file streams that yield more than one batch? The order within each stream comes from `ArrowReader`, not from DataFusion's merge, and no test exercises it across batch boundaries. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
