andygrove commented on code in PR #6241:
URL: https://github.com/apache/datafusion-comet/pull/6241#discussion_r4111609186


##########
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:
   Fixed in 3398c8827. The chooser now only reads leaves that have a dictionary 
choice to make. Booleans, fixed-length binary and columns whose dictionary is 
already off never become candidates, so a schema with nothing else holds no 
rows back and `choose` returns without touching the batch. Each candidate is 
walked one column at a time through a lazy iterator that stops where its first 
page ends, so no child index vectors are built. I also dropped the per-value id 
vector. The RLE-hybrid length only depends on whether each id repeats the one 
before it, so the encoder state machine now runs on that directly, and the only 
state kept per column is its dictionary.
   
   With your probe shape against the new head in release mode, the 1000 × 8192 
`array<boolean>` batch goes from 134,233,856 bytes of additional peak 
allocation to 25 bytes (3 µs). A 1000 × 8192 `array<int>` batch goes from 138 
MB to 52 KB, and the 20k-row, 75-column test corpus from 40 MiB to 3.2 MiB in 
the same time. The decisions don't change: the parquet-mr ground-truth tests 
pass as before, and so do the JVM parity tests and the rest of 
`CometIcebergWriteActionSuite`.
   
   I couldn't find a clean way to measure peak allocation inside the lib's test 
binary, since the crate already installs its process-wide accounting allocator. 
So the new tests pin the two properties that bound it instead. 
`columns_with_no_choice_are_never_read` checks that boolean lists never become 
candidates, including under a struct. 
`a_column_is_read_only_as_far_as_its_first_page` checks that a `list<int>` 
column with 8192 entries per row pulls exactly 101 of 400 top-level rows before 
its page closes at row 100. Collecting the entries eagerly, as before, pulls 
all 400 and fails it.



-- 
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