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]

Reply via email to