adriangb commented on code in PR #23565:
URL: https://github.com/apache/datafusion/pull/23565#discussion_r4118923989


##########
datafusion/physical-plan/src/spill/mod.rs:
##########
@@ -800,29 +804,133 @@ pub(crate) fn gc_view_arrays(batch: &RecordBatch) -> 
Result<RecordBatch> {
     }
 }
 
+/// Compacts the data buffers of a view array, or returns `None` when the
+/// array is too small for compaction to be worth it.
+///
+/// Arrow's `gc()` copies the bytes of every view separately, so repeated
+/// values each get their own copy. For a dictionary-encoded Parquet column
+/// this inflates the spilled data by the average repeat count: 1M rows of
+/// 1000 distinct 64 byte values spill as 64 MB
+/// (<https://github.com/apache/datafusion/issues/23564>).
+/// [`gc_dedup_view_array`] copies each distinct value once instead, and falls
+/// back to `gc()` when the values turn out to be distinct.
+fn gc_view_array<T: ByteViewType>(
+    array: &GenericByteViewArray<T>,
+) -> Option<GenericByteViewArray<T>> {
+    if !should_gc_view_array(array) {
+        return None;
+    }
+    Some(gc_dedup_view_array(array).unwrap_or_else(|| array.gc()))
+}
+
+/// Number of non-inline values [`gc_dedup_view_array`] deduplicates before
+/// it checks whether the array has enough repeats to be worth it.
+const DEDUP_SAMPLE_VALUES: usize = 256;
+
+/// Like `gc()`, but copies each distinct value once, and zeroes null views.
+///
+/// Returns `None` if the first [`DEDUP_SAMPLE_VALUES`] non-inline values
+/// hold almost no repeats, since hashing every value then costs more than
+/// deduplication saves.

Review Comment:
   The sample turns dedup off only when the first 256 non-inline values have 
fewer than 4 repeats (1.6%). If a batch has a few repeats but is mostly 
distinct, dedup hashes every value and saves few bytes. The `spill_views` suite 
has no query for this case: q02 (1000 distinct values) is above the threshold 
and saves 88% of the bytes, and q04 (all distinct) is below it.
   
   I ran the q04 shape (1M rows, 64-byte payload, `ORDER BY id`, 96M limit, 4 
partitions) with part of the rows set to one repeated value. `datafusion-cli` 
at merge-base `ed43e69` vs `7effeab`, 25 interleaved runs per side, Apple 
silicon laptop:
   
   | Data | PR / base (median of paired ratios) | PR slower in | Spilled bytes |
   | --- | --- | --- | --- |
   | All distinct (control) | 1.06x | 13 of 25 | 84.2 → 84.2 MB |
   | 5% one repeated value | **1.52x** | 23 of 25 | 84.2 → 81.1 MB |
   | 20% one repeated value | **1.32x** | 21 of 25 | 84.2 → 72.0 MB |
   | 1000 distinct values (40M limit) | 0.88x | 8 of 25 | 84.2 → 30.6 MB |
   
   The control moves by about ±10% on this machine. An earlier 9-run pass gave 
1.39x (5%) and 1.37x (20%). This agrees with q04 in the round before the sample 
(1.38–1.45x), when dedup always ran.
   
   Could the decision use the bytes that dedup saves in the sample, instead of 
the count of repeats? For example: sum the byte length of the sampled values 
and of the distinct sampled values, and fall back to `gc()` when the saving is 
below a threshold. Please set the threshold from a measured sweep of the 
repeated fraction (for example 5%, 20%, 50%, 80%). On my laptop the break-even 
is between 20% saved (1.32x slower) and 88% saved (0.88x).
   
   <details><summary>Reproducer (5% case)</summary>
   
   ```sql
   set datafusion.execution.target_partitions = 4;
   COPY (
     SELECT (value * 7919) % 1000003 AS id,
            'payload-' || lpad(CAST(CASE WHEN (value * 2654435761) % 100 < 5 
THEN 0 ELSE value END AS VARCHAR), 56, '0') AS s
     FROM generate_series(1, 1000000)
   ) TO 'hot5.parquet' STORED AS PARQUET;
   set datafusion.runtime.memory_limit = '96M';
   CREATE EXTERNAL TABLE t STORED AS PARQUET LOCATION 'hot5.parquet';
   EXPLAIN ANALYZE SELECT id, s FROM t ORDER BY id;
   ```
   
   For 20%, change `< 5` to `< 20`.
   
   </details>



