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


##########
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:
   Fixed in b986893. Partial-update row replacement now calls 
replace_with_delete on field aggregators, preserving Java first_non_null_value 
state. Added a persisted three-commit sequence-group regression for INSERT 100, 
DELETE NULL, INSERT 5; the result stays NULL.



##########
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:
   Fixed in b986893. Roaring64 decode/union now operates on per-container 
intervals and encodes run containers when smaller. The Java-generated 
[0,1000000) fixture is 542 bytes; self-union remains under 1024 bytes with the 
same contents.



##########
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:
   Fixed in b986893. Java HLL_4 updatable auxiliary hash slots are compacted 
before datasketches-rs deserialization, with count validation. Added Java 4.2.0 
compact/updatable fixtures for 0..199999 and asserted both self-unions produce 
Java estimate 200552.41133627715.



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