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


##########
crates/paimon/src/table/aggregator/mod.rs:
##########
@@ -67,6 +81,13 @@ pub(crate) trait FieldAggregator: Send + Sync + 
std::fmt::Debug {
     /// Accumulate one input cell.
     fn agg(&mut self, array: &dyn Array, row_idx: usize) -> crate::Result<()>;
 
+    /// Replace the accumulator cell with a DELETE row's payload. Java resets
+    /// the row here, while preserving field aggregator state until the next 
PK.
+    fn replace_with_delete(&mut self, array: &dyn Array, row_idx: usize) -> 
crate::Result<()> {
+        self.reset();
+        self.agg(array, row_idx)

Review Comment:
   [P2] Copy whole-row DELETE values without evaluating normal aggregation
   
   Java AggregateMergeFunction.initRow copies each DELETE payload cell 
verbatim; this default hook instead resets the accumulator and calls agg(), 
which can normalize or discard that payload. This is observable beyond the 
stateful pick aggregators already fixed. With merge-engine=aggregation, 
aggregation.remove-record-on-delete=true, 
fields.items.aggregate-function=nested_update and fields.items.count-limit=1 
(no nested key), separately commit INSERT(items=NULL), 
DELETE(items=[(1,10),(2,20)]), INSERT(items=NULL) for one PK. Java returns both 
nested rows; Rust returns only [(1,10)] because agg() applies count-limit while 
replacing the DELETE value. A separate persisted listagg control also fails: 
INSERT('old'), DELETE(' '), INSERT(NULL) returns a one-space string in Java but 
NULL in Rust, since agg() drops blank input. Both controls were executed in 
Java and through Rust's public table API. Please separate raw-cell replacement 
from ordinary aggregation for the affected aggregators; fi
 xing nested reverse alone does not make this shared DELETE initialization path 
equivalent to Java.



##########
crates/paimon/src/table/aggregator/nested.rs:
##########
@@ -0,0 +1,614 @@
+// 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, limit_new_keys: bool) -> 
crate::Result<()> {
+        if incoming.is_null(0) {
+            return Ok(());
+        }
+        let row = self.row(incoming.as_ref())?;
+        if self.key_indices.is_empty() {
+            if !limit_new_keys || 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 !limit_new_keys
+                || 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,
+        limit_new_keys: bool,
+    ) -> 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), limit_new_keys)?;
+            }
+        }
+        Ok(())
+    }
+
+    fn preserve_raw_accumulator(&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,
+                })?;
+        let values = list.value(row_idx);
+        self.rows = (0..values.len())

Review Comment:
   [P1] Handle the raw accumulator before subsequent nested 
aggregation/retraction
   
   Preserving this raw input fixes the preceding review case, but subsequent 
agg() only consumes the new rows, and retract() only retains/removes existing 
rows. Neither rebuilds the raw accumulator as Java's nested functions do when a 
non-null next operand arrives. Public three-commit reproducer for a 
partial-update table: fields.version.sequence-group=items, 
fields.items.aggregate-function=nested_update (also reproduced with 
nested_partial_update), fields.items.nested-key=item_id; for one PK write 
INSERT(version=3,items=NULL), INSERT(version=1,items=[(1,10),(1,20)]), then 
INSERT(version=4,items=[(2,30)]), where each item is (item_id,amount). Java 
returns [(1,20),(2,30)]; Rust returns [(1,10),(1,20),(2,30)], leaving duplicate 
nested keys and the superseded value. With nested_update, replacing the third 
INSERT by DELETE(version=4,items=[(2,30)]) also leaves both key-1 versions in 
Rust, whereas Java returns only [(1,20)]. Please distinguish raw from 
normalized accumulator state and app
 ly the respective Java agg/retract rules when the next operation requires 
normalization. Do not eagerly normalize during this raw-copy step, since that 
would reintroduce the previous NULL/reverse/count-limit issue.



-- 
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