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


##########
crates/paimon/src/table/aggregator/sketch.rs:
##########
@@ -0,0 +1,441 @@
+// 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::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, 
datasketches::error::Error> {
+    // 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 and HLL
+    // modes already allocate their full backing arrays.
+    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:
   Fixed in e45995a. datasketches-rs 0.2 skipped the register bytes in compact 
HLL arrays. The union path now normalizes compact HLL_4/6/8 inputs and rejects 
truncated arrays. Java 4.2.0 dense fixtures cover all three types; Java 
heapified the Rust HLL_4 output and reported 15148.816386062443, identical to 
Java 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