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]

Reply via email to