##########
datafusion/physical-plan/src/spill/mod.rs:
##########
@@ -800,29 +804,133 @@ pub(crate) fn gc_view_arrays(batch: &RecordBatch) -> 
Result<RecordBatch> {
     }
 }
 
+/// Compacts the data buffers of a view array, or returns `None` when the
+/// array is too small for compaction to be worth it.
+///
+/// Arrow's `gc()` copies the bytes of every view separately, so repeated
+/// values each get their own copy. For a dictionary-encoded Parquet column
+/// this inflates the spilled data by the average repeat count: 1M rows of
+/// 1000 distinct 64 byte values spill as 64 MB
+/// (<https://github.com/apache/datafusion/issues/23564>).
+/// [`gc_dedup_view_array`] copies each distinct value once instead, and falls
+/// back to `gc()` when the values turn out to be distinct.
+fn gc_view_array<T: ByteViewType>(
+    array: &GenericByteViewArray<T>,
+) -> Option<GenericByteViewArray<T>> {
+    if !should_gc_view_array(array) {
+        return None;
+    }
+    Some(gc_dedup_view_array(array).unwrap_or_else(|| array.gc()))
+}
+
+/// Number of non-inline values [`gc_dedup_view_array`] deduplicates before
+/// it checks whether the array has enough repeats to be worth it.
+const DEDUP_SAMPLE_VALUES: usize = 256;
+
+/// Like `gc()`, but copies each distinct value once, and zeroes null views.
+///
+/// Returns `None` if the first [`DEDUP_SAMPLE_VALUES`] non-inline values
+/// hold almost no repeats, since hashing every value then costs more than
+/// deduplication saves.
+fn gc_dedup_view_array<T: ByteViewType>(
+    array: &GenericByteViewArray<T>,
+) -> Option<GenericByteViewArray<T>> {
+    let buffers = array.data_buffers();
+    let bytes_of = |view: &ByteView| {
+        let start = view.offset as usize;
+        &buffers[view.buffer_index as usize][start..start + view.length as 
usize]
+    };
+    let hasher = DefaultHashBuilder::default();
+
+    // (input view, output view) of the first occurrence of each distinct value
+    let mut distinct: HashTable<(u128, u128)> = HashTable::new();
+    let mut completed: Vec<Buffer> = vec![];
+    let mut data: Vec<u8> = vec![];
+    let mut non_inline = 0;
+    let mut views = Vec::with_capacity(array.len());
+
+    for (i, &raw) in array.views().iter().enumerate() {
+        if array.is_null(i) {
+            views.push(0);
+            continue;
+        }
+        if (raw as u32) <= MAX_INLINE_VIEW_LEN {
+            views.push(raw);
+            continue;
+        }
+
+        let view = ByteView::from(raw);
+        let bytes = bytes_of(&view);
+        let hash = hasher.hash_one(bytes);
+        let found = distinct.find(hash, |(first, _)| {
+            *first == raw || bytes_of(&ByteView::from(*first)) == bytes
+        });
+        let new_view = match found {
+            Some((_, new_view)) => *new_view,
+            None => {
+                // A view offset is a `u32`, and a buffer may not exceed 
`i32::MAX`
+                if data.len() + bytes.len() > i32::MAX as usize {
+                    completed.push(Buffer::from_vec(std::mem::take(&mut 
data)));
+                }
+                let new_view = ByteView {
+                    buffer_index: completed.len() as u32,
+                    offset: data.len() as u32,
+                    ..view
+                }
+                .as_u128();
+                data.extend_from_slice(bytes);
+                distinct.insert_unique(hash, (raw, new_view), |(first, _)| {
+                    hasher.hash_one(bytes_of(&ByteView::from(*first)))
+                });
+                new_view
+            }
+        };
+        views.push(new_view);
+
+        non_inline += 1;
+        if non_inline == DEDUP_SAMPLE_VALUES
+            && non_inline - distinct.len() < DEDUP_SAMPLE_VALUES / 64
+        {
+            return None;
+        }

Review Comment:
   `distinct` and `data` start empty and grow while the loop runs. Each time 
the table grows, it hashes the bytes of every stored value again. In a 
microbenchmark of this function (8192-row batches, 64-byte values, views into 
one shared buffer), dedup takes 8–12x the time of `gc()` when the batch has few 
repeats, and 2–4x for 1 or 1000 distinct values.
   
   This change reserves capacity after the sample passes, scaled by the 
distinct ratio of the sample. With it, the few-repeats cases drop to about 
3.5–4x of `gc()`, and 1 or 1000 distinct values improve a little. The 
all-distinct case does not change, because the sample returns before the 
reserve. Do not reserve at the start: with a table sized for the full array, 
the all-distinct case went from 1.25x to 1.6x of `gc()`.
   
   | Batch | This PR vs `gc()` | With the reserve |
   | --- | --- | --- |
   | 4 repeats in the sample, rest distinct | 10.1–10.8x | 3.4–3.7x |
   | 5% one repeated value | 10.5–11.8x | 3.1–4.1x |
   | 20% one repeated value | 8.1–8.6x | 4.0–5.1x |
   | 1000 distinct values | 3.8–4.1x | 3.1–3.3x |
   | 1 distinct value | 1.9–2.1x | 1.8x |
   
   ```suggestion
           if non_inline == DEDUP_SAMPLE_VALUES {
               if non_inline - distinct.len() < DEDUP_SAMPLE_VALUES / 64 {
                   return None;
               }
               // Size for the rest of the array from the sample's distinct 
ratio
               let remaining = array.len() - i - 1;
               let expected = remaining * distinct.len() / DEDUP_SAMPLE_VALUES;
               distinct.reserve(expected, |(first, _)| {
                   hasher.hash_one(bytes_of(&ByteView::from(*first)))
               });
               data.reserve(expected * data.len() / distinct.len());
           }
   ```



##########
datafusion/physical-plan/src/spill/mod.rs:
##########
@@ -800,29 +804,133 @@ pub(crate) fn gc_view_arrays(batch: &RecordBatch) -> 
Result<RecordBatch> {
     }
 }
 
+/// Compacts the data buffers of a view array, or returns `None` when the
+/// array is too small for compaction to be worth it.
+///
+/// Arrow's `gc()` copies the bytes of every view separately, so repeated
+/// values each get their own copy. For a dictionary-encoded Parquet column
+/// this inflates the spilled data by the average repeat count: 1M rows of
+/// 1000 distinct 64 byte values spill as 64 MB
+/// (<https://github.com/apache/datafusion/issues/23564>).
+/// [`gc_dedup_view_array`] copies each distinct value once instead, and falls
+/// back to `gc()` when the values turn out to be distinct.
+fn gc_view_array<T: ByteViewType>(
+    array: &GenericByteViewArray<T>,
+) -> Option<GenericByteViewArray<T>> {
+    if !should_gc_view_array(array) {
+        return None;
+    }
+    Some(gc_dedup_view_array(array).unwrap_or_else(|| array.gc()))
+}
+
+/// Number of non-inline values [`gc_dedup_view_array`] deduplicates before
+/// it checks whether the array has enough repeats to be worth it.
+const DEDUP_SAMPLE_VALUES: usize = 256;
+
+/// Like `gc()`, but copies each distinct value once, and zeroes null views.
+///
+/// Returns `None` if the first [`DEDUP_SAMPLE_VALUES`] non-inline values
+/// hold almost no repeats, since hashing every value then costs more than
+/// deduplication saves.
+fn gc_dedup_view_array<T: ByteViewType>(
+    array: &GenericByteViewArray<T>,
+) -> Option<GenericByteViewArray<T>> {
+    let buffers = array.data_buffers();
+    let bytes_of = |view: &ByteView| {
+        let start = view.offset as usize;
+        &buffers[view.buffer_index as usize][start..start + view.length as 
usize]
+    };
+    let hasher = DefaultHashBuilder::default();
+
+    // (input view, output view) of the first occurrence of each distinct value
+    let mut distinct: HashTable<(u128, u128)> = HashTable::new();

Review Comment:
   This table is not reserved from the memory pool. `gc()` on `main` also 
allocates its output buffers without a reservation, so the new unaccounted 
memory is the table. Each entry is 32 bytes plus one control byte, and 
hashbrown rounds the bucket count up to a power of two. For an 8192-row batch 
that passes the sample with mostly distinct values, that is 16384 buckets × 33 
bytes ≈ 530 KiB. It is freed after each batch, but it grows with the rows per 
spilled batch, and each partition that spills at the same time holds one.
   
   Could you either reserve it (for example through a `MemoryReservation` from 
the `SpillManager`), or add a comment here that states the bound and why a 
reservation is not necessary? A comment is enough for me if the rows per 
spilled batch are bounded for all callers of `gc_view_arrays`.



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