JingsongLi commented on code in PR #955: URL: https://github.com/apache/paimon-rust/pull/955#discussion_r4105146981
########## 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: Fixed in 3128b69. Nested aggregators now retain a raw accumulator after the NULL/reverse case and normalize it only when a later non-NULL agg or retract requires Java's rebuild. This handles nested_update and nested_partial_update key deduplication; retract uses Java's unconditional key overwrite rule. Added persisted three-commit regressions for both functions and the delete variant. ########## 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: Fixed in 3128b69. Whole-row DELETE replacement now copies nested_update and listagg cells verbatim, preserving rows beyond count-limit and blank strings until the next operation. I also covered collect's raw DELETE payload and Java's subsequent distinct behavior. Added persisted nested/listagg controls and focused unit tests. -- 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]
