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


##########
crates/paimon/src/table/sort_merge.rs:
##########
@@ -683,12 +882,26 @@ impl AggregateMergeFunction {
             }
         }
 
+        let mut current_delete_row = false;
         for row in rows {
+            let kind = RowKind::from_value(row.value_kind)?;
+            if self.remove_record_on_delete && kind == RowKind::Delete {
+                for aggregator in aggregators.iter_mut().flatten() {
+                    aggregator.reset();
+                }
+                current_delete_row = true;
+                continue;

Review Comment:
   Fixed in e45995a. A DELETE payload now becomes the next aggregation 
accumulator, matching Java AggregateMergeFunction.add. The regression checks 
INSERT(100, old), DELETE(10, deleted), INSERT(5, NULL) => (15, deleted). A 
separate first_non_null_value test verifies Java's stateful initialization 
behavior across DELETE.



##########
crates/paimon/src/table/sort_merge.rs:
##########
@@ -270,6 +335,26 @@ impl PartialUpdateMergeFunction {
         let config = PartialUpdateConfig::new(table_options);
         config.validate_read_mode(true, table_name)?;
         let groups = config.validated_sequence_groups(table_fields, 
primary_keys)?;
+        let declared_group_sequences: HashSet<&str> = groups
+            .iter()
+            .flat_map(|group| group.sequence_fields.iter().map(String::as_str))
+            .collect();
+        let sequence_group_partial_delete = config
+            .remove_record_on_sequence_group()
+            .map(|fields| {
+                fields.split(',').map(str::trim).map(|field| {
+                    if !declared_group_sequences.contains(field) {
+                        return Err(Error::ConfigInvalid {
+                            message: format!("Field '{field}' in 
partial-update.remove-record-on-sequence-group must belong to a sequence 
group"),
+                        });
+                    }
+                    output_fields.iter().position(|candidate| candidate.name() 
== field)
+                        .ok_or_else(|| Error::ConfigInvalid {
+                            message: format!("Projected sequence group is 
missing field '{field}' required for partial delete"),
+                        })

Review Comment:
   Fixed in e45995a. required_sequence_fields now includes every sequence 
column of a group used by partial-update.remove-record-on-sequence-group, even 
for a PK-only projection. Added both composite-group unit coverage and a 
TableWrite/commit/scan/read-builder test with with_projection([id]).



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