leaves12138 commented on code in PR #955:
URL: https://github.com/apache/paimon-rust/pull/955#discussion_r4104584953
##########
crates/paimon/src/table/aggregator/mod.rs:
##########
@@ -121,13 +161,88 @@ pub(crate) fn new_aggregator(
data_type,
table_options,
)?)),
+ "collect" => Ok(Box::new(CollectAgg::new(
+ field_name,
+ data_type,
+ table_options
+ .get(&format!("fields.{field_name}.distinct"))
+ .is_some_and(|value| value.eq_ignore_ascii_case("true")),
+ )?)),
+ "merge_map" => Ok(Box::new(MergeMapAgg::new(field_name, data_type)?)),
+ "merge_map_with_keytime" => Ok(Box::new(MergeMapAgg::new_with_keytime(
+ field_name,
+ data_type,
+ table_options,
+ )?)),
+ "nested_update" | "nested_partial_update" =>
Ok(Box::new(NestedAgg::new(
+ name,
+ field_name,
+ data_type,
+ table_options,
+ )?)),
+ "rbm32" => Ok(Box::new(Roaring32Agg::new(field_name, data_type)?)),
+ "rbm64" => Ok(Box::new(Roaring64Agg::new(field_name, data_type)?)),
+ "hll_sketch" => Ok(Box::new(HllSketchAgg::new(field_name,
data_type)?)),
+ "theta_sketch" => Ok(Box::new(ThetaSketchAgg::new(field_name,
data_type)?)),
+ "primary-key" => Ok(Box::new(PrimaryKeyAgg::new(field_name,
data_type)?)),
_ => Err(crate::Error::ConfigInvalid {
message: format!(
"Unknown aggregate function '{name}' for field '{field_name}';
\
supported: sum, product, min, max, last_value, first_value, \
- last_non_null_value, first_non_null_value, bool_and, bool_or,
listagg"
+ last_non_null_value, first_non_null_value, bool_and, bool_or,
listagg, collect, merge_map, merge_map_with_keytime, nested_update,
nested_partial_update, rbm32, rbm64, hll_sketch, theta_sketch, primary-key"
),
}),
+ };
+ let aggregator = aggregator?;
+ let ignore_retract = table_options
+ .get(&format!("fields.{field_name}.ignore-retract"))
+ .is_some_and(|value| value.eq_ignore_ascii_case("true"));
+ if ignore_retract
+ && table_options
+ .get("aggregation.remove-record-on-delete")
+ .is_some_and(|value| value.eq_ignore_ascii_case("true"))
+ {
+ return Err(crate::Error::ConfigInvalid {
+ message: format!("aggregation.remove-record-on-delete conflicts
with fields.{field_name}.ignore-retract"),
+ });
+ }
+ if ignore_retract {
+ Ok(Box::new(IgnoreRetractAgg(aggregator)))
+ } else {
+ Ok(aggregator)
+ }
+}
+
+#[derive(Debug)]
+struct IgnoreRetractAgg(Box<dyn FieldAggregator>);
+
+impl FieldAggregator for IgnoreRetractAgg {
+ fn name(&self) -> &'static str {
+ self.0.name()
+ }
+ fn reset(&mut self) {
+ self.0.reset();
+ }
+ fn agg(&mut self, array: &dyn Array, row_idx: usize) -> crate::Result<()> {
+ self.0.agg(array, row_idx)
+ }
+ fn replace_with_delete(&mut self, array: &dyn Array, row_idx: usize) ->
crate::Result<()> {
+ self.0.replace_with_delete(array, row_idx)
+ }
+ fn agg_reversed(&mut self, array: &dyn Array, row_idx: usize) ->
crate::Result<()> {
+ // Java FieldIgnoreRetractAgg inherits FieldAggregator#aggReversed,
+ // which invokes wrapper.agg(input, accumulator) even if the wrapped
+ // function overrides its own reverse behavior.
+ let current = self.0.result()?;
+ self.0.reset();
Review Comment:
[P1] Preserve the wrapped field state when reversing ignore-retract
aggregation
Java FieldIgnoreRetractAgg inherits aggReversed(accumulator, input) =
agg(input, accumulator); it does not reset the wrapped aggregator. Resetting
here changes the initialized flag used by first_value / first_non_null_value.
Public two-commit reproducer with merge-engine=partial-update,
fields.version.sequence-group=value,
fields.value.aggregate-function=first_value and
fields.value.ignore-retract=true: commit DELETE(version=3,value=20), then
INSERT(version=2,value=5) for the same key. Java returns (version=3,value=20),
while Rust returns (3,5). There is also an insert-only case with
first_non_null_value: INSERT(3,10), then INSERT(1,NULL); Java returns (3,NULL),
Rust (3,10). I executed both Java reducer controls and both Rust persisted
reads. Please implement the operand swap without resetting the field-function
state, and cover both initially uninitialized state after a retract and
initialized state with a NULL older operand.
##########
crates/paimon/src/table/aggregator/nested.rs:
##########
@@ -0,0 +1,559 @@
+// 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.
+
+//! `nested_update` and `nested_partial_update` for ARRAY<ROW> fields.
+
+use std::cmp::Ordering;
+use std::collections::HashMap;
+use std::sync::Arc;
+
+use arrow_array::{new_empty_array, new_null_array, Array, ArrayRef, ListArray,
StructArray};
+use arrow_buffer::{OffsetBuffer, ScalarBuffer};
+use arrow_ord::ord::make_comparator;
+use arrow_schema::{DataType as ArrowDataType, FieldRef, Fields, SortOptions};
+use arrow_select::concat::concat;
+
+use super::{unsupported_type_error, FieldAggregator};
+use crate::arrow::paimon_type_to_arrow;
+use crate::spec::DataType;
+use crate::Error;
+
+#[derive(Clone, Copy, Debug)]
+enum NestedMode {
+ Update,
+ PartialUpdate,
+}
+
+#[derive(Clone, Copy, Debug)]
+enum NullKeyStrategy {
+ Merge,
+ Ignore,
+ Error,
+}
+
+#[derive(Debug)]
+pub(crate) struct NestedAgg {
+ field_name: String,
+ mode: NestedMode,
+ element_field: FieldRef,
+ row_fields: Fields,
+ key_indices: Vec<usize>,
+ sequence_indices: Vec<usize>,
+ null_key_strategy: NullKeyStrategy,
+ count_limit: usize,
+ seen_input: bool,
+ rows: Vec<ArrayRef>,
+}
+
+impl NestedAgg {
+ pub(crate) fn new(
+ name: &str,
+ field_name: &str,
+ data_type: &DataType,
+ options: &HashMap<String, String>,
+ ) -> crate::Result<Self> {
+ let mode = match name {
+ "nested_update" => NestedMode::Update,
+ "nested_partial_update" => NestedMode::PartialUpdate,
+ _ => unreachable!("Only nested aggregators call this constructor"),
+ };
+ let DataType::Array(array_type) = data_type else {
+ return Err(unsupported_type_error(name, field_name, data_type));
+ };
+ let DataType::Row(row_type) = array_type.element_type() else {
+ return Err(unsupported_type_error(name, field_name, data_type));
+ };
+ let ArrowDataType::List(element_field) =
paimon_type_to_arrow(data_type)? else {
+ unreachable!("ARRAY must map to Arrow List")
+ };
+ let ArrowDataType::Struct(row_fields) = element_field.data_type() else
{
+ unreachable!("ARRAY<ROW> must map to List<Struct>")
+ };
+ let row_fields = row_fields.clone();
+ let key_option = format!("fields.{field_name}.nested-key");
+ let sequence_option =
format!("fields.{field_name}.nested-sequence-field");
+ let strategy_option =
format!("fields.{field_name}.nested-key-null-strategy");
+ let key_indices =
+ parse_nested_fields(options.get(&key_option), row_type.fields(),
&key_option)?;
+ let sequence_indices = parse_nested_fields(
+ options.get(&sequence_option),
+ row_type.fields(),
+ &sequence_option,
+ )?;
+ if matches!(mode, NestedMode::PartialUpdate) && key_indices.is_empty()
{
+ return Err(Error::ConfigInvalid {
+ message: format!(
+ "nested_partial_update field '{field_name}' requires
'{key_option}'"
+ ),
+ });
+ }
+ if key_indices.is_empty()
+ && (options.contains_key(&strategy_option) ||
!sequence_indices.is_empty())
+ {
+ return Err(Error::ConfigInvalid {
+ message: format!(
+ "Nested key strategy and sequence fields for
'{field_name}' require '{key_option}'"
+ ),
+ });
+ }
+ let null_key_strategy = match
options.get(&strategy_option).map(String::as_str) {
+ None => NullKeyStrategy::Merge,
+ Some(value) if value.eq_ignore_ascii_case("merge") =>
NullKeyStrategy::Merge,
+ Some(value) if value.eq_ignore_ascii_case("ignore") =>
NullKeyStrategy::Ignore,
+ Some(value) if value.eq_ignore_ascii_case("error") =>
NullKeyStrategy::Error,
+ Some(value) => {
+ return Err(Error::ConfigInvalid {
+ message: format!("Invalid nested-key-null-strategy
'{value}'"),
+ })
+ }
+ };
+ let count_limit = options
+ .get(&format!("fields.{field_name}.count-limit"))
+ .map(|value| {
+ value
+ .parse::<i32>()
+ .map(|n| n.max(0) as usize)
+ .map_err(|_| Error::ConfigInvalid {
+ message: format!("Invalid nested_update count-limit
'{value}'"),
+ })
+ })
+ .transpose()?
+ .unwrap_or(i32::MAX as usize);
+ Ok(Self {
+ field_name: field_name.to_string(),
+ mode,
+ element_field,
+ row_fields,
+ key_indices,
+ sequence_indices,
+ null_key_strategy,
+ count_limit,
+ seen_input: false,
+ rows: Vec::new(),
+ })
+ }
+
+ fn row<'a>(&self, value: &'a dyn Array) -> crate::Result<&'a StructArray> {
+ value
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Nested aggregator for '{}' requires Arrow Struct
elements",
+ self.field_name
+ ),
+ source: None,
+ })
+ }
+
+ fn key_is_valid(&self, row: &StructArray) -> crate::Result<bool> {
+ if self
+ .key_indices
+ .iter()
+ .all(|&index| row.column(index).is_valid(0))
+ {
+ return Ok(true);
+ }
+ match self.null_key_strategy {
+ NullKeyStrategy::Merge => Ok(true),
+ NullKeyStrategy::Ignore => Ok(false),
+ NullKeyStrategy::Error => Err(Error::DataInvalid {
+ message: "Nested key contains null values. Primary key fields
must not be null."
+ .into(),
+ source: None,
+ }),
+ }
+ }
+
+ fn same_key(&self, left: &StructArray, right: &StructArray) -> bool {
+ self.key_indices
+ .iter()
+ .all(|&index| left.column(index).as_ref() ==
right.column(index).as_ref())
+ }
+
+ fn compare_sequence(&self, left: &StructArray, right: &StructArray) ->
crate::Result<Ordering> {
+ for &index in &self.sequence_indices {
+ let comparator = make_comparator(
+ left.column(index).as_ref(),
+ right.column(index).as_ref(),
+ SortOptions {
+ descending: false,
+ nulls_first: true,
+ },
+ )
+ .map_err(|e| Error::DataInvalid {
+ message: format!(
+ "Failed to compare nested sequence for '{}': {e}",
+ self.field_name
+ ),
+ source: Some(Box::new(e)),
+ })?;
+ let ordering = comparator(0, 0);
+ if !ordering.is_eq() {
+ return Ok(ordering);
+ }
+ }
+ Ok(Ordering::Equal)
+ }
+
+ fn partial_update(&self, old: &StructArray, new: &StructArray) ->
crate::Result<ArrayRef> {
+ let columns = (0..self.row_fields.len())
+ .map(|index| {
+ let column = new.column(index);
+ if column.is_valid(0) {
+ column.clone()
+ } else {
+ old.column(index).clone()
+ }
+ })
+ .collect();
+ let result = StructArray::try_new(self.row_fields.clone(), columns,
None).map_err(|e| {
+ Error::DataInvalid {
+ message: format!(
+ "Failed to build nested partial update for '{}': {e}",
+ self.field_name
+ ),
+ source: Some(Box::new(e)),
+ }
+ })?;
+ Ok(Arc::new(result))
+ }
+
+ fn add_row(&mut self, incoming: ArrayRef) -> crate::Result<()> {
+ let row = self.row(incoming.as_ref())?;
+ if self.key_indices.is_empty() {
+ if self.rows.len() < self.count_limit {
+ self.rows.push(incoming);
+ }
+ return Ok(());
+ }
+ if !self.key_is_valid(row)? {
+ return Ok(());
+ }
+ let position = self.rows.iter().position(|existing| {
+ self.same_key(self.row(existing.as_ref()).expect("stored Struct
row"), row)
+ });
+ match position {
+ Some(index) => {
+ let existing = self.row(self.rows[index].as_ref())?;
+ let replacement = match self.mode {
+ NestedMode::PartialUpdate => self.partial_update(existing,
row)?,
+ NestedMode::Update
+ if self.sequence_indices.is_empty()
+ || self.compare_sequence(row, existing)?.is_ge() =>
+ {
+ incoming
+ }
+ NestedMode::Update => return Ok(()),
+ };
+ self.rows[index] = replacement;
+ }
+ None if matches!(self.mode, NestedMode::PartialUpdate)
+ || self.rows.len() < self.count_limit =>
+ {
+ self.rows.push(incoming)
+ }
+ None => {}
+ }
+ Ok(())
+ }
+
+ fn consume(&mut self, array: &dyn Array, row_idx: usize) ->
crate::Result<()> {
+ if array.is_null(row_idx) {
+ return Ok(());
+ }
+ let list =
+ array
+ .as_any()
+ .downcast_ref::<ListArray>()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Nested aggregator for '{}' requires Arrow List",
+ self.field_name
+ ),
+ source: None,
+ })?;
+ self.seen_input = true;
+ let values = list.value(row_idx);
+ for index in 0..values.len() {
+ if values.is_valid(index) {
+ self.add_row(values.slice(index, 1))?;
+ }
+ }
+ Ok(())
+ }
+}
+
+fn parse_nested_fields(
+ option: Option<&String>,
+ fields: &[crate::spec::DataField],
+ option_name: &str,
+) -> crate::Result<Vec<usize>> {
+ option
+ .map(|value| {
+ value
+ .split(',')
+ .map(str::trim)
+ .map(|name| {
+ fields
+ .iter()
+ .position(|field| field.name() == name)
+ .ok_or_else(|| Error::ConfigInvalid {
+ message: format!(
+ "Nested field '{name}' referenced by
'{option_name}' does not exist"
+ ),
+ })
+ })
+ .collect()
+ })
+ .unwrap_or_else(|| Ok(Vec::new()))
+}
+
+impl FieldAggregator for NestedAgg {
+ fn name(&self) -> &'static str {
+ match self.mode {
+ NestedMode::Update => "nested_update",
+ NestedMode::PartialUpdate => "nested_partial_update",
+ }
+ }
+
+ fn reset(&mut self) {
+ self.seen_input = false;
+ self.rows.clear();
+ }
+
+ fn agg(&mut self, array: &dyn Array, row_idx: usize) -> crate::Result<()> {
+ self.consume(array, row_idx)
+ }
+
+ fn agg_reversed(&mut self, array: &dyn Array, row_idx: usize) ->
crate::Result<()> {
+ let current = std::mem::take(&mut self.rows);
+ let current_seen = self.seen_input;
+ self.seen_input = false;
+ self.consume(array, row_idx)?;
Review Comment:
[P2] Do not reapply nested count-limit to the accumulator side of a reverse
merge
Calling consume on the older operand normalizes/truncates it as a new input.
Java's default aggReversed instead calls agg(older_input, current_accumulator):
FieldNestedUpdateAgg preserves the existing accumulator when the new input is
NULL, and applies count-limit to newly appended rows rather than retroactively
truncating that accumulator. Public two-commit reproducer: partial-update table
with fields.version.sequence-group=items,
fields.items.aggregate-function=nested_update, fields.items.count-limit=1, and
no nested key; commit INSERT(version=3,items=NULL), then
INSERT(version=1,items=[ROW(1),ROW(2)]). Singleton writes preserve both source
elements. Java returns both item IDs [1,2]; Rust returns only [1]. This is a
cardinality difference, not unspecified hash-map ordering. Please preserve the
raw accumulator-side semantics of the reverse call and add a persisted
out-of-order/NULL-current regression.
##########
crates/paimon/src/table/aggregator/value.rs:
##########
@@ -165,10 +165,40 @@ macro_rules! pick_agg {
self.0.agg(array, row_idx);
Ok(())
}
+ fn replace_with_delete(
+ &mut self,
+ array: &dyn Array,
+ row_idx: usize,
+ ) -> crate::Result<()> {
+ self.0.winner = Some(array.slice(row_idx, 1));
Review Comment:
[P1] Make a NULL delete payload eligible for last_non_null_value reverse
fallback
The new raw replacement stores a NULL cell as Some(null_array), but
PickPolicy::LastNonNull::agg_reversed still treats only winner.is_none() as an
empty accumulator. These two representations of SQL NULL now produce different
results. Reproducer through three separate commits:
merge-engine=partial-update, fields.version.sequence-group=value,
fields.value.aggregate-function=last_non_null_value,
partial-update.remove-record-on-sequence-group=version; write
INSERT(version=1,value=100), DELETE(version=2,value=NULL), then
INSERT(version=1,value=5). Java's reducer returns (version=2,value=5), because
aggReversed(NULL,5) calls agg(5,NULL) and keeps 5. Rust returns (2,NULL),
silently losing the non-null fallback. Please make the reverse path inspect the
stored cell's nullness as well as Option::None (or maintain an equivalent NULL
representation), without resetting the stateful first-value aggregators.
--
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]