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


##########
crates/paimon/src/table/sort_merge.rs:
##########
@@ -760,13 +981,17 @@ impl MergeFunction for AggregateMergeFunction {
         let mut ordered: Vec<_> = rows.iter().collect();
         // Stable sorting keeps input order for equal sequences.
         ordered.sort_by(|lhs, rhs| compare_sequence_order(lhs, rhs));
-        self.merge_ordered(
+        let (batch, delete) = self.merge_ordered(

Review Comment:
   [P1] Preserve singleton tombstones in the read-side reducer as well
   
   This still calls merge_ordered for a key with exactly one record. Now that 
the writer correctly preserves singleton DELETEs, the reader evaluates such a 
tombstone against an empty field accumulator and emits a live MaterializedRow 
instead. Public write/commit/scan/read reproducer: create a bucket=1 
aggregation table with id as PK and nullable INT value, commit only 
DELETE(id=1, value=20), then read it. With fields.value.aggregate-function=sum, 
Rust returns (1, -20); with product it returns (1, NULL). Both should return no 
rows. I executed Java's AggregateMergeFunction inside 
ReducerMergeFunctionWrapper: it returns the original DELETE unchanged for both 
functions, which MergeFileSplitRead's DropDeleteReader filters out. Please 
apply the same singleton-reducer semantics on the read path and add persisted 
delete-only regressions. Do not drop retract inputs globally, since groups with 
multiple versions still need to aggregate them.



##########
crates/paimon/src/table/aggregator/sketch.rs:
##########
@@ -0,0 +1,586 @@
+// 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 reads HLL_4 auxiliary coupons as a packed list.
+    // Java's updatable form instead stores them in a sparse hash table. It
+    // also skips the registers of compact HLL arrays. Normalize both forms
+    // before handing them to the Rust reader.
+    if bytes.len() >= 40 && bytes[2] == 7 && bytes[7] & 3 == 2 {
+        let lg_k = bytes[3];
+        if (4..=21).contains(&lg_k) {
+            let k = 1usize << lg_k;
+            let hll_type = bytes[7] >> 2;
+            let register_len = match hll_type {
+                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 {
+                let compact = bytes[5] & 8 != 0;
+                let payload_start = 40 + register_len;
+                let slots = if !compact && hll_type == 0 && aux_count > 0 {
+                    1usize.checked_shl(bytes[4] as u32).ok_or_else(|| {
+                        SketchError::new(
+                            ErrorKind::InvalidData,
+                            "invalid HLL_4 auxiliary table size",
+                        )
+                    })?
+                } else {
+                    aux_count
+                };
+                let payload_end = slots
+                    .checked_mul(4)
+                    .and_then(|len| payload_start.checked_add(len));
+                if payload_end.is_none_or(|end| bytes.len() < end) {
+                    return Err(SketchError::new(
+                        ErrorKind::InvalidData,
+                        "HLL array ends before its registers or auxiliary 
entries",
+                    ));
+                }
+                if !compact && hll_type == 0 && aux_count > 0 {
+                    let mut normalized = bytes[..payload_start].to_vec();
+                    let mut found = 0;
+                    for coupon in bytes[payload_start..payload_end.unwrap()]
+                        .as_chunks::<4>()
+                        .0
+                    {
+                        if *coupon != [0; 4] {
+                            normalized.extend_from_slice(coupon);
+                            found += 1;
+                        }
+                    }
+                    if found != aux_count {
+                        return Err(SketchError::new(
+                            ErrorKind::InvalidData,
+                            "HLL_4 auxiliary table count does not match 
coupons",
+                        ));
+                    }
+                    return HllSketch::deserialize(&normalized);
+                }
+                if !compact {
+                    return HllSketch::deserialize(bytes);
+                }
+                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 hll_type == 1 && expanded.len() == payload_start {
+                    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)
+}
+
+impl FieldAggregator for HllSketchAgg {
+    fn name(&self) -> &'static str {
+        "hll_sketch"
+    }
+
+    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!(
+                    "hll_sketch column '{}' requires Arrow Binary, got {:?}",
+                    self.field_name,
+                    array.data_type()
+                ),
+                source: None,
+            })?
+            .value(row_idx);
+        let Some(accumulator) = self.value.as_ref() else {
+            self.value = Some(bytes.to_vec());
+            return Ok(());
+        };
+        let left = deserialize_for_union(accumulator).map_err(|e| 
Error::DataInvalid {
+            message: format!(
+                "Invalid hll_sketch accumulator for '{}': {e}",
+                self.field_name
+            ),
+            source: Some(Box::new(e)),
+        })?;
+        let right = deserialize_for_union(bytes).map_err(|e| 
Error::DataInvalid {
+            message: format!("Invalid hll_sketch input for '{}': {e}", 
self.field_name),
+            source: Some(Box::new(e)),
+        })?;
+        let mut union = 
HllUnion::new(left.lg_config_k().max(right.lg_config_k()));
+        union.update(&right);
+        union.update(&left);
+        self.value = Some(union.get_result(HllType::Hll4).serialize());

Review Comment:
   [P1] Emit a self-consistent HLL auxiliary layout for subsequent unions
   
   The datasketches 0.2 HLL_4 serializer used here writes packed auxiliary 
coupons, but leaves the compact flag clear and byte[4] (lgAux) at zero. This 
conflicts with the new deserialize_for_union normalizer, which interprets this 
as a one-slot updatable hash table. Feed the committed hll_java_compact_aux.bin 
fixture to one HllSketchAgg three times: the second call succeeds, but the 
third fails with `HLL_4 auxiliary table count does not match coupons` (the 
output contains 12 auxiliary entries). The same failure occurs after three 
separate commits for one key through the public table read API. The output also 
breaks Java interoperability: HllSketchUtil.union(rust_result, 
original_fixture) throws `SketchesStateException: SlotNo not found: 6365`; 
merely heapifying the output and checking its stored estimate does not detect 
this. Marking this packed output as compact makes the Java control succeed. 
Please emit a valid compact or updatable header/layout and test a third 
aggregation, a per
 sisted three-commit read, and a subsequent Java union, rather than only a 
two-input union.



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