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