leaves12138 commented on code in PR #955:
URL: https://github.com/apache/paimon-rust/pull/955#discussion_r4103799075


##########
crates/paimon/src/table/kv_file_writer.rs:
##########
@@ -832,7 +838,14 @@ impl KeyValueFileWriter {
                 end += 1;
             }
             let group = rows[start..end].iter().collect::<Vec<_>>();
-            merged.push(merge.merge_ordered(&group, &buffers, &output_indices, 
&output_schema)?);
+            let (row, delete) =
+                merge.merge_ordered(&group, &buffers, &output_indices, 
&output_schema)?;
+            merged.push(row);
+            merged_kinds.push(if delete {
+                RowKind::Delete as i8
+            } else {
+                RowKind::Insert as i8

Review Comment:
   [P1] Preserve singleton retract records when flushing both merge engines
   
   Every key group is reduced here, even when it contains just one DELETE. Java 
SortBufferWriteBuffer uses ReducerMergeFunctionWrapper, which returns the 
original KeyValue unchanged for a singleton. That distinction is necessary when 
the matching INSERT is in an older file. Public write/commit/read 
reproductions: (1) aggregation/product: commit INSERT(amount=100), then 
separately commit DELETE(amount=20); expected 5, actual 100. The second flush 
evaluates retract against an empty product accumulator and writes INSERT(NULL), 
losing the deletion. (2) aggregation/collect: commit [1,2,2], separately 
retract [2]; expected [1,2], actual [1,2,2]. (3) partial-update with 
fields.version.sequence-group=amount and sum(amount): commit INSERT(version=1, 
amount=100), separately commit DELETE(version=2, amount=20); expected 80, 
actual 100. The corresponding merge_partial_update_rows call at line 940 
changes the lone DELETE payload from 20 to 0 before it reaches the reader. 
Please preserve singleton
  row kind, values and sequence in BOTH write-buffer merge paths, and add 
persisted cross-commit regressions. Results must not depend on whether the 
retract shares a flush with its earlier INSERT.



##########
crates/paimon/src/table/sort_merge.rs:
##########
@@ -246,6 +250,64 @@ struct RuntimeSequenceGroup {
 }
 
 impl PartialUpdateMergeFunction {
+    fn materialize_source_row(
+        batch_idx: usize,
+        row_idx: usize,
+        batch_buffer: &[BufferedBatch],
+        source_output_col_indices: &[usize],
+        output_schema: &SchemaRef,
+    ) -> crate::Result<RecordBatch> {
+        let columns: Vec<ArrayRef> = output_schema
+            .fields()
+            .iter()
+            .enumerate()
+            .map(|(index, field)| {
+                let column = batch_buffer[batch_idx]
+                    .column_for_output(index, source_output_col_indices)
+                    .slice(row_idx, 1);
+                if !field.is_nullable() && column.is_null(0) {
+                    return Err(Error::DataInvalid {
+                        message: format!(
+                            "Partial-update delete row has NULL for 
non-nullable field '{}'",
+                            field.name()
+                        ),
+                        source: None,
+                    });
+                }
+                Ok(column)
+            })
+            .collect::<crate::Result<_>>()?;
+        RecordBatch::try_new(output_schema.clone(), columns).map_err(|error| 
Error::DataInvalid {
+            message: format!("Failed to materialize partial-update delete row: 
{error}"),
+            source: Some(Box::new(error)),
+        })
+    }
+
+    /// Java replaces the current row with the DELETE value. Later inserts may
+    /// leave some columns null, so their merge must start from that value.
+    fn reset_to_source(
+        selected_by_col: &mut [Option<(usize, usize)>],
+        aggregators: Option<&mut [Option<Box<dyn FieldAggregator>>]>,
+        batch_idx: usize,
+        row_idx: usize,
+        batch_buffer: &[BufferedBatch],
+        source_output_col_indices: &[usize],
+    ) -> crate::Result<()> {
+        selected_by_col.fill(Some((batch_idx, row_idx)));
+        if let Some(aggregators) = aggregators {
+            for (index, aggregator) in aggregators.iter_mut().enumerate() {
+                if let Some(aggregator) = aggregator {
+                    aggregator.reset();
+                    aggregator.agg(
+                        batch_buffer[batch_idx].column_for_output(index, 
source_output_col_indices),
+                        row_idx,
+                    )?;

Review Comment:
   [P2] Preserve stateful field-aggregator state on partial-update whole-row 
deletion
   
   The aggregation-engine fix now preserves state through replace_with_delete, 
but the partial-update counterpart still calls reset followed by agg. Java 
PartialUpdateMergeFunction.initRow replaces the accumulator row without 
resetting its field aggregators. With fields.version.sequence-group=amount, 
fields.amount.aggregate-function=first_non_null_value and 
partial-update.remove-record-on-sequence-group=version, process 
INSERT(version=1, amount=100), DELETE(version=2, amount=NULL), 
INSERT(version=3, amount=5). Executing Java returns an INSERT with amount=NULL; 
Rust returns amount=5 because this helper clears the initialized state. Please 
align this helper with Java's raw row replacement semantics as well; the 
regression added for aggregation does not cover partial-update.



##########
crates/paimon/src/table/aggregator/sketch.rs:
##########
@@ -0,0 +1,524 @@
+// 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.
+
+//! Apache DataSketches HLL union, matching Java `FieldHllSketchAgg`.
+
+use std::collections::BTreeSet;
+use std::sync::Arc;
+
+use arrow_array::{Array, ArrayRef, BinaryArray};
+use datasketches::error::{Error as SketchError, ErrorKind};
+use datasketches::hll::{HllSketch, HllType, HllUnion};
+
+use super::{unsupported_type_error, FieldAggregator};
+use crate::spec::DataType;
+use crate::Error;
+
+#[derive(Debug)]
+pub(crate) struct HllSketchAgg {
+    field_name: String,
+    value: Option<Vec<u8>>,
+}
+
+impl HllSketchAgg {
+    pub(crate) fn new(field_name: &str, data_type: &DataType) -> 
crate::Result<Self> {
+        if !matches!(data_type, DataType::VarBinary(_)) {
+            return Err(unsupported_type_error("hll_sketch", field_name, 
data_type));
+        }
+        Ok(Self {
+            field_name: field_name.to_string(),
+            value: None,
+        })
+    }
+}
+
+fn deserialize_for_union(bytes: &[u8]) -> Result<HllSketch, SketchError> {
+    // datasketches-rs 0.2 skips the register bytes of compact HLL arrays,
+    // although Java's compact array has the same register layout as its
+    // noncompact form. Clear the flag so the registers are actually read.
+    if bytes.len() >= 40 && bytes[2] == 7 && bytes[7] & 3 == 2 && bytes[5] & 8 
!= 0 {
+        let lg_k = bytes[3];
+        if (4..=21).contains(&lg_k) {
+            let k = 1usize << lg_k;
+            let register_len = match bytes[7] >> 2 {
+                0 => k / 2,
+                1 => k * 3 / 4,
+                2 => k,
+                _ => 0,
+            };
+            let aux_count = 
u32::from_le_bytes(bytes[36..40].try_into().unwrap()) as usize;
+            if register_len > 0 {
+                if bytes.len() < 40 + register_len + 
aux_count.saturating_mul(4) {
+                    return Err(SketchError::new(
+                        ErrorKind::InvalidData,
+                        "compact HLL array ends before its registers or 
auxiliary entries",
+                    ));
+                }
+                let mut expanded = bytes.to_vec();
+                expanded[5] &= !8;
+                // The Rust HLL_6 reader requests one extra byte for a safe
+                // packed-register window at the end of the array.
+                if bytes[7] >> 2 == 1 && expanded.len() == 40 + register_len {
+                    expanded.push(0);
+                }
+                return HllSketch::deserialize(&expanded);
+            }
+        }
+    }
+    // datasketches-rs 0.2 deserializes a compact LIST into a backing array of
+    // exactly coupon_count slots. A subsequent union silently drops every new
+    // coupon because List::update sees no empty slot. Expand only that mode to
+    // the Java noncompact representation before deserializing; SET already
+    // allocates its full backing array.
+    if bytes.len() >= 8 && bytes[2] == 7 && bytes[7] & 3 == 0 && bytes[5] & 8 
!= 0 {
+        let lg_arr = bytes[4] as u32;
+        let count = bytes[6] as usize;
+        if lg_arr == 3 {
+            let capacity = 1usize << lg_arr;
+            if capacity >= count && bytes.len() >= 8 + 4 * count {
+                let mut expanded = bytes[..8 + 4 * count].to_vec();
+                expanded[5] &= !8;
+                expanded.resize(8 + 4 * capacity, 0);
+                return HllSketch::deserialize(&expanded);
+            }
+        }
+    }
+    HllSketch::deserialize(bytes)

Review Comment:
   [P2] Normalize Java updatable HLL_4 auxiliary tables before deserializing
   
   Compact dense sketches are now handled, but valid toUpdatableByteArray 
inputs still reach this fallback unchanged. datasketches-rs 0.2 
Array4::deserialize reads aux_count consecutive coupons, whereas Java's 
updatable representation contains the sparse auxiliary hash table, including 
empty slots. Reproducer with Java DataSketches 4.2.0: create HllSketch(16, 
HLL_4), update integers 0..199999 (12 auxiliary entries), and union the sketch 
with itself. Java HllSketchUtil.union returns estimate 200552.41133627715. Rust 
returns that same result for toCompactByteArray, but 200387.20969631049 for the 
equivalent toUpdatableByteArray. Java heapifies both Rust outputs and confirms 
the discrepancy, so equivalent encodings of the same sketch produce different 
persisted results. Please normalize the auxiliary-table layout too, or reject 
that representation explicitly instead of accepting and misreading it.



##########
crates/paimon/src/table/aggregator/bitmap64.rs:
##########
@@ -0,0 +1,422 @@
+// 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.
+
+//! Java Roaring64Bitmap 1.2.1 wire format: ART over high 48 bits followed by
+//! serialized 16-bit containers. This is distinct from Rust RoaringTreemap.
+
+use std::collections::{BTreeMap, BTreeSet};
+use std::sync::Arc;
+
+use arrow_array::{Array, ArrayRef, BinaryArray};
+
+use super::{unsupported_type_error, FieldAggregator};
+use crate::spec::DataType;
+use crate::Error;
+
+#[derive(Debug)]
+pub(crate) struct Roaring64Agg {
+    field_name: String,
+    value: Option<Vec<u8>>,
+}
+
+impl Roaring64Agg {
+    pub(crate) fn new(field_name: &str, data_type: &DataType) -> 
crate::Result<Self> {
+        if !matches!(data_type, DataType::VarBinary(_)) {
+            return Err(unsupported_type_error("rbm64", field_name, data_type));
+        }
+        Ok(Self {
+            field_name: field_name.to_string(),
+            value: None,
+        })
+    }
+}
+
+impl FieldAggregator for Roaring64Agg {
+    fn name(&self) -> &'static str {
+        "rbm64"
+    }
+
+    fn reset(&mut self) {
+        self.value = None;
+    }
+
+    fn agg(&mut self, array: &dyn Array, row_idx: usize) -> crate::Result<()> {
+        if array.is_null(row_idx) {
+            return Ok(());
+        }
+        let bytes = array
+            .as_any()
+            .downcast_ref::<BinaryArray>()
+            .ok_or_else(|| Error::DataInvalid {
+                message: format!("rbm64 column '{}' requires Arrow Binary", 
self.field_name),
+                source: None,
+            })?
+            .value(row_idx);
+        let Some(accumulator) = &self.value else {
+            // Java preserves the first non-null byte sequence unchanged.
+            self.value = Some(bytes.to_vec());
+            return Ok(());
+        };
+        let mut values = decode(accumulator).map_err(|message| 
Error::DataInvalid {
+            message: format!(
+                "Invalid rbm64 accumulator for '{}': {message}",
+                self.field_name
+            ),
+            source: None,
+        })?;
+        values.extend(decode(bytes).map_err(|message| Error::DataInvalid {
+            message: format!("Invalid rbm64 input for '{}': {message}", 
self.field_name),
+            source: None,
+        })?);
+        self.value = Some(encode(&values)?);

Review Comment:
   [P2] Keep Roaring64 unions compressed instead of expanding every set bit
   
   The new union decodes each operand into BTreeSet<u64> and then builds 
another per-value grouped representation. This makes resource use proportional 
to cardinality rather than compressed container size. A concrete valid Java 
fixture containing the contiguous range [0, 1,000,000), runOptimize'd as 
Paimon's serializer does, is only 542 bytes. Unioning this bitmap with itself 
in Rust emits 131,531 bytes (~243x larger); the isolated debug test process 
peaks around 71 MiB RSS. Java verifies that the output still has one million 
elements, so this is a resource/representation regression, not a 
wrong-cardinality claim. Larger run-compressed ranges can exhaust memory 
despite tiny input cells, and the encoder never emits run containers. Please 
operate on compressed containers/a compatible compressed bitmap representation, 
preserve run compression, and add a dense-range resource regression before 
enabling this in production.



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

Reply via email to