sunchao commented on code in PR #6241: URL: https://github.com/apache/datafusion-comet/pull/6241#discussion_r4111185721
########## native/core/src/execution/operators/iceberg_dictionary.rs: ########## @@ -0,0 +1,1132 @@ +// 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. + +//! Which columns of a native Iceberg data file are dictionary-encoded. +//! +//! iceberg-java writes Parquet through parquet-mr, which judges dictionary encoding on the first +//! data page of every column chunk (`FallbackValuesWriter`): unless the page's dictionary-encoded +//! values plus its dictionary come out smaller than the page's plain encoding, the chunk is written +//! plain from its first page on and has no dictionary page. parquet-rs never makes that +//! comparison. It keeps a column dictionary-encoded until the dictionary reaches +//! `dictionary_page_size_limit`, then falls back to plain for the rest of the chunk but still +//! writes the dictionary page. A high-cardinality column therefore carries a dictionary page of up +//! to `write.parquet.dict-size-bytes` that every selective read of the chunk has to fetch, where +//! iceberg-java writes none (apache/datafusion-comet#6114). +//! +//! parquet-rs cannot change a column's encoding once a file is open, so the native writer takes +//! parquet-mr's decision before opening one: [`DictionaryChooser`] replays parquet-mr's size +//! accounting over the rows each column's first data page would hold, and turns dictionary encoding +//! off for the columns parquet-mr would write plain. The accounting follows `FallbackValuesWriter`, +//! `DictionaryValuesWriter` and `RunLengthBitPackingHybridEncoder` in parquet-column, and the page +//! boundary follows `ColumnWriteStoreBase.sizeCheck`. + +use std::borrow::Cow; + +use arrow::array::{downcast_primitive_array, Array, AsArray, OffsetSizeTrait, RecordBatch}; +use arrow::datatypes::{DataType, Schema as ArrowSchema, ToByteSlice}; +use datafusion::common::hash_map::Entry; +use datafusion::common::HashMap; +use datafusion::error::{DataFusionError, Result as DFResult}; +use parquet::arrow::ArrowSchemaConverter; +use parquet::basic::Type as PhysicalType; +use parquet::file::properties::WriterProperties; +use parquet::schema::types::ColumnPath; + +/// The row count of parquet-mr's first page size check. Iceberg passes +/// `write.parquet.row-group-check-min-record-count` to parquet-mr as the minimum row count between +/// page size checks, and the native writer declines tables that change it from this default, so no +/// first page ends before this row even when `write.parquet.page-row-limit` is lower. +const FIRST_PAGE_SIZE_CHECK_ROWS: usize = 100; + +/// parquet-mr ends a page once less than this share of the page size is left +/// (`ColumnWriteStoreBase.THRESHOLD_TOLERANCE_RATIO`). +const PAGE_SIZE_TOLERANCE_RATIO: f32 = 0.1; + +/// Decides, from the rows a partition's data starts with, which of its columns keep dictionary +/// encoding. +/// +/// The choice is made once per partition and applies to every file and row group the task writes +/// for it, where parquet-mr decides again for every column chunk. A column's first data page is +/// the table's `write.parquet.page-row-limit` rows unless the page size ends it sooner. parquet-mr +/// only notices a full page at its periodic size checks, whose timing depends on every column, so +/// where the page size decides, this ends the page at the first row past parquet-mr's threshold and +/// can judge a column on somewhat fewer rows than parquet-mr would. +pub(super) struct DictionaryChooser { + /// The table's writer properties, which every choice starts from. + base: WriterProperties, + /// One entry per Parquet leaf column, in schema order. `None` for the columns there is nothing + /// to decide for: parquet-rs never dictionary-encodes booleans, nor fixed-length byte arrays + /// under Parquet format v1, and a column whose dictionary is already off stays off. + columns: Vec<Option<Column>>, + /// Rows the first data page of every column holds at most. + page_rows: usize, + /// Plain-encoded bytes at which parquet-mr's size check ends a page. + page_bytes: u64, + /// Memory at which a partition stops holding rows back, however few it has: the row group + /// size, since parquet-mr's first page never outlasts its first row group. + max_held_bytes: Option<usize>, +} + +struct Column { + path: ColumnPath, + /// Plain-encoded width of one value, or `None` for a byte array, which parquet-mr accounts as + /// a 4-byte length plus the bytes. + width: Option<u64>, + /// `write.parquet.dict-size-bytes`. parquet-mr abandons a dictionary as soon as it grows past + /// this, which on a first page leaves no dictionary page at all. + dictionary_limit: u64, +} + +impl DictionaryChooser { + /// `schema` is the Arrow schema every batch is written with, which fixes the Parquet leaf + /// columns the same way the parquet writer derives them. + pub(super) fn try_new(base: WriterProperties, schema: &ArrowSchema) -> DFResult<Self> { + let parquet_schema = ArrowSchemaConverter::new() + .with_coerce_types(base.coerce_types()) + .convert(schema) + .map_err(DataFusionError::from)?; + let columns = parquet_schema + .columns() + .iter() + .map(|descr| { + let path = descr.path(); + let width = match descr.physical_type() { + PhysicalType::INT32 | PhysicalType::FLOAT => Some(4), + PhysicalType::INT64 | PhysicalType::DOUBLE => Some(8), + PhysicalType::BYTE_ARRAY => None, + // INT96 is not an Iceberg type. + PhysicalType::BOOLEAN + | PhysicalType::FIXED_LEN_BYTE_ARRAY + | PhysicalType::INT96 => return None, + }; + base.dictionary_enabled(path).then(|| Column { + path: path.clone(), + width, + dictionary_limit: base.column_dictionary_page_size_limit(path) as u64, + }) + }) + .collect(); + let page_size = base.data_page_size_limit(); + let tolerance = (page_size as f32 * PAGE_SIZE_TOLERANCE_RATIO) as usize; + Ok(Self { + columns, + page_rows: base + .data_page_row_count_limit() + .max(FIRST_PAGE_SIZE_CHECK_ROWS), + page_bytes: page_size.saturating_sub(tolerance) as u64, + max_held_bytes: base.max_row_group_bytes(), + base, + }) + } + + /// The table's writer properties, before any choice. + pub(super) fn base(&self) -> &WriterProperties { + &self.base + } + + /// Whether a partition that has shown `rows` rows, held in `bytes` of memory, should hold back + /// more before choosing. + pub(super) fn wants_more(&self, rows: usize, bytes: usize) -> bool { + self.columns.iter().any(Option::is_some) && rows < self.page_rows && self.may_hold(bytes) + } + + /// Whether rows taking `bytes` of memory may stay held back. + pub(super) fn may_hold(&self, bytes: usize) -> bool { + self.max_held_bytes.is_none_or(|max| bytes < max) + } + + /// The writer properties for a partition whose rows start with `sample`: the table's, with + /// dictionary encoding turned off for every column parquet-mr would write plain. + pub(super) fn choose(&self, sample: &[RecordBatch]) -> WriterProperties { + let mut pages: Vec<Option<FirstPage>> = self + .columns + .iter() + .map(|column| { + column + .as_ref() + .map(|column| FirstPage::new(column, self.page_bytes)) + }) + .collect(); + let mut rows_seen = 0; + for batch in sample { + let rows = batch.num_rows().min(self.page_rows - rows_seen); + let slots: Vec<(u32, usize)> = (0..rows).map(|row| (row as u32, row)).collect(); + let mut walk = LeafWalk { + pages: &mut pages, + next: 0, + rows, + }; + for column in batch.columns() { + walk.visit(column.as_ref(), &slots); + } + if walk.next != pages.len() { + // The batch does not have the leaves the schema promised, so there is no telling + // which page a value belongs to. Leave every column as the table configured it. + return self.base.clone(); + } + rows_seen += rows; + if rows_seen == self.page_rows || pages.iter().flatten().all(FirstPage::is_closed) { + break; + } + } + let plain: Vec<&ColumnPath> = self + .columns + .iter() + .zip(&pages) + .filter_map(|(column, page)| match (column, page) { + (Some(column), Some(page)) if !page.keeps_dictionary() => Some(&column.path), + _ => None, + }) + .collect(); + if rows_seen == 0 || plain.is_empty() { + return self.base.clone(); + } + plain + .into_iter() + .fold(self.base.clone().into_builder(), |builder, path| { + builder.set_column_dictionary_enabled(path.clone(), false) + }) + .build() + } +} + +/// Hands one batch's values to the first page of the leaf column they belong to, visiting the +/// leaves in Parquet schema order. +struct LeafWalk<'p, 'a> { + pages: &'p mut [Option<FirstPage<'a>>], + /// The leaf the next primitive array belongs to. + next: usize, + /// Rows of the batch that fall within the first page. + rows: usize, +} + +impl<'a> LeafWalk<'_, 'a> { + /// `slots` holds a `(row, index)` pair for every entry of `array` whose ancestors are all + /// present, in the order the column chunk stores them. + fn visit(&mut self, array: &'a dyn Array, slots: &[(u32, usize)]) { + match array.data_type() { + DataType::Struct(_) => { + let slots = present(array, slots); + for child in array.as_struct().columns() { + self.visit(child.as_ref(), &slots); + } + } + DataType::List(_) => { + let list = array.as_list::<i32>(); + self.visit( + list.values().as_ref(), + &children(array, list.value_offsets(), slots), Review Comment: [P2] Could the walker skip ineligible or closed subtrees before expanding child slots, and traverse active children with bounded storage? For a first batch containing 1,000 `array<boolean>` rows with 8,192 elements each, `children()` materializes millions of `(u32, usize)` tuples even though boolean leaves have no dictionary decision. `PartitionFeed::push` still invokes `choose()` when `wants_more()` returns false. The expected behavior is to return the unchanged properties without visiting these elements. Instead, the exact-head allocation probe measured about 128 MiB of additional temporary memory for roughly 1 MiB of input. These allocations are outside both `held_bytes` and pool reservations, creating substantial avoidable memory pressure during native writes. Skipping inactive subtrees and replacing eager child-index vectors with bounded traversal would address this. Evidence: Ran `cargo run --offline --manifest-path /tmp/comet6241-review/Cargo.toml`. The inspected harness imports the checkout’s exact `iceberg_dictionary.rs`, uses Arrow/Parquet 59.3.0, constructs the boolean-list batch described above, and measures live allocation around `choose()`. Output: `rows=1000 bools_per_row=8192 held_bytes=1028204 wants_more=false additional_peak_bytes=134233872`. No executor exhaustion was attempted. The eager allocation comes from `children()` at lines 283–293, before the leaf eligibility check at line 261. -- 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]
