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);
+}