JingsongLi commented on code in PR #955:
URL: https://github.com/apache/paimon-rust/pull/955#discussion_r4104336110
##########
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:
Fixed in c5d3db1. Read-side aggregation now applies Java
ReducerMergeFunctionWrapper singleton behavior: one add reuses its source row,
one retract is omitted. The same shared behavior covers partial-update. Added
persisted delete-only sum/product regressions and a partial-update singleton
delete regression; multi-version aggregation still runs normally.
##########
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:
Fixed in c5d3db1. Rust HLL_4 results with packed auxiliary coupons are now
marked compact, so subsequent Rust/Java unions use the correct layout. Added
three-input unit coverage and a persisted three-commit read. I also passed the
Rust second-union bytes to Java DataSketches 4.2.0, then ran the same
HllSketch.heapify + Union.heapify + update sequence as HllSketchUtil.union
against the original Java fixture; it succeeds with 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]