leaves12138 commented on code in PR #955: URL: https://github.com/apache/paimon-rust/pull/955#discussion_r4103202949
########## 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: [P1] Dense HLL unions emit an incorrect cardinality estimate The selected datasketches 0.2 union path is not Java-compatible once both inputs reach HLL mode. Reproducer: generate two Java HllSketch(12, HLL_4) compact sketches by updating integer values 0..9999 and 5000..14999 respectively, then feed those bytes through this aggregator. Java HllSketchUtil.union estimates 15148.816386062443 for their 15,000-element union, whereas Rust emits a sketch estimating 9966.96796876342. Heapifying those emitted bytes with Java HllSketch gives the same incorrect 9966.96796876342, so this also affects Java readers of Rust-written data. The existing small LIST/SET fixtures do not cover this mode. Please fix or replace the dense union implementation and add a Java-to-Rust-to-Java HLL-mode regression before enabling this aggregator. ########## crates/paimon/src/table/sort_merge.rs: ########## @@ -270,6 +335,26 @@ impl PartialUpdateMergeFunction { let config = PartialUpdateConfig::new(table_options); config.validate_read_mode(true, table_name)?; let groups = config.validated_sequence_groups(table_fields, primary_keys)?; + let declared_group_sequences: HashSet<&str> = groups + .iter() + .flat_map(|group| group.sequence_fields.iter().map(String::as_str)) + .collect(); + let sequence_group_partial_delete = config + .remove_record_on_sequence_group() + .map(|fields| { + fields.split(',').map(str::trim).map(|field| { + if !declared_group_sequences.contains(field) { + return Err(Error::ConfigInvalid { + message: format!("Field '{field}' in partial-update.remove-record-on-sequence-group must belong to a sequence group"), + }); + } + output_fields.iter().position(|candidate| candidate.name() == field) + .ok_or_else(|| Error::ConfigInvalid { + message: format!("Projected sequence group is missing field '{field}' required for partial delete"), + }) Review Comment: [P2] Carry whole-row-delete sequence fields through every read projection This lookup requires the deletion sequence field in output_fields, but required_sequence_fields / widen_partial_update_sequence_group_fields only add a group's sequence fields when that group intersects the user's projection. With schema (id PK, version, price), fields.version.sequence-group=price and partial-update.remove-record-on-sequence-group=version, writing one INSERT and then reading with with_projection(&["id"]) fails in new_read with ConfigInvalid: Projected sequence group is missing field 'version' required for partial delete. I reproduced this through TableWrite, commit, scan and the public read builder, as well as directly through the projection helper. Please widen the internal projection to include all sequence columns needed by whole-row deletion regardless of selected output fields; merely suppressing this error would lose deletion semantics for unprojected groups. ########## crates/paimon/src/table/sort_merge.rs: ########## @@ -683,12 +882,26 @@ impl AggregateMergeFunction { } } + let mut current_delete_row = false; for row in rows { + let kind = RowKind::from_value(row.value_kind)?; + if self.remove_record_on_delete && kind == RowKind::Delete { + for aggregator in aggregators.iter_mut().flatten() { + aggregator.reset(); + } + current_delete_row = true; + continue; Review Comment: [P1] Preserve the DELETE payload as the next aggregation accumulator With aggregation.remove-record-on-delete=true, Java AggregateMergeFunction.add calls initRow with the DELETE value; subsequent inputs aggregate against that value. Resetting every aggregator and immediately continuing instead discards it. For one key, use sum(amount) and last_non_null_value(tag), then process INSERT(100, 'old'), DELETE(10, 'deleted'), INSERT(5, NULL). Executing Java yields (15, 'deleted'), while this implementation yields (5, NULL). Since merge_ordered is shared by buffered writes and merge reads, both paths can produce different stored/query results from Java. Please preserve the DELETE accumulator semantics and cover a DELETE followed by another input, not only a terminal tombstone. -- 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]
