This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new 9908c479 fix(partial-update): support ignore-delete semantics (#573)
9908c479 is described below

commit 9908c479537d714bf922cc57b0548df1412eefcc
Author: shyjsarah <[email protected]>
AuthorDate: Tue Jul 21 20:49:44 2026 +0800

    fix(partial-update): support ignore-delete semantics (#573)
---
 crates/integrations/datafusion/tests/pk_tables.rs  |  71 +++++++-
 crates/paimon/src/spec/core_options.rs             |  17 ++
 crates/paimon/src/spec/partial_update.rs           |  36 +++-
 crates/paimon/src/spec/schema.rs                   |  97 ++++++++++-
 crates/paimon/src/table/kv_file_writer.rs          | 194 ++++++++++++++++++++-
 crates/paimon/src/table/sort_merge.rs              | 174 +++++++++++++++++-
 crates/paimon/tests/incremental_batch_scan_test.rs |  41 +++++
 docs/src/sql.md                                    |   9 +
 8 files changed, 615 insertions(+), 24 deletions(-)

diff --git a/crates/integrations/datafusion/tests/pk_tables.rs 
b/crates/integrations/datafusion/tests/pk_tables.rs
index 4b0840c0..0bad6e3b 100644
--- a/crates/integrations/datafusion/tests/pk_tables.rs
+++ b/crates/integrations/datafusion/tests/pk_tables.rs
@@ -32,9 +32,13 @@ use common::{
     collect_id_name, collect_id_value, collect_int_int_str, 
create_sql_context, create_test_env,
     row_count, setup_sql_context, string_value,
 };
-use datafusion::arrow::array::{Array, Int32Array, Int64Array};
+use datafusion::arrow::array::{Array, Int32Array, Int64Array, Int8Array, 
RecordBatch};
+use datafusion::arrow::datatypes::{
+    DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema,
+};
 use paimon::catalog::Identifier;
 use paimon::Catalog;
+use std::sync::Arc;
 
 // ======================= Basic PK Write + Read =======================
 
@@ -174,6 +178,71 @@ async fn test_pk_partial_update_fixed_bucket_e2e() {
     );
 }
 
+#[tokio::test]
+async fn test_pk_partial_update_ignore_delete_alias_e2e() {
+    let (_tmp, catalog) = create_test_env();
+    let sql_context = create_sql_context(catalog.clone()).await;
+    sql_context
+        .sql("CREATE SCHEMA paimon.test_db")
+        .await
+        .unwrap();
+    sql_context
+        .sql(
+            "CREATE TABLE paimon.test_db.t_partial_update_ignore_delete (
+                id INT NOT NULL, value INT,
+                PRIMARY KEY (id)
+            ) WITH (
+                'bucket' = '1',
+                'merge-engine' = 'partial-update',
+                'partial-update.ignore-delete' = 'true'
+            )",
+        )
+        .await
+        .unwrap();
+
+    let table = catalog
+        .get_table(&Identifier::new(
+            "test_db",
+            "t_partial_update_ignore_delete",
+        ))
+        .await
+        .unwrap();
+    let batch = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+            ArrowField::new("_VALUE_KIND", ArrowDataType::Int8, false),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1, 1, 2, 1])),
+            Arc::new(Int32Array::from(vec![
+                Some(10),
+                Some(999),
+                Some(200),
+                Some(20),
+            ])),
+            Arc::new(Int8Array::from(vec![0, 3, 1, 2])),
+        ],
+    )
+    .unwrap();
+    let write_builder = table.new_write_builder();
+    let mut write = write_builder.new_write().unwrap();
+    write.write_arrow_batch(&batch).await.unwrap();
+    let messages = write.prepare_commit().await.unwrap();
+    write_builder.new_commit().commit(messages).await.unwrap();
+
+    assert_eq!(
+        collect_id_value(
+            &sql_context,
+            "SELECT id, value
+             FROM paimon.test_db.t_partial_update_ignore_delete
+             ORDER BY id",
+        )
+        .await,
+        vec![(1, 20)]
+    );
+}
+
 /// Partial updates of one key within a single INSERT are merged at flush
 /// (mirrors Java MergeTreeWriter#flushWriteBuffer): the flushed file holds
 /// one row per key, so SELECT and COUNT(*) agree.
