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 50725c84 fix(core): complete native update partition and deletion 
vector paths (#969)
50725c84 is described below

commit 50725c8451e399ea4a4085dfa7f4ef35435a3c9b
Author: Jingsong Lee <[email protected]>
AuthorDate: Sat Sep 26 21:41:28 2026 +0800

    fix(core): complete native update partition and deletion vector paths (#969)
---
 crates/paimon/Cargo.toml                          |   2 +-
 crates/paimon/src/spec/murmur_hash.rs             |  18 +
 crates/paimon/src/table/data_evolution_writer.rs  | 121 ++++-
 crates/paimon/src/table/external_path.rs          | 294 ++++++++++++
 crates/paimon/src/table/index_file_path.rs        |   2 +-
 crates/paimon/src/table/mod.rs                    |   1 +
 crates/paimon/src/table/table_update.rs           |   6 +-
 crates/paimon/src/table/table_update_by_row_id.rs |   4 +-
 crates/paimon/src/table/table_upsert.rs           |  22 +-
 crates/paimon/tests/table_update_paths_test.rs    | 521 ++++++++++++++++++++++
 10 files changed, 960 insertions(+), 31 deletions(-)

diff --git a/crates/paimon/Cargo.toml b/crates/paimon/Cargo.toml
index 0d408ca5..7c971dcd 100644
--- a/crates/paimon/Cargo.toml
+++ b/crates/paimon/Cargo.toml
@@ -63,6 +63,7 @@ storage-gcs = ["dep:opendal-http-transport-reqwest", 
"dep:opendal-service-gcs"]
 storage-hdfs = ["dep:opendal-service-hdfs-native"]
 
 [dependencies]
+rand = "0.8.5"
 url = "2.5.2"
 # Already in the tree via `url`; direct dep for `RESTUtil::decode_string`.
 percent-encoding = "2.3"
@@ -151,4 +152,3 @@ unicode-segmentation = "=1.13.2"
 
 [dev-dependencies]
 axum = { version = "0.7", features = ["macros", "tokio", "http1", "http2"] }
-rand = "0.8.5"
diff --git a/crates/paimon/src/spec/murmur_hash.rs 
b/crates/paimon/src/spec/murmur_hash.rs
index 08f228ce..70c35ba5 100644
--- a/crates/paimon/src/spec/murmur_hash.rs
+++ b/crates/paimon/src/spec/murmur_hash.rs
@@ -86,6 +86,24 @@ pub(crate) fn hash_bytes(data: &[u8]) -> i32 {
     fmix(h1 ^ data.len() as u32) as i32
 }
 
+/// Guava's canonical Murmur3_32 with seed zero, used by Java's external
+/// entropy paths. Unlike Paimon row hashing, the tail forms a single word.
+pub(crate) fn hash_bytes_guava(data: &[u8]) -> i32 {
+    let (words, tail) = data.as_chunks::<4>();
+    let mut hash = 0;
+    for word in words {
+        hash = mix_h1(hash, mix_k1(u32::from_le_bytes(*word)));
+    }
+    let mut last = 0;
+    for (index, byte) in tail.iter().enumerate() {
+        last |= u32::from(*byte) << (8 * index);
+    }
+    if !tail.is_empty() {
+        hash ^= mix_k1(last);
+    }
+    fmix(hash ^ data.len() as u32) as i32
+}
+
 #[cfg(test)]
 mod tests {
     use super::*;
diff --git a/crates/paimon/src/table/data_evolution_writer.rs 
b/crates/paimon/src/table/data_evolution_writer.rs
index b4d1b6bd..f070fe4e 100644
--- a/crates/paimon/src/table/data_evolution_writer.rs
+++ b/crates/paimon/src/table/data_evolution_writer.rs
@@ -50,7 +50,7 @@ use bytes::Bytes;
 use futures::TryStreamExt;
 use indexmap::IndexMap;
 use roaring::RoaringBitmap;
-use std::collections::{HashMap, HashSet};
+use std::collections::{BTreeMap, HashMap, HashSet};
 use std::sync::Arc;
 use uuid::Uuid;
 
@@ -92,6 +92,21 @@ impl DataEvolutionWriter {
     /// - No primary keys
     /// - Update columns don't include partition keys
     pub fn new(table: &Table, update_columns: Vec<String>) -> Result<Self> {
+        Self::with_partition_columns(table, update_columns, false)
+    }
+
+    /// Row-ID inputs may carry unchanged partition values, as complete-row
+    /// upserts do. Validate their values against the pinned file index before
+    /// writing; changing a partition requires delete + insert.
+    pub(super) fn for_row_id(table: &Table, update_columns: Vec<String>) -> 
Result<Self> {
+        Self::with_partition_columns(table, update_columns, true)
+    }
+
+    fn with_partition_columns(
+        table: &Table,
+        update_columns: Vec<String>,
+        allow_partition_columns: bool,
+    ) -> Result<Self> {
         let schema = table.schema();
         let core_options = CoreOptions::new(schema.options());
 
@@ -136,7 +151,7 @@ impl DataEvolutionWriter {
         let blob_descriptor_fields = core_options.blob_descriptor_fields();
         for col in &update_columns {
             let top_level = 
DataEvolutionPartialWriter::top_level_write_name(col, schema.fields());
-            if partition_keys.iter().any(|key| key == top_level) {
+            if !allow_partition_columns && partition_keys.iter().any(|key| key 
== top_level) {
                 return Err(crate::Error::Unsupported {
                     message: format!("Cannot update partition column '{col}' 
in MERGE INTO"),
                 });
@@ -164,6 +179,51 @@ impl DataEvolutionWriter {
         })
     }
 
+    fn validate_partition_values(
+        &self,
+        files: &[FileRowRange],
+        matches: &HashMap<usize, Vec<MatchedRow>>,
+    ) -> Result<()> {
+        for (partition_index, name) in 
self.table.schema().partition_keys().iter().enumerate() {
+            let Some(field) = self.write_fields.iter().find(|field| 
field.name() == name) else {
+                continue;
+            };
+            let target_type = 
crate::arrow::paimon_type_to_arrow(field.data_type())?;
+            let values = self
+                .matched_batches
+                .iter()
+                .map(|batch| {
+                    super::update_input::cast_update_value(
+                        &matched_column(batch, name)?,
+                        &target_type,
+                        super::update_input::CastMode::RowUpdate,
+                    )
+                })
+                .collect::<Result<Vec<_>>>()?;
+            for (&file_pos, rows) in matches {
+                let partition = 
BinaryRow::from_serialized_bytes(&files[file_pos].partition)?;
+                let expected = crate::arrow::partition::partition_array(
+                    &partition,
+                    partition_index,
+                    field.data_type(),
+                    1,
+                )?
+                .to_data();
+                for row in rows {
+                    if values[row.batch_idx].slice(row.row_idx, 1).to_data() 
!= expected {
+                        return Err(crate::Error::DataInvalid {
+                            message: format!(
+                                "Cannot change partition column '{name}' in a 
row-ID update"
+                            ),
+                            source: None,
+                        });
+                    }
+                }
+            }
+        }
+        Ok(())
+    }
+
     /// Add a batch of matched rows.
     ///
     /// The batch must contain:
@@ -277,6 +337,7 @@ impl DataEvolutionWriter {
             &self.update_columns,
         )?;
         let file_matches = group_matched_rows_by_file(&self.matched_batches, 
file_index)?;
+        self.validate_partition_values(file_index, &file_matches)?;
 
         // 3. For each affected file: read original columns, apply updates, 
write partial files
         let mut writer = DataEvolutionPartialWriter::new(&self.table, 
self.update_columns.clone())?;
@@ -545,7 +606,7 @@ impl DataEvolutionDeleteWriter {
             });
         }
 
-        let mut deletes_by_bucket: HashMap<(Vec<u8>, i32), BucketDeletePlan> = 
HashMap::new();
+        let mut deletes_by_bucket: BTreeMap<(Vec<u8>, i32), BucketDeletePlan> 
= BTreeMap::new();
         for row_id in &self.row_ids {
             let (file_pos, file_range) =
                 find_delete_owning_file(&file_index, *row_id).ok_or_else(|| {
@@ -580,11 +641,21 @@ impl DataEvolutionDeleteWriter {
 
         let mut messages = Vec::new();
         for ((partition, bucket), delete_plan) in deletes_by_bucket {
-            if let Some(message) = self
+            match self
                 .prepare_bucket_delete_message(partition, bucket, delete_plan, 
&snapshot)
-                .await?
+                .await
             {
-                messages.push(message);
+                Ok(Some(message)) => messages.push(message),
+                Ok(None) => {}
+                Err(error) => {
+                    let _ = self
+                        .table
+                        .new_write_builder()
+                        .new_commit()
+                        .abort(&messages)
+                        .await;
+                    return Err(error);
+                }
             }
         }
 
@@ -683,7 +754,7 @@ impl DataEvolutionDeleteWriter {
         let mut bitmaps = IndexMap::new();
         let mut deleted_index_files = Vec::new();
 
-        for entry in index_entries {
+        for mut entry in index_entries {
             if entry.kind != FileKind::Add
                 || entry.bucket != bucket
                 || entry.partition != partition
@@ -692,6 +763,12 @@ impl DataEvolutionDeleteWriter {
                 continue;
             }
             deleted_index_files.push(entry.index_file.clone());
+            // Preserve the original manifest identity in deleted_index_files.
+            // Only the copy used for reading acquires the legacy physical 
path.
+            layout
+                .location()
+                .resolve_legacy_deletion_vector(self.table.file_io(), &mut 
entry.index_file)
+                .await?;
             let Some(ranges) = 
entry.index_file.deletion_vectors_ranges.as_ref() else {
                 continue;
             };
@@ -725,12 +802,21 @@ impl DataEvolutionDeleteWriter {
         bitmaps.sort_keys();
 
         let file_name = format!("index-{}-1", Uuid::new_v4());
-        // Write where the reader resolves it, so a deletion vector written 
here is
-        // found again on the next scan.
         let layout = self.deletion_vector_layout(partition, bucket)?;
-        let location = layout.location();
-        let path = location.resolve(&file_name, None);
-        self.table.file_io().mkdirs(&location.directory()).await?;
+        let external_path = super::external_path::new_index_external_path(
+            self.table.schema().options(),
+            layout.index_file_in_data_file_dir,
+            layout
+                .bucket_path
+                .strip_prefix(&format!("{}/", layout.table_path))
+                .expect("bucket path is under the table"),
+            &file_name,
+        )?;
+        let path = layout
+            .location()
+            .resolve(&file_name, external_path.as_deref());
+        let directory = path.rsplit_once('/').expect("index file has a 
parent").0;
+        self.table.file_io().mkdirs(directory).await?;
 
         let mut bytes = vec![DELETION_VECTORS_INDEX_VERSION_V1];
         let mut ranges = IndexMap::new();
@@ -768,11 +854,16 @@ impl DataEvolutionDeleteWriter {
             message: "Deletion-vector index file has too many 
entries".to_string(),
             source: None,
         })?;
-        self.table
+        if let Err(error) = self
+            .table
             .file_io()
             .new_output(&path)?
             .write(Bytes::from(bytes))
-            .await?;
+            .await
+        {
+            let _ = self.table.file_io().delete_file(&path).await;
+            return Err(error);
+        }
 
         Ok(IndexFileMeta {
             index_type: DELETION_VECTORS_INDEX_TYPE.to_string(),
@@ -780,7 +871,7 @@ impl DataEvolutionDeleteWriter {
             file_size,
             row_count: i64::from(row_count),
             deletion_vectors_ranges: Some(ranges),
-            external_path: None,
+            external_path,
             global_index_meta: None,
         })
     }
diff --git a/crates/paimon/src/table/external_path.rs 
b/crates/paimon/src/table/external_path.rs
new file mode 100644
index 00000000..2b694d15
--- /dev/null
+++ b/crates/paimon/src/table/external_path.rs
@@ -0,0 +1,294 @@
+// 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.
+
+//! External index placement follows Java's FileStorePathFactory. Existing
+//! files always resolve using their recorded path, never current write 
options.
+
+use crate::Result;
+use rand::Rng;
+use std::collections::HashMap;
+
+fn invalid(message: impl Into<String>) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.into(),
+        source: None,
+    }
+}
+
+pub(super) fn new_index_external_path(
+    options: &HashMap<String, String>,
+    in_data_directory: bool,
+    relative_bucket: &str,
+    file_name: &str,
+) -> Result<Option<String>> {
+    if !in_data_directory {
+        return Ok(options
+            .get("global-index.external-path")
+            .map(|path| format!("{}/{file_name}", 
path.trim_end_matches('/'))));
+    }
+    Ok(ExternalPathProvider::new(options, relative_bucket)?
+        .map(|mut provider| provider.next_path(file_name)))
+}
+
+/// Mirrors Java's per-bucket ExternalPathProvider. The random starting point
+/// avoids concentrating single-file buckets on the first configured root.
+struct ExternalPathProvider {
+    paths: Vec<String>,
+    bucket: String,
+    position: usize,
+    entropy: bool,
+    cumulative_weights: Vec<u64>,
+}
+
+impl ExternalPathProvider {
+    fn new(options: &HashMap<String, String>, bucket: &str) -> 
Result<Option<Self>> {
+        let strategy = options
+            .get("data-file.external-paths.strategy")
+            .map(|value| value.to_ascii_lowercase())
+            .unwrap_or_else(|| "none".into());
+        let Some(paths) = options
+            .get("data-file.external-paths")
+            .filter(|paths| !paths.is_empty())
+        else {
+            return Ok(None);
+        };
+        if strategy == "none" {
+            return Ok(None);
+        }
+        if !matches!(
+            strategy.as_str(),
+            "round-robin" | "specific-fs" | "weight-robin" | "entropy-inject"
+        ) {
+            return Err(invalid(format!(
+                "Unsupported external path strategy: {strategy}"
+            )));
+        }
+        let specific_fs = if strategy == "specific-fs" {
+            Some(
+                options
+                    .get("data-file.external-paths.specific-fs")
+                    .ok_or_else(|| invalid("External path specific-fs is 
required"))?,
+            )
+        } else {
+            None
+        };
+        let mut roots = Vec::new();
+        // Java String.split discards trailing empty entries.
+        for path in paths.trim_end_matches(',').split(',').map(str::trim) {
+            let uri = url::Url::parse(path)
+                .map_err(|_| invalid(format!("External path must have a URI 
scheme: {path}")))?;
+            if specific_fs.is_none_or(|scheme| 
uri.scheme().eq_ignore_ascii_case(scheme)) {
+                roots.push(path.trim_end_matches('/').to_string());
+            }
+        }
+        if roots.is_empty() {
+            return Err(invalid("External paths should not be empty"));
+        }
+        let mut cumulative_weights = Vec::new();
+        if strategy == "weight-robin" && roots.len() > 1 {
+            if let Some(weights) = options
+                .get("data-file.external-paths.weights")
+                .filter(|weights| !weights.trim().is_empty())
+            {
+                let mut total = 0_u64;
+                for weight in weights.trim_end_matches(',').split(',') {
+                    let weight = weight
+                        .trim()
+                        .parse::<i32>()
+                        .ok()
+                        .filter(|weight| *weight > 0)
+                        .ok_or_else(|| {
+                            invalid("External path weights must be positive 
integers")
+                        })?;
+                    total = total
+                        .checked_add(weight as u64)
+                        .ok_or_else(|| invalid("External path weight 
overflow"))?;
+                    cumulative_weights.push(total);
+                }
+                if cumulative_weights.len() != roots.len() {
+                    return Err(invalid(
+                        "The number of external paths and weights should be 
the same",
+                    ));
+                }
+            }
+        }
+        let entropy = strategy == "entropy-inject";
+        let position = if entropy {
+            0
+        } else {
+            rand::thread_rng().gen_range(0..roots.len())
+        };
+        Ok(Some(Self {
+            paths: roots,
+            bucket: bucket.trim_matches('/').into(),
+            position,
+            entropy,
+            cumulative_weights,
+        }))
+    }
+
+    fn next_path(&mut self, file_name: &str) -> String {
+        let index = if let Some(total) = self.cumulative_weights.last() {
+            let value = rand::thread_rng().gen_range(0..*total);
+            self.cumulative_weights
+                .partition_point(|weight| *weight <= value)
+        } else {
+            self.position = (self.position + 1) % self.paths.len();
+            self.position
+        };
+        let mut path = self.paths[index].clone();
+        if !self.bucket.is_empty() {
+            path.push('/');
+            path.push_str(&self.bucket);
+        }
+        if self.entropy {
+            let hash =
+                
crate::spec::murmur_hash::hash_bytes_guava(file_name.as_bytes()) as u32 & 
0xfffff;
+            path.push_str(&format!(
+                "/{:04b}/{:04b}/{:04b}/{:08b}",
+                hash >> 16,
+                (hash >> 12) & 15,
+                (hash >> 8) & 15,
+                hash & 255
+            ));
+        }
+        format!("{path}/{file_name}")
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+
+    fn options(strategy: &str) -> HashMap<String, String> {
+        HashMap::from([
+            (
+                "data-file.external-paths".into(),
+                "file:///a/, file:///b".into(),
+            ),
+            ("data-file.external-paths.strategy".into(), strategy.into()),
+            (
+                "global-index.external-path".into(),
+                "file:///global/".into(),
+            ),
+        ])
+    }
+
+    #[test]
+    fn round_robin_and_location_precedence() {
+        let settings = options("round-robin");
+        let mut provider = ExternalPathProvider::new(&settings, "p=x/bucket-0")
+            .unwrap()
+            .unwrap();
+        let paths = [provider.next_path("index-1"), 
provider.next_path("index-1")];
+        assert!(paths.contains(&"file:///a/p=x/bucket-0/index-1".into()));
+        assert!(paths.contains(&"file:///b/p=x/bucket-0/index-1".into()));
+        assert_eq!(
+            new_index_external_path(&settings, false, "p=x/bucket-0", 
"index-1")
+                .unwrap()
+                .as_deref(),
+            Some("file:///global/index-1")
+        );
+        assert!(
+            new_index_external_path(&options("none"), true, "bucket-0", 
"index-1")
+                .unwrap()
+                .is_none()
+        );
+    }
+
+    #[test]
+    fn specific_fs_and_invalid_options() {
+        let mut options = options("specific-fs");
+        assert!(ExternalPathProvider::new(&options, "").is_err());
+        options.insert("data-file.external-paths.specific-fs".into(), 
"FILE".into());
+        options.insert(
+            "data-file.external-paths".into(),
+            "s3://bucket/path,file:///a".into(),
+        );
+        assert_eq!(
+            ExternalPathProvider::new(&options, "bucket-0")
+                .unwrap()
+                .unwrap()
+                .next_path("index-1"),
+            "file:///a/bucket-0/index-1"
+        );
+        options.insert("data-file.external-paths.specific-fs".into(), 
"oss".into());
+        assert!(ExternalPathProvider::new(&options, "").is_err());
+        options.insert("data-file.external-paths".into(), "/no/scheme".into());
+        assert!(ExternalPathProvider::new(&options, "").is_err());
+    }
+
+    #[test]
+    fn weighted_validation_and_fallback() {
+        let mut options = options("weight-robin");
+        assert!(ExternalPathProvider::new(&options, "")
+            .unwrap()
+            .unwrap()
+            .cumulative_weights
+            .is_empty());
+        for weights in ["0,1", "-1,2", "x,2", "1", "2147483648,1"] {
+            options.insert("data-file.external-paths.weights".into(), 
weights.into());
+            assert!(
+                ExternalPathProvider::new(&options, "").is_err(),
+                "{weights}"
+            );
+        }
+        options.insert("data-file.external-paths.weights".into(), 
"1,2".into());
+        let mut provider = ExternalPathProvider::new(&options, "bucket-0")
+            .unwrap()
+            .unwrap();
+        assert_eq!(provider.cumulative_weights, vec![1, 3]);
+        for _ in 0..10 {
+            assert!(["file:///a/bucket-0/index", "file:///b/bucket-0/index"]
+                .contains(&provider.next_path("index").as_str()));
+        }
+    }
+
+    #[test]
+    fn comma_lists_follow_java_trailing_empty_semantics() {
+        let mut options = options("weight-robin");
+        options.insert(
+            "data-file.external-paths".into(),
+            "file:///a,file:///b,,".into(),
+        );
+        options.insert("data-file.external-paths.weights".into(), 
"1,2,,".into());
+        let provider = ExternalPathProvider::new(&options, 
"").unwrap().unwrap();
+        assert_eq!(provider.paths, vec!["file:///a", "file:///b"]);
+        assert_eq!(provider.cumulative_weights, vec![1, 3]);
+        options.insert(
+            "data-file.external-paths".into(),
+            "file:///a,,file:///b".into(),
+        );
+        assert!(ExternalPathProvider::new(&options, "").is_err());
+    }
+
+    #[test]
+    fn entropy_uses_guava_hash_and_rotates_from_second_root() {
+        let mut provider = 
ExternalPathProvider::new(&options("entropy-inject"), "p=x/bucket-0")
+            .unwrap()
+            .unwrap();
+        // Guava murmur3_32(0), UTF-8 "hello": 0x248bfa47.
+        assert_eq!(
+            provider.next_path("hello"),
+            "file:///b/p=x/bucket-0/1011/1111/1010/01000111/hello"
+        );
+        assert_eq!(
+            provider.next_path("hello"),
+            "file:///a/p=x/bucket-0/1011/1111/1010/01000111/hello"
+        );
+    }
+}
diff --git a/crates/paimon/src/table/index_file_path.rs 
b/crates/paimon/src/table/index_file_path.rs
index d96743fc..a4a3cbd4 100644
--- a/crates/paimon/src/table/index_file_path.rs
+++ b/crates/paimon/src/table/index_file_path.rs
@@ -83,7 +83,7 @@ impl IndexFileLocation<'_> {
     /// Older Python DV writers ignored the bucket-directory option. Resolve
     /// their existing files without changing the manifest or masking missing
     /// canonical files with a path that does not exist either.
-    async fn resolve_legacy_deletion_vector(
+    pub(super) async fn resolve_legacy_deletion_vector(
         &self,
         file_io: &FileIO,
         file: &mut IndexFileMeta,
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 1b978643..c195deda 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -47,6 +47,7 @@ mod data_file_writer;
 mod de_vector_read;
 mod de_vector_scan;
 mod dedicated_format_file_writer;
+mod external_path;
 mod format_partition;
 mod format_partition_location;
 mod format_partition_stats;
diff --git a/crates/paimon/src/table/table_update.rs 
b/crates/paimon/src/table/table_update.rs
index c9cb3fc1..444c7064 100644
--- a/crates/paimon/src/table/table_update.rs
+++ b/crates/paimon/src/table/table_update.rs
@@ -44,7 +44,7 @@ fn invalid(message: impl Into<String>) -> crate::Error {
 ///
 /// Row-ID updates infer columns from the input unless configured with
 /// [`with_update_type`](Self::with_update_type). Key-based upserts currently
-/// require complete Arrow rows and an unpartitioned table.
+/// require complete Arrow rows and match keys within each partition.
 #[derive(Clone)]
 pub struct TableUpdate {
     table: Table,
@@ -200,7 +200,9 @@ impl TableUpdate {
     }
 
     /// Upsert complete Arrow rows by composite key through the core upsert
-    /// writer. Existing keys update every matching row ID; new keys append.
+    /// writer. Keys match within each partition, including when partition
+    /// columns are omitted from `upsert_keys`. Existing keys update every
+    /// matching row ID; new keys append.
     pub async fn upsert_by_arrow_with_key(
         &self,
         batches: Vec<RecordBatch>,
diff --git a/crates/paimon/src/table/table_update_by_row_id.rs 
b/crates/paimon/src/table/table_update_by_row_id.rs
index cb5e0322..6d79e8de 100644
--- a/crates/paimon/src/table/table_update_by_row_id.rs
+++ b/crates/paimon/src/table/table_update_by_row_id.rs
@@ -64,6 +64,8 @@ impl TableUpdateByRowId {
     }
 
     /// Stage one logical Arrow table. Its chunks may share a file group.
+    /// Partition columns may carry their existing values. Changing them is
+    /// rejected before writing, because a move requires delete + insert.
     pub async fn update_columns(
         &mut self,
         batches: Vec<RecordBatch>,
@@ -77,7 +79,7 @@ impl TableUpdateByRowId {
             .into_iter()
             .filter(|name| seen.insert(name.clone()))
             .collect::<Vec<_>>();
-        let mut writer = DataEvolutionWriter::new(&self.table, 
columns.clone())?;
+        let mut writer = DataEvolutionWriter::for_row_id(&self.table, 
columns.clone())?;
         let batches = batches
             .into_iter()
             .map(super::update_input::normalize_row_ids)
diff --git a/crates/paimon/src/table/table_upsert.rs 
b/crates/paimon/src/table/table_upsert.rs
index c626816f..3a62311a 100644
--- a/crates/paimon/src/table/table_upsert.rs
+++ b/crates/paimon/src/table/table_upsert.rs
@@ -71,16 +71,18 @@ impl TableUpsert {
     pub(super) fn new(
         table: &Table,
         commit_user: String,
-        keys: Vec<String>,
+        mut keys: Vec<String>,
         update_columns: Vec<String>,
     ) -> crate::Result<Self> {
         if keys.is_empty() {
             return Err(invalid("upsert keys must not be empty"));
         }
-        if !table.schema().partition_keys().is_empty() {
-            return Err(crate::Error::Unsupported {
-                message: "native upsert currently requires an unpartitioned 
table".to_string(),
-            });
+        // PyPaimon matches independently within each source partition, even
+        // when callers omit partition columns from their upsert keys.
+        for partition in table.schema().partition_keys() {
+            if !keys.contains(partition) {
+                keys.push(partition.clone());
+            }
         }
         let fields = table.schema().fields();
         for key in &keys {
@@ -105,7 +107,8 @@ impl TableUpsert {
             return Err(invalid("upsert update columns must not be empty"));
         }
         // Reuse the row-ID writer's precondition and column-path checks.
-        let _validated_update = super::DataEvolutionWriter::new(table, 
update_columns.clone())?;
+        let _validated_update =
+            super::DataEvolutionWriter::for_row_id(table, 
update_columns.clone())?;
         Ok(Self {
             table: table.clone(),
             commit_user,
@@ -182,11 +185,8 @@ impl TableUpsert {
                     .map_err(|error| {
                         invalid(format!("cannot build matched upsert rows: 
{error}"))
                     })?;
-                let mut update = self
-                    .table
-                    .new_write_builder()
-                    .with_commit_user(self.commit_user.clone())?
-                    .new_data_evolution_writer(self.update_columns)?;
+                let mut update =
+                    super::DataEvolutionWriter::for_row_id(&self.table, 
self.update_columns)?;
                 if let Some(snapshot_id) = plan.snapshot_id() {
                     update.pin_read_snapshot(snapshot_id);
                 }
diff --git a/crates/paimon/tests/table_update_paths_test.rs 
b/crates/paimon/tests/table_update_paths_test.rs
new file mode 100644
index 00000000..2c7590d0
--- /dev/null
+++ b/crates/paimon/tests/table_update_paths_test.rs
@@ -0,0 +1,521 @@
+// 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.
+
+mod common;
+
+use arrow_array::{Array, ArrayRef, Int32Array, Int64Array, RecordBatch, 
StringArray};
+use common::incremental_helpers::{memory_table, persist_table_schema, 
setup_dirs, write_batch};
+use futures::TryStreamExt;
+use paimon::spec::{DataType, IntType, Schema, TableSchema, VarCharType};
+use paimon::table::{CommitMessage, DataEvolutionWriter, Table};
+use std::collections::HashMap;
+use std::sync::Arc;
+
+async fn table(options: &[(&str, &str)]) -> Table {
+    let mut schema = Schema::builder()
+        .column("p", DataType::VarChar(VarCharType::string_type()))
+        .column("q", DataType::Int(IntType::new()))
+        .column("id", DataType::Int(IntType::new()))
+        .column("v", DataType::Int(IntType::new()))
+        .partition_keys(["p", "q"])
+        .option("row-tracking.enabled", "true")
+        .option("data-evolution.enabled", "true")
+        .option("deletion-vectors.enabled", "true");
+    for (key, value) in options {
+        schema = schema.option(*key, *value);
+    }
+    let schema = schema.build().unwrap();
+    let path = "memory:/update_paths";
+    let (io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&io, path).await;
+    persist_table_schema(&io, path, table.schema()).await;
+    table
+}
+
+fn batch(p: Vec<Option<&str>>, q: Vec<i32>, id: Vec<i32>, v: Vec<i32>) -> 
RecordBatch {
+    RecordBatch::try_from_iter([
+        ("p", Arc::new(StringArray::from(p)) as ArrayRef),
+        ("q", Arc::new(Int32Array::from(q)) as ArrayRef),
+        ("id", Arc::new(Int32Array::from(id)) as ArrayRef),
+        ("v", Arc::new(Int32Array::from(v)) as ArrayRef),
+    ])
+    .unwrap()
+}
+
+async fn seed(table: &Table) {
+    write_batch(
+        table,
+        &batch(
+            vec![Some("a"), Some("a"), None, None],
+            vec![1, 1, 2, 2],
+            vec![0, 1, 2, 3],
+            vec![10, 11, 12, 13],
+        ),
+    )
+    .await;
+}
+
+async fn read(table: &Table) -> Vec<RecordBatch> {
+    let mut read = table.new_read_builder();
+    read.with_projection(&["p", "q", "id", "v", "_ROW_ID"])
+        .unwrap();
+    let plan = read.new_scan().plan().await.unwrap();
+    read.new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap()
+}
+
+async fn row_ids(table: &Table) -> HashMap<i32, i64> {
+    let mut result = HashMap::new();
+    for batch in read(table).await {
+        let ids = batch
+            .column_by_name("id")
+            .unwrap()
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let row_ids = batch
+            .column_by_name("_ROW_ID")
+            .unwrap()
+            .as_any()
+            .downcast_ref::<Int64Array>()
+            .unwrap();
+        for row in 0..batch.num_rows() {
+            result.insert(ids.value(row), row_ids.value(row));
+        }
+    }
+    result
+}
+
+async fn values(table: &Table) -> Vec<(Option<String>, i32, i32, i32)> {
+    let mut result = Vec::new();
+    for batch in read(table).await {
+        let p = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<StringArray>()
+            .unwrap();
+        let q = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let id = batch
+            .column(2)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let v = batch
+            .column(3)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        for row in 0..batch.num_rows() {
+            result.push((
+                (!p.is_null(row)).then(|| p.value(row).to_string()),
+                q.value(row),
+                id.value(row),
+                v.value(row),
+            ));
+        }
+    }
+    result.sort_by_key(|row| row.2);
+    result
+}
+
+fn matched(ids: Vec<i64>, p: Vec<Option<&str>>, q: Vec<i64>, v: Vec<i32>) -> 
RecordBatch {
+    RecordBatch::try_from_iter([
+        ("_ROW_ID", Arc::new(Int64Array::from(ids)) as ArrayRef),
+        ("p", Arc::new(StringArray::from(p)) as ArrayRef),
+        // The table uses INT: compare after the ordinary row-update coercion.
+        ("q", Arc::new(Int64Array::from(q)) as ArrayRef),
+        ("v", Arc::new(Int32Array::from(v)) as ArrayRef),
+    ])
+    .unwrap()
+}
+
+async fn commit(table: &Table, messages: Vec<CommitMessage>) {
+    table
+        .new_write_builder()
+        .new_commit()
+        .commit(messages)
+        .await
+        .unwrap();
+}
+
+async fn files(table: &Table) -> Vec<String> {
+    let mut files = table
+        .file_io()
+        .list_status_recursive("memory:/")
+        .await
+        .unwrap()
+        .into_iter()
+        .map(|status| status.path)
+        .filter(|path| path.ends_with(".parquet") || path.contains("/index-"))
+        .collect::<Vec<_>>();
+    files.sort();
+    files
+}
+
+#[tokio::test]
+async fn row_id_updates_carry_unchanged_composite_and_null_partitions() {
+    let table = table(&[]).await;
+    seed(&table).await;
+    let ids = row_ids(&table).await;
+    let update = table.new_write_builder().new_update().unwrap();
+    let mut writer = update.new_update_by_row_id().await.unwrap();
+    let input = matched(
+        vec![ids[&3], ids[&0]],
+        vec![None, Some("a")],
+        vec![2, 1],
+        vec![130, 100],
+    );
+    let messages = writer
+        .update_columns(vec![input], vec!["p".into(), "q".into(), "v".into()])
+        .await
+        .unwrap();
+    assert!(messages
+        .iter()
+        .flat_map(|message| &message.new_files)
+        .all(|file| file.write_cols.as_ref().unwrap().contains(&"p".into())));
+    let before = files(&table).await;
+    let error = writer
+        .update_columns(
+            vec![matched(vec![ids[&0]], vec![Some("a")], vec![1], vec![0])],
+            vec!["p".into()],
+        )
+        .await
+        .unwrap_err();
+    assert!(error.to_string().contains("overlap"), "{error}");
+    assert_eq!(files(&table).await, before);
+    commit(&table, messages).await;
+    assert_eq!(
+        values(&table).await,
+        vec![
+            (Some("a".into()), 1, 0, 100),
+            (Some("a".into()), 1, 1, 11),
+            (None, 2, 2, 12),
+            (None, 2, 3, 130)
+        ]
+    );
+
+    // SQL assignment entry points retain their partition-column restriction.
+    assert!(DataEvolutionWriter::new(&table, vec!["p".into()]).is_err());
+    let mut writer = update.new_update_by_row_id().await.unwrap();
+    let messages = writer
+        .update_columns(
+            vec![matched(vec![ids[&0]], vec![Some("a")], vec![1], vec![0])],
+            vec!["p".into(), "q".into()],
+        )
+        .await
+        .unwrap();
+    commit(&table, messages).await;
+    assert_eq!(values(&table).await[0].3, 100);
+}
+
+#[tokio::test]
+async fn partition_changes_fail_before_writes_and_preserve_prior_messages() {
+    let table = table(&[]).await;
+    seed(&table).await;
+    let ids = row_ids(&table).await;
+    let update = table.new_write_builder().new_update().unwrap();
+    let mut writer = update.new_update_by_row_id().await.unwrap();
+    // Unselected input fields do not constrain the update.
+    let saved = writer
+        .update_columns(
+            vec![matched(
+                vec![ids[&0]],
+                vec![Some("ignored")],
+                vec![999],
+                vec![100],
+            )],
+            vec!["v".into()],
+        )
+        .await
+        .unwrap();
+    let before = files(&table).await;
+    for (p, q) in [(Some("changed"), 2), (Some("a"), 2), (None, 3)] {
+        let input = matched(
+            vec![ids[&0], ids[&2]],
+            vec![Some("a"), p],
+            vec![1, q],
+            vec![1, 2],
+        );
+        let error = writer
+            .update_columns(vec![input], vec!["p".into(), "q".into()])
+            .await
+            .unwrap_err();
+        assert!(
+            error.to_string().contains("Cannot change partition column"),
+            "{error}"
+        );
+        assert_eq!(files(&table).await, before);
+        assert_eq!(writer.commit_messages().len(), saved.len());
+    }
+    let error = writer
+        .update_columns(
+            vec![matched(vec![ids[&0]], vec![None], vec![1], vec![0])],
+            vec!["p".into()],
+        )
+        .await
+        .unwrap_err();
+    assert!(error.to_string().contains("Cannot change partition column"));
+    commit(&table, saved).await;
+    assert_eq!(values(&table).await[0].3, 100);
+}
+
+#[tokio::test]
+async fn upsert_matches_within_partition_even_when_keys_omit_partition() {
+    let table = table(&[]).await;
+    write_batch(
+        &table,
+        &batch(
+            vec![Some("a"), Some("b"), None],
+            vec![1, 1, 2],
+            vec![7, 7, 7],
+            vec![10, 20, 30],
+        ),
+    )
+    .await;
+    let update = table.new_write_builder().new_update().unwrap();
+    let input = batch(
+        vec![Some("a"), Some("c"), None, Some("a")],
+        vec![1, 1, 2, 1],
+        vec![7, 7, 7, 7],
+        vec![11, 40, 31, 12],
+    );
+    let messages = update
+        .upsert_by_arrow_with_key(vec![input], vec!["id".into()])
+        .await
+        .unwrap();
+    commit(&table, messages).await;
+    let mut rows = values(&table).await;
+    rows.sort();
+    assert_eq!(
+        rows,
+        vec![
+            (None, 2, 7, 31),
+            (Some("a".into()), 1, 7, 12),
+            (Some("b".into()), 1, 7, 20),
+            (Some("c".into()), 1, 7, 40)
+        ]
+    );
+    // A partition-only key targets all rows in that partition; duplicates use
+    // the last source row, just as Python's per-partition matching does.
+    let input = batch(vec![Some("a")], vec![1], vec![8], vec![80]);
+    commit(
+        &table,
+        update
+            .upsert_by_arrow_with_key(vec![input], vec!["p".into(), 
"q".into()])
+            .await
+            .unwrap(),
+    )
+    .await;
+    assert!(values(&table).await.contains(&(Some("a".into()), 1, 8, 80)));
+}
+
+async fn delete(table: &Table, ids: Vec<i64>) -> Vec<CommitMessage> {
+    let mut writer = table.new_write_builder().new_delete().unwrap();
+    writer.add_row_ids(ids).unwrap();
+    writer.prepare_commit().await.unwrap()
+}
+
+#[tokio::test]
+async fn external_deletion_vectors_repeat_time_travel_and_abort() {
+    for (in_bucket, strategy) in [
+        (true, "round-robin"),
+        (true, "weight-robin"),
+        (true, "entropy-inject"),
+        (true, "specific-fs"),
+        (false, "round-robin"),
+    ] {
+        let table = table(&[
+            (
+                "index-file-in-data-file-dir",
+                if in_bucket { "true" } else { "false" },
+            ),
+            (
+                "data-file.external-paths",
+                "memory:/external-data/a,memory:/external-data/b",
+            ),
+            ("data-file.external-paths.strategy", strategy),
+            ("data-file.external-paths.weights", "1,3"),
+            ("data-file.external-paths.specific-fs", "memory"),
+            ("global-index.external-path", "memory:/external-index"),
+        ])
+        .await;
+        seed(&table).await;
+        let ids = row_ids(&table).await;
+        let messages = delete(&table, vec![ids[&0], ids[&2]]).await;
+        let paths = messages
+            .iter()
+            .flat_map(|message| &message.new_index_files)
+            .map(|file| file.external_path.clone().unwrap())
+            .collect::<Vec<_>>();
+        assert_eq!(paths.len(), 2);
+        for path in &paths {
+            assert!(path.starts_with(if in_bucket {
+                "memory:/external-data/"
+            } else {
+                "memory:/external-index/"
+            }));
+            assert!(table.file_io().exists(path).await.unwrap());
+            assert_eq!(path.contains("/bucket-"), in_bucket);
+        }
+        commit(&table, messages).await;
+        assert_eq!(row_ids(&table).await.len(), 2);
+        // A changed output root must not redirect existing explicit paths.
+        let changed = table.copy_with_options(HashMap::from([
+            ("data-file.external-paths".into(), "memory:/new-data".into()),
+            (
+                "global-index.external-path".into(),
+                "memory:/new-index".into(),
+            ),
+        ]));
+        let messages = delete(&changed, vec![ids[&1]]).await;
+        let staged = messages[0].new_index_files[0]
+            .external_path
+            .clone()
+            .unwrap();
+        assert!(staged.contains("/new-"));
+        changed
+            .new_write_builder()
+            .new_commit()
+            .abort(&messages)
+            .await
+            .unwrap();
+        assert!(!table.file_io().exists(&staged).await.unwrap());
+        for path in &paths {
+            assert!(table.file_io().exists(path).await.unwrap());
+        }
+        commit(&changed, delete(&changed, vec![ids[&1]]).await).await;
+        assert_eq!(
+            row_ids(&changed).await.keys().copied().collect::<Vec<_>>(),
+            vec![3]
+        );
+        let historical = table
+            .copy_with_time_travel(HashMap::from([("scan.snapshot-id".into(), 
"2".into())]))
+            .await
+            .unwrap();
+        assert_eq!(row_ids(&historical).await.len(), 2);
+    }
+}
+
+#[tokio::test]
+async fn 
delete_reads_legacy_index_directory_without_changing_manifest_identity() {
+    let table = table(&[("index-file-in-data-file-dir", "true")]).await;
+    seed(&table).await;
+    let ids = row_ids(&table).await;
+    let messages = delete(&table, vec![ids[&0]]).await;
+    let old = messages[0].new_index_files[0].clone();
+    assert!(old.external_path.is_none());
+    let bucket_path = format!("{}/p=a/q=1/bucket-0/{}", table.location(), 
old.file_name);
+    let legacy_path = format!("{}/index/{}", table.location(), old.file_name);
+    let bytes = table
+        .file_io()
+        .new_input(&bucket_path)
+        .unwrap()
+        .read()
+        .await
+        .unwrap();
+    table
+        .file_io()
+        .new_output(&legacy_path)
+        .unwrap()
+        .write(bytes)
+        .await
+        .unwrap();
+    table.file_io().delete_file(&bucket_path).await.unwrap();
+    commit(&table, messages).await;
+    let messages = delete(&table, vec![ids[&1]]).await;
+    assert_eq!(messages[0].deleted_index_files, vec![old]);
+    commit(&table, messages).await;
+    assert_eq!(row_ids(&table).await.len(), 2);
+    assert!(table.file_io().exists(&legacy_path).await.unwrap());
+    let historical = table
+        .copy_with_time_travel(HashMap::from([("scan.snapshot-id".into(), 
"2".into())]))
+        .await
+        .unwrap();
+    assert_eq!(row_ids(&historical).await.len(), 3);
+}
+
+#[tokio::test]
+async fn missing_explicit_deletion_vector_never_falls_back_to_local_decoys() {
+    let table = table(&[
+        ("index-file-in-data-file-dir", "true"),
+        ("data-file.external-paths", "memory:/external"),
+        ("data-file.external-paths.strategy", "round-robin"),
+    ])
+    .await;
+    seed(&table).await;
+    let ids = row_ids(&table).await;
+    let messages = delete(&table, vec![ids[&0]]).await;
+    let file = messages[0].new_index_files[0].clone();
+    let path = file.external_path.as_ref().unwrap();
+    let bytes = table
+        .file_io()
+        .new_input(path)
+        .unwrap()
+        .read()
+        .await
+        .unwrap();
+    for dir in ["index", "p=a/q=1/bucket-0"] {
+        table
+            .file_io()
+            .new_output(&format!("{}/{dir}/{}", table.location(), 
file.file_name))
+            .unwrap()
+            .write(bytes.clone())
+            .await
+            .unwrap();
+    }
+    commit(&table, messages).await;
+    table.file_io().delete_file(path).await.unwrap();
+    let mut writer = table.new_write_builder().new_delete().unwrap();
+    writer.add_row_ids(vec![ids[&1]]).unwrap();
+    let before = files(&table).await;
+    assert!(writer.prepare_commit().await.is_err());
+    assert_eq!(files(&table).await, before);
+}
+
+#[tokio::test]
+async fn later_bucket_failure_aborts_earlier_external_deletion_vectors() {
+    let table = table(&[
+        ("index-file-in-data-file-dir", "false"),
+        ("global-index.external-path", "memory:/external"),
+    ])
+    .await;
+    seed(&table).await;
+    let ids = row_ids(&table).await;
+    let messages = delete(&table, vec![ids[&0], ids[&2]]).await;
+    // The writer processes bucket keys in order. Removing the last bucket's
+    // old DV forces failure after staging a replacement for the first bucket.
+    let missing = messages.last().unwrap().new_index_files[0]
+        .external_path
+        .clone()
+        .unwrap();
+    commit(&table, messages).await;
+    table.file_io().delete_file(&missing).await.unwrap();
+    let before = files(&table).await;
+    let mut writer = table.new_write_builder().new_delete().unwrap();
+    writer.add_row_ids(vec![ids[&1], ids[&3]]).unwrap();
+    assert!(writer.prepare_commit().await.is_err());
+    assert_eq!(files(&table).await, before);
+}

Reply via email to