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


##########
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:
   Fixed in 48bb7f1. IgnoreRetractAgg now reverses without resetting the 
wrapped aggregator, preserving first_value / first_non_null_value state. Added 
direct and persisted two-commit regressions for both reported cases; also 
covered collect's Java-specific reverse override.



##########
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:
   Fixed in 48bb7f1. LastNonNull reverse now treats a stored NULL Arrow cell as 
an empty winner, so the older non-NULL value supplies the fallback. Added 
direct and three-commit persisted regression for the reported sequence-group 
delete case.



##########
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:
   Fixed in 48bb7f1. Nested reverse preserves the raw older accumulator, 
including rows beyond count-limit when the current input is NULL; when current 
rows exist, it limits only additions to the older accumulator. Added direct and 
persisted two-commit regressions.



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