diff --git a/crates/paimon/src/spec/core_options.rs 
b/crates/paimon/src/spec/core_options.rs
index e5a61b54..fd4809bc 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -1966,9 +1966,26 @@ mod tests {
         let opts = CoreOptions::new(&fallback);
         assert!(opts.ignore_delete());
 
+        let partial_update = HashMap::from([(
+            "partial-update.ignore-delete".to_string(),
+            "true".to_string(),
+        )]);
+        let opts = CoreOptions::new(&partial_update);
+        assert!(opts.ignore_delete());
+
         let primary = HashMap::from([("ignore-delete".to_string(), 
"true".to_string())]);
         let opts = CoreOptions::new(&primary);
         assert!(opts.ignore_delete());
+
+        let primary_precedence = HashMap::from([
+            ("ignore-delete".to_string(), "false".to_string()),
+            (
+                "partial-update.ignore-delete".to_string(),
+                "true".to_string(),
+            ),
+        ]);
+        let opts = CoreOptions::new(&primary_precedence);
+        assert!(!opts.ignore_delete());
     }
 
     #[test]
diff --git a/crates/paimon/src/spec/partial_update.rs 
b/crates/paimon/src/spec/partial_update.rs
index b7ae1b6d..07a57d73 100644
--- a/crates/paimon/src/spec/partial_update.rs
+++ b/crates/paimon/src/spec/partial_update.rs
@@ -21,6 +21,7 @@ const MERGE_ENGINE_OPTION: &str = "merge-engine";
 const PARTIAL_UPDATE_ENGINE: &str = "partial-update";
 const IGNORE_DELETE_OPTION: &str = "ignore-delete";
 const IGNORE_DELETE_SUFFIX: &str = ".ignore-delete";
+const PARTIAL_UPDATE_IGNORE_DELETE_OPTION: &str = 
"partial-update.ignore-delete";
 const PARTIAL_UPDATE_REMOVE_RECORD_ON_DELETE_OPTION: &str =
     "partial-update.remove-record-on-delete";
 const PARTIAL_UPDATE_REMOVE_RECORD_ON_SEQUENCE_GROUP_OPTION: &str =
@@ -116,8 +117,9 @@ impl<'a> PartialUpdateConfig<'a> {
 }
 
 fn is_unsupported_partial_update_option(key: &str) -> bool {
-    key == IGNORE_DELETE_OPTION
-        || key.ends_with(IGNORE_DELETE_SUFFIX)
+    (key.ends_with(IGNORE_DELETE_SUFFIX)
+        && key != IGNORE_DELETE_OPTION
+        && key != PARTIAL_UPDATE_IGNORE_DELETE_OPTION)
         || key == PARTIAL_UPDATE_REMOVE_RECORD_ON_DELETE_OPTION
         || key == PARTIAL_UPDATE_REMOVE_RECORD_ON_SEQUENCE_GROUP_OPTION
         || key == FIELDS_DEFAULT_AGG_FUNCTION_OPTION
@@ -157,6 +159,32 @@ mod tests {
         );
     }
 
+    #[test]
+    fn test_validate_create_mode_accepts_partial_update_ignore_delete() {
+        for value in ["true", "false"] {
+            let options = 
partial_update_options(&[(PARTIAL_UPDATE_IGNORE_DELETE_OPTION, value)]);
+            let config = PartialUpdateConfig::new(&options);
+
+            assert_eq!(
+                config.validate_create_mode(true).unwrap(),
+                Some(PartialUpdateMode::Basic)
+            );
+        }
+    }
+
+    #[test]
+    fn test_validate_create_mode_accepts_ignore_delete() {
+        for value in ["true", "false"] {
+            let options = partial_update_options(&[(IGNORE_DELETE_OPTION, 
value)]);
+            let config = PartialUpdateConfig::new(&options);
+
+            assert_eq!(
+                config.validate_create_mode(true).unwrap(),
+                Some(PartialUpdateMode::Basic)
+            );
+        }
+    }
+
     #[test]
     fn test_validate_create_mode_ignores_non_pk_tables() {
         let options = partial_update_options(&[(IGNORE_DELETE_OPTION, 
"true")]);
@@ -168,10 +196,10 @@ mod tests {
     #[test]
     fn test_validate_create_mode_rejects_unsupported_partial_update_options() {
         for key in [
-            IGNORE_DELETE_OPTION,
-            "partial-update.ignore-delete",
             PARTIAL_UPDATE_REMOVE_RECORD_ON_DELETE_OPTION,
             PARTIAL_UPDATE_REMOVE_RECORD_ON_SEQUENCE_GROUP_OPTION,
+            "deduplicate.ignore-delete",
+            "fields.price.ignore-delete",
             "fields.price.sequence-group",
             "fields.price.aggregate-function",
             FIELDS_DEFAULT_AGG_FUNCTION_OPTION,
diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs
index 413c0d99..0245cf41 100644
--- a/crates/paimon/src/spec/schema.rs
+++ b/crates/paimon/src/spec/schema.rs
@@ -453,6 +453,15 @@ impl TableSchema {
         new_schema.highest_field_id =
             
highest_field_id.max(Self::current_highest_field_id(&new_schema.fields));
 
+        if PartialUpdateConfig::new(&new_schema.options).is_enabled()
+            && self.core_options().ignore_delete()
+            && !CoreOptions::new(&new_schema.options).ignore_delete()
+        {
+            return Err(crate::Error::Unsupported {
+                message: "Cannot change ignore-delete from true to 
false.".to_string(),
+            });
+        }
+
         // Re-run create-time validations on the final schema, mirroring Java
         // `SchemaValidation.validateTableSchema` after applying changes.
         Schema::validate_key_field_types(
@@ -2099,7 +2108,7 @@ mod tests {
     #[test]
     fn test_partial_update_schema_validation_rejects_unsupported_options() {
         for (key, value) in [
-            ("ignore-delete", "true"),
+            ("fields.value.ignore-delete", "true"),
             ("fields.value.sequence-group", "g1"),
             ("fields.default-aggregate-function", "last_non_null"),
         ] {
@@ -2119,6 +2128,27 @@ mod tests {
         }
     }
 
+    #[test]
+    fn test_partial_update_schema_validation_accepts_ignore_delete_options() {
+        for (key, value) in [
+            ("ignore-delete", "true"),
+            ("ignore-delete", "false"),
+            ("partial-update.ignore-delete", "true"),
+            ("partial-update.ignore-delete", "false"),
+        ] {
+            let schema = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("value", DataType::Int(IntType::new()))
+                .primary_key(["id"])
+                .option("merge-engine", "partial-update")
+                .option(key, value)
+                .build()
+                .unwrap();
+
+            assert_eq!(schema.options().get(key).map(String::as_str), 
Some(value));
+        }
+    }
+
     #[test]
     fn test_aggregation_schema_validation_accepts_basic_options() {
         let schema = Schema::builder()
@@ -2709,7 +2739,7 @@ mod tests {
     }
 
     #[test]
-    fn test_partial_update_apply_changes_rejects_unsupported_option() {
+    fn test_partial_update_apply_changes_accepts_ignore_delete_option() {
         let table_schema = TableSchema::new(
             0,
             &Schema::builder()
@@ -2721,21 +2751,70 @@ mod tests {
                 .unwrap(),
         );
 
-        let err = table_schema
+        let new_schema = table_schema
             .apply_changes(vec![crate::spec::SchemaChange::set_option(
                 "ignore-delete".to_string(),
                 "true".to_string(),
             )])
-            .unwrap_err();
+            .unwrap();
 
-        assert!(
-            matches!(err, crate::Error::ConfigInvalid { ref message }
-                if message.contains("merge-engine=partial-update")
-                    && message.contains("ignore-delete")),
-            "partial-update alter should reject unsupported option, got 
{err:?}"
+        assert_eq!(
+            new_schema
+                .options()
+                .get("ignore-delete")
+                .map(String::as_str),
+            Some("true")
         );
     }
 
+    #[test]
+    fn test_partial_update_apply_changes_rejects_disabling_ignore_delete() {
+        for (existing_options, change) in [
+            (
+                vec![("ignore-delete", "true")],
+                crate::spec::SchemaChange::set_option(
+                    "ignore-delete".to_string(),
+                    "false".to_string(),
+                ),
+            ),
+            (
+                vec![("ignore-delete", "true")],
+                
crate::spec::SchemaChange::remove_option("ignore-delete".to_string()),
+            ),
+            (
+                vec![("partial-update.ignore-delete", "true")],
+                crate::spec::SchemaChange::set_option(
+                    "ignore-delete".to_string(),
+                    "false".to_string(),
+                ),
+            ),
+            (
+                vec![("partial-update.ignore-delete", "true")],
+                crate::spec::SchemaChange::remove_option(
+                    "partial-update.ignore-delete".to_string(),
+                ),
+            ),
+        ] {
+            let mut builder = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("value", DataType::Int(IntType::new()))
+                .primary_key(["id"])
+                .option("merge-engine", "partial-update");
+            for (key, value) in existing_options {
+                builder = builder.option(key, value);
+            }
+            let table_schema = TableSchema::new(0, &builder.build().unwrap());
+
+            let err = table_schema.apply_changes(vec![change]).unwrap_err();
+
+            assert!(
+                matches!(err, crate::Error::Unsupported { ref message }
+                    if message.contains("Cannot change ignore-delete from true 
to false")),
+                "got {err:?}"
+            );
+        }
+    }
+
     #[test]
     fn test_aggregation_apply_changes_accepts_valid_option() {
         let table_schema = TableSchema::new(
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 54ffa463..86873e8a 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -30,13 +30,13 @@ use crate::arrow::format::create_format_writer;
 use crate::io::FileIO;
 use crate::spec::stats::{compute_column_stats, BinaryTableStats};
 use crate::spec::{
-    extract_datum_from_arrow, AggregationConfig, BinaryRowBuilder, 
DataFileMeta, DataType,
-    MergeEngine, PartialUpdateConfig, RowKind, EMPTY_SERIALIZED_ROW, 
SEQUENCE_NUMBER_FIELD_NAME,
-    VALUE_KIND_FIELD_NAME,
+    extract_datum_from_arrow, AggregationConfig, BinaryRowBuilder, 
CoreOptions, DataFileMeta,
+    DataType, MergeEngine, PartialUpdateConfig, RowKind, EMPTY_SERIALIZED_ROW,
+    SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_NAME,
 };
 use crate::table::prepared_files::PreparedFiles;
 use crate::Result;
-use arrow_array::{Array, Int64Array, Int8Array, RecordBatch, UInt32Array};
+use arrow_array::{Array, BooleanArray, Int64Array, Int8Array, RecordBatch, 
UInt32Array};
 use arrow_ord::sort::{lexsort_to_indices, SortColumn, SortOptions};
 use arrow_row::{RowConverter, SortField};
 use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
@@ -49,6 +49,7 @@ use std::sync::Arc;
 pub(crate) struct KeyValueFileWriter {
     file_io: FileIO,
     config: KeyValueWriteConfig,
+    ignore_delete: bool,
     /// Next sequence number to assign (bucket-local, always auto-incremented).
     next_sequence_number: i64,
     /// Buffered batches (user schema).
@@ -105,6 +106,8 @@ impl KeyValueFileWriter {
         config: KeyValueWriteConfig,
         next_sequence_number: i64,
     ) -> Result<Self> {
+        let ignore_delete = config.merge_engine == MergeEngine::PartialUpdate
+            && CoreOptions::new(&config.table_options).ignore_delete();
         if config.merge_engine == MergeEngine::PartialUpdate {
             PartialUpdateConfig::new(&config.table_options)
                 .validate_runtime_mode(true, &config.table_name)?;
@@ -136,6 +139,7 @@ impl KeyValueFileWriter {
         Ok(Self {
             file_io,
             config,
+            ignore_delete,
             next_sequence_number,
             buffer: Vec::new(),
             buffer_bytes: 0,
@@ -147,6 +151,13 @@ impl KeyValueFileWriter {
     /// Buffer a RecordBatch. Flushes when buffer exceeds write_buffer_size.
     /// Sequence numbers are assigned per-bucket on flush, matching Java 
Paimon behavior.
     pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
+        // Filter before byte accounting and buffering so ignored retracts do
+        // not consume write-buffer memory or automatic sequence numbers.
+        let batch = if self.ignore_delete {
+            Self::filter_retract_rows(batch)?
+        } else {
+            batch.clone()
+        };
         if batch.num_rows() == 0 {
             return Ok(());
         }
@@ -155,7 +166,7 @@ impl KeyValueFileWriter {
             .iter()
             .map(|c| c.get_buffer_memory_size())
             .sum();
-        self.buffer.push(batch.clone());
+        self.buffer.push(batch);
         self.buffer_bytes += batch_bytes;
 
         if self.buffer_bytes as i64 >= self.config.write_buffer_size {
@@ -310,6 +321,44 @@ impl KeyValueFileWriter {
         Ok(())
     }
 
+    fn filter_retract_rows(batch: &RecordBatch) -> Result<RecordBatch> {
+        let Some(vk_idx) = batch
+            .schema()
+            .fields()
+            .iter()
+            .position(|field| field.name() == VALUE_KIND_FIELD_NAME)
+        else {
+            return Ok(batch.clone());
+        };
+        let value_kinds = batch
+            .column(vk_idx)
+            .as_any()
+            .downcast_ref::<Int8Array>()
+            .ok_or_else(|| crate::Error::DataInvalid {
+                message: "_VALUE_KIND column must be Int8".to_string(),
+                source: None,
+            })?;
+        let keep = BooleanArray::from(
+            (0..batch.num_rows())
+                .map(|row| {
+                    let value = if value_kinds.is_null(row) {
+                        RowKind::Insert.to_value()
+                    } else {
+                        value_kinds.value(row)
+                    };
+                    RowKind::from_value(value).map(|kind| kind.is_add())
+                })
+                .collect::<Result<Vec<_>>>()?,
+        );
+
+        arrow_select::filter::filter_record_batch(batch, &keep).map_err(|e| {
+            crate::Error::DataInvalid {
+                message: format!("Failed to filter ignored retract rows: {e}"),
+                source: None,
+            }
+        })
+    }
+
     async fn write_indexed_file(
         &self,
         batch: &RecordBatch,
@@ -551,7 +600,8 @@ impl KeyValueFileWriter {
     /// the read-side `PartialUpdateMergeFunction`: rows are visited in
     /// ascending (sequence fields, auto-seq) order and every column keeps its
     /// latest non-null value; a column that is null in every row stays null.
-    /// DELETE / UPDATE_BEFORE rows are rejected, matching the read side.
+    /// DELETE / UPDATE_BEFORE rows are rejected defensively. When
+    /// `ignore-delete=true`, `write` filters them before buffering.
     ///
     /// Returns the merged batch (user schema, in primary-key order) and its
     /// `_SEQUENCE_NUMBER` column; each merged row keeps the highest sequence
@@ -866,6 +916,138 @@ mod tests {
         .unwrap()
     }
 
+    #[tokio::test]
+    async fn 
test_flush_partial_update_ignore_delete_skips_retract_only_batch() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            Arc::new(ArrowField::new("id", ArrowDataType::Int32, false)),
+            Arc::new(ArrowField::new("seq", ArrowDataType::Int64, false)),
+            Arc::new(ArrowField::new(
+                VALUE_KIND_FIELD_NAME,
+                ArrowDataType::Int8,
+                false,
+            )),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1, 2])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int64Array::from(vec![10, 20])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int8Array::from(vec![1, 3])) as Arc<dyn 
arrow_array::Array>,
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::PartialUpdate);
+        config
+            .table_options
+            .insert("ignore-delete".to_string(), "true".to_string());
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
7)
+                .unwrap();
+
+        writer.write(&batch).await.unwrap();
+        assert!(writer.buffer.is_empty());
+        assert_eq!(writer.buffer_bytes, 0);
+        assert_eq!(writer.next_sequence_number, 7);
+
+        let prepared = writer.prepare_commit().await.unwrap();
+
+        assert!(prepared.data_files.is_empty());
+        assert!(prepared.changelog_files.is_empty());
+        assert_eq!(writer.next_sequence_number, 7);
+    }
+
+    #[tokio::test]
+    async fn 
test_flush_partial_update_ignore_delete_filters_data_and_changelog() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            Arc::new(ArrowField::new("id", ArrowDataType::Int32, false)),
+            Arc::new(ArrowField::new("seq", ArrowDataType::Int64, false)),
+            Arc::new(ArrowField::new(
+                VALUE_KIND_FIELD_NAME,
+                ArrowDataType::Int8,
+                false,
+            )),
+            Arc::new(ArrowField::new("value", ArrowDataType::Int32, true)),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1, 1, 2])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int64Array::from(vec![10, 20, 30])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int8Array::from(vec![0, 3, 1])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int32Array::from(vec![Some(100), None, None]))
+                    as Arc<dyn arrow_array::Array>,
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::PartialUpdate);
+        config.input_changelog = true;
+        config.table_options.insert(
+            "partial-update.ignore-delete".to_string(),
+            "true".to_string(),
+        );
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
7)
+                .unwrap();
+
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+
+        assert_eq!(prepared.data_files.len(), 1);
+        assert_eq!(prepared.changelog_files.len(), 1);
+        for file in prepared
+            .data_files
+            .iter()
+            .chain(prepared.changelog_files.iter())
+        {
+            assert_eq!(file.row_count, 1);
+            assert_eq!(file.min_sequence_number, 7);
+            assert_eq!(file.max_sequence_number, 7);
+            assert_eq!(file.delete_row_count, Some(0));
+        }
+        assert_eq!(writer.next_sequence_number, 8);
+    }
+
+    #[tokio::test]
+    async fn test_flush_partial_update_explicit_false_rejects_retract() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            Arc::new(ArrowField::new("id", ArrowDataType::Int32, false)),
+            Arc::new(ArrowField::new("seq", ArrowDataType::Int64, false)),
+            Arc::new(ArrowField::new(
+                VALUE_KIND_FIELD_NAME,
+                ArrowDataType::Int8,
+                false,
+            )),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![1])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int64Array::from(vec![10])) as Arc<dyn 
arrow_array::Array>,
+                Arc::new(Int8Array::from(vec![3])) as Arc<dyn 
arrow_array::Array>,
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::PartialUpdate);
+        config
+            .table_options
+            .insert("ignore-delete".to_string(), "false".to_string());
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
0)
+                .unwrap();
+
+        writer.write(&batch).await.unwrap();
+        let err = match writer.prepare_commit().await {
+            Ok(_) => panic!("explicit ignore-delete=false must reject retract 
rows"),
+            Err(err) => err,
+        };
+
+        assert!(matches!(
+            err,
+            crate::Error::Unsupported { message }
+            if message.contains("does not support DELETE or UPDATE_BEFORE")
+        ));
+    }
+
     /// Partial-update merges each key group down to one row at flush: every
     /// column keeps its latest non-null value (different columns may come
     /// from different source rows) and the merged row carries the group's
diff --git a/crates/paimon/src/table/sort_merge.rs 
b/crates/paimon/src/table/sort_merge.rs
index 523f9552..abbb6513 100644
--- a/crates/paimon/src/table/sort_merge.rs
+++ b/crates/paimon/src/table/sort_merge.rs
@@ -26,7 +26,7 @@
 //! - DataFusion: `SortPreservingMergeStream` (LoserTree layout)
 //! - Arrow-row: `RowConverter` for efficient key comparison
 
-use crate::spec::{AggregationConfig, DataField, PartialUpdateConfig, RowKind};
+use crate::spec::{AggregationConfig, CoreOptions, DataField, 
PartialUpdateConfig, RowKind};
 use crate::table::aggregator::{new_aggregator, FieldAggregator};
 use crate::table::ArrowRecordBatchStream;
 use crate::Error;
@@ -184,9 +184,12 @@ impl MergeFunction for DeduplicateMergeFunction {
 /// Basic partial-update merge: for each non-key column, keep the latest
 /// non-null value ordered by user sequence (if configured) then system 
sequence.
 ///
-/// DELETE / UPDATE_BEFORE rows are treated as unsupported in this mode.
+/// DELETE / UPDATE_BEFORE rows are ignored when `ignore-delete=true` and
+/// treated as unsupported otherwise.
 #[derive(Debug, Clone, Copy)]
-pub(crate) struct PartialUpdateMergeFunction(());
+pub(crate) struct PartialUpdateMergeFunction {
+    ignore_delete: bool,
+}
 
 impl PartialUpdateMergeFunction {
     pub(crate) fn new(
@@ -194,7 +197,9 @@ impl PartialUpdateMergeFunction {
         table_name: &str,
     ) -> crate::Result<Self> {
         PartialUpdateConfig::new(table_options).validate_runtime_mode(true, 
table_name)?;
-        Ok(Self(()))
+        Ok(Self {
+            ignore_delete: CoreOptions::new(table_options).ignore_delete(),
+        })
     }
 }
 
@@ -221,14 +226,19 @@ impl MergeFunction for PartialUpdateMergeFunction {
 
         let mut latest_non_null_by_col: Vec<Option<(usize, usize)>> =
             vec![None; output_schema.fields().len()];
+        let mut saw_add = false;
 
         for row_idx in ordered_row_indices {
             let row = &rows[row_idx];
             if !RowKind::from_value(row.value_kind)?.is_add() {
+                if self.ignore_delete {
+                    continue;
+                }
                 return Err(crate::Error::Unsupported {
                     message: "merge-engine=partial-update basic mode does not 
support DELETE or UPDATE_BEFORE rows".to_string(),
                 });
             }
+            saw_add = true;
 
             for (output_col_idx, latest_non_null) in 
latest_non_null_by_col.iter_mut().enumerate() {
                 let source_array = batch_buffer[row.batch_idx]
@@ -239,6 +249,10 @@ impl MergeFunction for PartialUpdateMergeFunction {
             }
         }
 
+        if !saw_add {
+            return Ok(MergeResult::Omit);
+        }
+
         let output_columns: Vec<ArrayRef> = output_schema
             .fields()
             .iter()
@@ -1911,6 +1925,158 @@ mod tests {
         ));
     }
 
+    #[tokio::test]
+    async fn test_partial_update_merge_ignores_delete_when_configured() {
+        let schema = make_schema();
+        let output_schema = make_output_schema();
+        let s0 = stream_from_batches(vec![make_batch_with_kind(
+            &schema,
+            vec![1],
+            vec![1],
+            vec![0],
+            vec![Some("old")],
+        )]);
+        let s1 = stream_from_batches(vec![make_batch_with_kind(
+            &schema,
+            vec![1],
+            vec![2],
+            vec![3],
+            vec![Some("delete")],
+        )]);
+        let options = HashMap::from([
+            ("merge-engine".to_string(), "partial-update".to_string()),
+            ("ignore-delete".to_string(), "true".to_string()),
+        ]);
+
+        let result = SortMergeReaderBuilder::new(
+            vec![s0, s1],
+            schema,
+            vec![0],
+            1,
+            2,
+            vec![],
+            vec![3],
+            output_schema,
+            Box::new(PartialUpdateMergeFunction::new(&options, 
"test_table").unwrap()),
+        )
+        .build()
+        .unwrap()
+        .try_collect::<Vec<_>>()
+        .await
+        .unwrap();
+
+        assert_eq!(result.len(), 1);
+        assert_eq!(result[0].num_rows(), 1);
+        assert_eq!(
+            result[0]
+                .column(1)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap()
+                .value(0),
+            "old"
+        );
+    }
+
+    #[tokio::test]
+    async fn test_partial_update_merge_alias_ignores_update_before() {
+        let schema = make_schema();
+        let output_schema = make_output_schema();
+        let s0 = stream_from_batches(vec![make_batch_with_kind(
+            &schema,
+            vec![1],
+            vec![1],
+            vec![0],
+            vec![Some("old")],
+        )]);
+        let s1 = stream_from_batches(vec![make_batch_with_kind(
+            &schema,
+            vec![1],
+            vec![2],
+            vec![1],
+            vec![Some("before")],
+        )]);
+        let s2 = stream_from_batches(vec![make_batch_with_kind(
+            &schema,
+            vec![1],
+            vec![3],
+            vec![2],
+            vec![Some("new")],
+        )]);
+        let options = HashMap::from([
+            ("merge-engine".to_string(), "partial-update".to_string()),
+            (
+                "partial-update.ignore-delete".to_string(),
+                "true".to_string(),
+            ),
+        ]);
+
+        let result = SortMergeReaderBuilder::new(
+            vec![s0, s1, s2],
+            schema,
+            vec![0],
+            1,
+            2,
+            vec![],
+            vec![3],
+            output_schema,
+            Box::new(PartialUpdateMergeFunction::new(&options, 
"test_table").unwrap()),
+        )
+        .build()
+        .unwrap()
+        .try_collect::<Vec<_>>()
+        .await
+        .unwrap();
+
+        assert_eq!(result.len(), 1);
+        assert_eq!(result[0].num_rows(), 1);
+        assert_eq!(
+            result[0]
+                .column(1)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap()
+                .value(0),
+            "new"
+        );
+    }
+
+    #[tokio::test]
+    async fn 
test_partial_update_merge_omits_retract_only_key_when_configured() {
+        let schema = make_schema();
+        let output_schema = make_output_schema();
+        let stream = stream_from_batches(vec![make_batch_with_kind(
+            &schema,
+            vec![1],
+            vec![1],
+            vec![3],
+            vec![Some("delete")],
+        )]);
+        let options = HashMap::from([
+            ("merge-engine".to_string(), "partial-update".to_string()),
+            ("ignore-delete".to_string(), "true".to_string()),
+        ]);
+
+        let result = SortMergeReaderBuilder::new(
+            vec![stream],
+            schema,
+            vec![0],
+            1,
+            2,
+            vec![],
+            vec![3],
+            output_schema,
+            Box::new(PartialUpdateMergeFunction::new(&options, 
"test_table").unwrap()),
+        )
+        .build()
+        .unwrap()
+        .try_collect::<Vec<_>>()
+        .await
+        .unwrap();
+
+        assert_eq!(result.iter().map(RecordBatch::num_rows).sum::<usize>(), 0);
+    }
+
     #[test]
     fn test_partial_update_merge_function_new_rejects_unsupported_options() {
         let options = HashMap::from([
diff --git a/crates/paimon/tests/incremental_batch_scan_test.rs 
b/crates/paimon/tests/incremental_batch_scan_test.rs
index c1a0a8c2..bc007851 100644
--- a/crates/paimon/tests/incremental_batch_scan_test.rs
+++ b/crates/paimon/tests/incremental_batch_scan_test.rs
@@ -69,6 +69,19 @@ async fn read_incremental_pairs(
     collect_pairs(&batches)
 }
 
+async fn read_current_pairs(table: &paimon::table::Table) -> Vec<(i32, i32)> {
+    let builder = table.new_read_builder();
+    let plan = builder.new_scan().plan().await.unwrap();
+    let read = table.new_read_builder().new_read().unwrap();
+    let batches: Vec<RecordBatch> = read
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    collect_pairs(&batches)
+}
+
 async fn plan_incremental(
     table: &paimon::table::Table,
     mode: IncrementalScanMode,
@@ -269,6 +282,34 @@ async fn 
changelog_between_snapshots_reads_changelog_manifest_files() {
     assert_eq!(rows, vec![(1, 10), (1, 20)]);
 }
 
+#[tokio::test]
+async fn partial_update_ignore_delete_filters_data_and_input_changelog() {
+    let table_path = "memory:/incremental_batch/partial_update_ignore_delete";
+    let (file_io, table) = memory_table(
+        table_path,
+        pk_schema(&[
+            ("changelog-producer", "input"),
+            ("merge-engine", "partial-update"),
+            ("partial-update.ignore-delete", "true"),
+            ("bucket", "1"),
+        ]),
+    );
+    setup_dirs(&file_io, table_path).await;
+    persist_table_schema(&file_io, table_path, table.schema()).await;
+
+    write_batch(
+        &table,
+        &make_batch_with_kinds(vec![1, 1, 2, 1], vec![10, 999, 200, 20], 
vec![0, 3, 1, 2]),
+    )
+    .await;
+
+    assert_eq!(read_current_pairs(&table).await, vec![(1, 20)]);
+    assert_eq!(
+        read_incremental_pairs(&table, IncrementalScanMode::Changelog, 0, 
1).await,
+        vec![(1, 10), (1, 20)]
+    );
+}
+
 /// Multi-snapshot changelog range is left-open / right-closed and ordered by 
snapshot id.
 #[tokio::test]
 async fn changelog_multi_snapshot_range_is_ordered_and_left_open() {
diff --git a/docs/src/sql.md b/docs/src/sql.md
index 0cd5e255..ac3d0ee9 100644
--- a/docs/src/sql.md
+++ b/docs/src/sql.md
@@ -1880,6 +1880,15 @@ Set via `WITH ('key' = 'value')` at table creation time, 
or dynamically via `SET
 | `'merge-engine' = 'partial-update'` | Basic partial-update engine for PK 
tables |
 | `'merge-engine' = 'aggregation'` | Basic aggregation engine for PK tables |
 
+Rust supports the basic partial-update engine with latest-non-null semantics.
+Set either `'ignore-delete' = 'true'` or
+`'partial-update.ignore-delete' = 'true'` to ignore `DELETE` and
+`UPDATE_BEFORE` rows during writes and when reading existing files. The default
+and an explicit `false` continue to reject these retract rows. Once enabled on
+an existing partial-update table, `ignore-delete` cannot be changed back to
+`false`. Advanced partial-update features such as sequence groups, partial
+aggregation, and remove-record-on-delete are not supported.
+
 Rust currently supports `merge-engine=aggregation` in basic mode only. It works
 with fixed buckets and ordinary dynamic buckets (`'bucket' = '-1'`) when the
 primary key includes all partition columns. It supports per-field aggregate


Reply via email to