From 71bbdb4ede73737fc6d231617e09a3f5a9cfecf3 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Sun, 27 Sep 2026 23:46:11 +0800 Subject: [PATCH 1/2] Support configured data directories throughout table IO --- .../src/system_tables/physical_files_size.rs | 12 +- .../system_tables/referenced_files_size.rs | 7 +- crates/paimon/src/catalog/filesystem.rs | 60 ++- crates/paimon/src/spec/core_options.rs | 17 + crates/paimon/src/spec/mod.rs | 3 +- crates/paimon/src/spec/partition_utils.rs | 227 +++++++++ crates/paimon/src/spec/schema.rs | 25 +- .../src/table/bucket_assigner_dynamic.rs | 15 +- .../paimon/src/table/data_evolution_writer.rs | 17 +- crates/paimon/src/table/data_file_writer.rs | 6 +- .../src/table/dedicated_format_file_writer.rs | 2 +- crates/paimon/src/table/external_path.rs | 43 +- .../table/global_index_build_common/vector.rs | 4 +- crates/paimon/src/table/index_file_path.rs | 2 +- crates/paimon/src/table/kv_file_writer.rs | 10 +- .../paimon/src/table/managed_blob_writer.rs | 5 +- crates/paimon/src/table/mod.rs | 8 + crates/paimon/src/table/read_builder.rs | 1 + crates/paimon/src/table/referenced_files.rs | 170 +++++-- .../planning.rs | 4 +- crates/paimon/src/table/table_commit.rs | 8 +- crates/paimon/src/table/table_scan.rs | 20 +- crates/paimon/src/table/table_write.rs | 34 +- .../paimon/tests/data_file_directory_test.rs | 470 ++++++++++++++++++ 24 files changed, 1044 insertions(+), 126 deletions(-) create mode 100644 crates/paimon/tests/data_file_directory_test.rs diff --git a/crates/integrations/datafusion/src/system_tables/physical_files_size.rs b/crates/integrations/datafusion/src/system_tables/physical_files_size.rs index 01ef43c1b..0c46ab4d5 100644 --- a/crates/integrations/datafusion/src/system_tables/physical_files_size.rs +++ b/crates/integrations/datafusion/src/system_tables/physical_files_size.rs @@ -28,7 +28,9 @@ use datafusion::datasource::{TableProvider, TableType}; use datafusion::error::Result as DFResult; use datafusion::logical_expr::Expr; use datafusion::physical_plan::ExecutionPlan; -use paimon::table::referenced_files::{collect_physical_files_summary, PhysicalFilesSummary}; +use paimon::table::referenced_files::{ + collect_physical_files_summary_with_options, PhysicalFilesSummary, +}; use paimon::table::Table; use crate::error::to_datafusion_error; @@ -78,7 +80,13 @@ impl TableProvider for PhysicalFilesSizeTable { let table = self.table.clone(); let summary = crate::runtime::await_with_runtime(async move { let partition_depth = table.schema().partition_keys().len(); - collect_physical_files_summary(table.file_io(), table.location(), partition_depth).await + collect_physical_files_summary_with_options( + table.file_io(), + table.location(), + partition_depth, + &table.schema().core_options(), + ) + .await }) .await .map_err(to_datafusion_error)?; diff --git a/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs b/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs index 568663ca3..db09d0e72 100644 --- a/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs +++ b/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs @@ -28,7 +28,9 @@ use datafusion::datasource::{TableProvider, TableType}; use datafusion::error::Result as DFResult; use datafusion::logical_expr::Expr; use datafusion::physical_plan::ExecutionPlan; -use paimon::table::referenced_files::{collect_referenced_files_summary, ReferencedFilesSummary}; +use paimon::table::referenced_files::{ + collect_referenced_files_summary_with_options, ReferencedFilesSummary, +}; use paimon::table::Table; use crate::error::to_datafusion_error; @@ -81,11 +83,12 @@ impl TableProvider for ReferencedFilesSizeTable { let schema = table.schema(); let partition_keys = schema.partition_keys(); let partition_fields = schema.partition_fields(); - collect_referenced_files_summary( + collect_referenced_files_summary_with_options( table.file_io(), table.location(), partition_keys, &partition_fields, + &schema.core_options(), ) .await }) diff --git a/crates/paimon/src/catalog/filesystem.rs b/crates/paimon/src/catalog/filesystem.rs index d660697f7..23195c781 100644 --- a/crates/paimon/src/catalog/filesystem.rs +++ b/crates/paimon/src/catalog/filesystem.rs @@ -704,8 +704,8 @@ fn local_datetime_to_millis( /// snapshot is not required: a format table holds data without ever writing /// one, so a snapshot check would not catch it. Java special-cases `type` the /// same way (`SchemaManager.generateTableSchema`). -/// * `index-file-in-data-file-dir` picks the directory every bucket-local index -/// file is written to and read from, so flipping it hides every index file the +/// * `index-file-in-data-file-dir` and `data-file.path-directory` pick the directories +/// files are written to and read from, so flipping them hides files the /// table already has. A change is rejected once the table exists, and a /// `SetOption` repeating the stored value or a `RemoveOption` for an option that /// is not set is let through, since neither moves anything. @@ -745,25 +745,27 @@ fn reject_immutable_option_changes( }); } crate::spec::SchemaChange::SetOption { key, value } - if key == INDEX_FILE_IN_DATA_FILE_DIR_OPTION + if (key == INDEX_FILE_IN_DATA_FILE_DIR_OPTION + || key == "data-file.path-directory") && current_options.get(key.as_str()) != Some(value) => { return Err(Error::Unsupported { message: format!( - "changing '{INDEX_FILE_IN_DATA_FILE_DIR_OPTION}' after the table exists \ - is not supported: it selects the directory index files are written to, \ + "changing '{key}' after the table exists \ + is not supported: it selects the directory files are written to, \ so the files already written would no longer be found" ), }); } crate::spec::SchemaChange::RemoveOption { key } - if key == INDEX_FILE_IN_DATA_FILE_DIR_OPTION + if (key == INDEX_FILE_IN_DATA_FILE_DIR_OPTION + || key == "data-file.path-directory") && current_options.contains_key(key.as_str()) => { return Err(Error::Unsupported { message: format!( - "removing '{INDEX_FILE_IN_DATA_FILE_DIR_OPTION}' is not supported: \ - it selects the directory index files are written to, so the files \ + "removing '{key}' is not supported: \ + it selects the directory files are written to, so the files \ already written would no longer be found" ), }); @@ -1277,6 +1279,48 @@ mod tests { assert!(matches!(err, Error::Unsupported { .. }), "{err:?}"); } + #[tokio::test] + async fn test_alter_table_cannot_change_data_directory() { + use crate::spec::SchemaChange; + for stored in [None, Some("data/nested")] { + let (_temp_dir, catalog) = create_test_catalog(); + let options = stored + .map(|value| { + HashMap::from([("data-file.path-directory".to_string(), value.to_string())]) + }) + .unwrap_or_default(); + let identifier = create_table_for_alter(&catalog, options).await; + give_the_table_a_snapshot(&catalog, &identifier).await; + let key = "data-file.path-directory".to_string(); + let original = catalog.get_table(&identifier).await.unwrap(); + let mut changes = vec![SchemaChange::set_option(key.clone(), "other".to_string())]; + if stored.is_some() { + changes.push(SchemaChange::remove_option(key.clone())); + } + for change in changes { + let error = catalog + .alter_table(&identifier, vec![change], false) + .await + .unwrap_err(); + assert!(matches!(error, Error::Unsupported { .. }), "{error:?}"); + } + let no_op = match stored { + Some(value) => SchemaChange::set_option(key, value.to_string()), + None => SchemaChange::remove_option(key), + }; + catalog + .alter_table(&identifier, vec![no_op], false) + .await + .unwrap(); + let reloaded = catalog.get_table(&identifier).await.unwrap(); + assert_eq!( + reloaded.schema().core_options().data_file_path_directory(), + stored + ); + assert_eq!(reloaded.schema().options(), original.schema().options()); + } + } + #[tokio::test] async fn test_alter_table_index_layout_no_ops_are_allowed() { use crate::spec::SchemaChange; diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 984e13806..d474479b1 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -1404,6 +1404,23 @@ impl<'a> CoreOptions<'a> { .unwrap_or(1) } + pub(crate) fn validate_data_file_path_directory(&self) -> crate::Result<()> { + if self.data_file_path_directory() == Some("") { + return Err(crate::Error::ConfigInvalid { + message: "data-file.path-directory must not be empty".into(), + }); + } + Ok(()) + } + + /// Directory containing bucket data, resolved against the table location. + /// Supports relative paths, absolute paths, and URIs, like Java `Path`. + pub fn data_file_path_directory(&self) -> Option<&str> { + self.options + .get("data-file.path-directory") + .map(String::as_str) + } + /// File name prefix for data files. Default is `"data-"`. pub fn data_file_prefix(&self) -> &str { self.options diff --git a/crates/paimon/src/spec/mod.rs b/crates/paimon/src/spec/mod.rs index f1da81a42..5fe8faa11 100644 --- a/crates/paimon/src/spec/mod.rs +++ b/crates/paimon/src/spec/mod.rs @@ -104,7 +104,8 @@ mod partition; pub use partition::Partition; mod partition_utils; pub(crate) use partition_utils::{ - bucket_path, bucket_path_under, escape_path_name, unescape_path_name, PartitionComputer, + bucket_path, bucket_path_under, data_file_path, escape_path_name, relative_bucket_path, + unescape_path_name, PartitionComputer, }; mod predicate; pub(crate) use predicate::datum_cmp; diff --git a/crates/paimon/src/spec/partition_utils.rs b/crates/paimon/src/spec/partition_utils.rs index 99391ef1c..96c0d78ac 100644 --- a/crates/paimon/src/spec/partition_utils.rs +++ b/crates/paimon/src/spec/partition_utils.rs @@ -31,6 +31,134 @@ use chrono::{Local, NaiveDate, NaiveDateTime, TimeZone, Timelike}; const MILLIS_PER_DAY: i64 = 86_400_000; +/// Java `FileStorePathFactory.dataFilePath`: resolve the configured directory +/// against the table root. This applies only to data, never table metadata. +/// Like Java Path, preserve literal percent signs, spaces, `?` and `#` in names. +pub(crate) fn data_file_path(table_path: &str, directory: Option<&str>) -> String { + let Some(directory) = directory.filter(|directory| !directory.is_empty()) else { + return if table_path == "/" { + "/" + } else { + table_path.trim_end_matches('/') + } + .to_string(); + }; + resolve_data_path(table_path, directory, cfg!(windows)) +} + +fn resolve_data_path(parent: &str, child: &str, windows: bool) -> String { + fn has_drive(path: &str) -> bool { + let path = path.strip_prefix('/').unwrap_or(path).as_bytes(); + path.len() >= 2 && path[0].is_ascii_alphabetic() && path[1] == b':' + } + fn normalize(path: &str, windows: bool) -> String { + let mut components = Vec::new(); + for component in path.split('/') { + match component { + "" | "." => {} + ".." if components.last().is_some_and(|last| *last != "..") => { + components.pop(); + } + component => components.push(component), + } + } + let leading = if path.starts_with('/') { "/" } else { "" }; + let normalized = format!("{leading}{}", components.join("/")); + if windows && has_drive(&normalized) && normalized.len() == 3 && path != normalized { + format!("{normalized}/") + } else { + normalized + } + } + // Keep scheme and authority separate: a port or user-info can contain ':'. + fn parts(input: &str, windows: bool) -> (String, String, String) { + let input = if windows && has_drive(input) && !input.starts_with('/') { + format!("/{input}") + } else { + input.to_string() + }; + let scheme_end = input + .find(':') + .filter(|&i| !input[..i].contains('/')) + .map_or(0, |i| i + 1); + let scheme = &input[..scheme_end]; + let (authority, path) = + if input[scheme_end..].starts_with("//") && input.len() - scheme_end > 2 { + let start = scheme_end + 2; + let end = input[start..].find('/').map_or(input.len(), |i| start + i); + let authority = if end == start { + "" + } else { + &input[scheme_end..end] + }; + (authority, &input[end..]) + } else { + ("", &input[scheme_end..]) + }; + // Java collapses forward slashes before converting Windows separators. + let mut previous_slash = false; + let path: String = path + .chars() + .filter(|&c| { + let duplicate = c == '/' && previous_slash; + previous_slash = c == '/'; + !duplicate + }) + .collect(); + let path = if windows && (has_drive(&path) || scheme.is_empty() || scheme == "file:") { + path.replace('\\', "/") + } else { + path.to_string() + }; + // Java URI recognizes a UNC authority after Windows separator conversion. + let (authority, path) = if authority.is_empty() && path.starts_with("//") && path.len() > 2 + { + let end = path[2..].find('/').map_or(path.len(), |i| i + 2); + (&path[..end], &path[end..]) + } else { + (authority, path.as_str()) + }; + ( + scheme.to_string(), + authority.to_string(), + normalize(path, windows), + ) + } + let (parent_scheme, parent_authority, parent_path) = parts(parent, windows); + let (child_scheme, child_authority, child_path) = parts(child, windows); + let (scheme, authority, path) = if !child_scheme.is_empty() { + (child_scheme, child_authority, child_path) + } else if !child_authority.is_empty() { + (parent_scheme, child_authority, child_path) + } else if child_path.starts_with('/') { + (parent_scheme, parent_authority, child_path) + } else { + let path = + if parent_path.is_empty() && parent_scheme.is_empty() && parent_authority.is_empty() { + child_path + } else { + format!("{}/{child_path}", parent_path.trim_end_matches('/')) + }; + (parent_scheme, parent_authority, path) + }; + let path = normalize(&path, windows); + let path = if windows && scheme.is_empty() && authority.is_empty() && has_drive(&path) { + path.strip_prefix('/').unwrap_or(&path).to_string() + } else if scheme.is_empty() + && authority.is_empty() + && !path.starts_with('/') + && path + .split('/') + .next() + .is_some_and(|part| part.contains(':')) + { + format!("./{path}") + } else { + path + }; + format!("{scheme}{authority}{path}") +} + /// The directory holding one bucket's data files, `/[/]bucket-N`. /// /// Mirrors Java `FileStorePathFactory.bucketPath`. Every consumer that needs a @@ -53,6 +181,20 @@ pub(crate) fn bucket_path( Ok(bucket_path_under(table_path, &partition_path, bucket)) } +/// Java `relativeBucketPath`, which may itself be absolute when the configured +/// data directory is absolute. External path providers must resolve, not append it. +pub(crate) fn relative_bucket_path( + partition_path: &str, + bucket: i32, + directory: Option<&str>, +) -> String { + let relative = format!("{partition_path}{}", crate::spec::bucket_dir_name(bucket)); + match directory.filter(|directory| !directory.is_empty()) { + Some(directory) => data_file_path(directory, Some(&relative)), + None => relative, + } +} + /// [`bucket_path`] for a partition directory that is already computed. /// /// `partition_path` is empty for an unpartitioned table and otherwise ends with `/`, @@ -652,6 +794,91 @@ fn needs_escaping(c: char) -> bool { #[cfg(test)] mod tests { + #[test] + fn data_directories_match_java_path_resolution() { + // Expected strings produced by the Java org.apache.paimon.fs.Path. + for (parent, child, expected) in [ + ("/warehouse/t", "data/nested", "/warehouse/t/data/nested"), + ( + "/warehouse/t/", + "data//discard/../nested/.", + "/warehouse/t/data/nested", + ), + ( + "s3://bucket/warehouse/t", + "/shared/data", + "s3://bucket/shared/data", + ), + ("s3://bucket/t", "s3://other/data", "s3://other/data"), + ( + "file:/warehouse/t", + "file:///other/data", + "file:/other/data", + ), + ("memory:/t", "../other", "memory:/other"), + ("warehouse/t", "../data", "warehouse/data"), + ("warehouse/t", "data", "warehouse/t/data"), + ("s3://bucket/t", "//other/data", "s3://other/data"), + ( + "s3://bucket/t", + "data 100%?#/你好", + "s3://bucket/t/data 100%?#/你好", + ), + ("file:/t", "./data", "file:/t/data"), + ("/", "data", "/data"), + ("/t", ".", "/t"), + ("/t", "..", "/"), + ("/t", "../../data", "/../data"), + ("s3://bucket/", "data", "s3://bucket/data"), + ("warehouse", "a/../b:c", "warehouse/b:c"), + ("memory:/t", "a/../../b", "memory:/b"), + ("/t", "///data", "/data"), + ("/t", "////data", "/data"), + ("file:/t", "file:////data", "file:/data"), + ("/t", "//", "/"), + ("file:/t", "//", "file:/"), + ("s3://bucket/t", "//other:9000/data", "s3://other:9000/data"), + ("//host:8020/table", "//other/data", "//other/data"), + ("relative", "../b:c", "./b:c"), + ("file:/t", "file:/data", "file:/data"), + ] { + assert_eq!( + data_file_path(parent, Some(child)), + expected, + "{parent} + {child}" + ); + } + } + + #[test] + fn windows_data_directories_match_java() { + for (parent, child, expected) in [ + (r"C:\warehouse\table", "../data", "C:/warehouse/data"), + (r"C:\warehouse\table", r"..\data", "C:/warehouse/data"), + ("C:/warehouse/table", "/data", "/data"), + ("C:/warehouse/table", "../..", "C:/"), + ("file:/C:/warehouse/table", "../..", "file:/C:/"), + (r"\\server\share\table", "data", "//server/share/table/data"), + ("s3://bucket/table", r"\\other\share", "s3://other/share"), + ( + "file:/C:/warehouse/table", + r"..\data", + "file:/C:/warehouse/data", + ), + ( + "s3://bucket/t", + r"data\literal", + "s3://bucket/t/data/literal", + ), + ] { + assert_eq!( + resolve_data_path(parent, child, true), + expected, + "{parent} + {child}" + ); + } + } + use super::*; use crate::spec::types::*; use crate::spec::DataField; diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs index c7d157365..7386a625b 100644 --- a/crates/paimon/src/spec/schema.rs +++ b/crates/paimon/src/spec/schema.rs @@ -149,10 +149,10 @@ impl TableSchema { /// override, and the declared `type` can't be changed by one: an override /// could re-route foreign data through the Paimon reader. /// - /// `index-file-in-data-file-dir` can't be changed by one either. It selects - /// the directory every bucket-local index file is written to and read from, - /// while an index manifest records only the file name, so a copy carrying an - /// overridden value would write hash and deletion-vector index files where a + /// `index-file-in-data-file-dir` and `data-file.path-directory` also retain + /// their stored values. They select where data and bucket-local indexes are + /// written and read, while manifests record file names, so a copy carrying an + /// overridden value would write files where a /// normally loaded table cannot find them. Java rejects such an override /// outright in `AbstractFileStoreTable.checkImmutability`; this copy cannot /// fail, so the stored value wins instead, and an absent one leaves the @@ -165,13 +165,15 @@ impl TableSchema { Some(declared) => extra.insert(TABLE_TYPE_OPTION.to_string(), declared.clone()), None => extra.remove(TABLE_TYPE_OPTION), }; - match self.options.get(INDEX_FILE_IN_DATA_FILE_DIR_OPTION) { - Some(stored) => extra.insert( - INDEX_FILE_IN_DATA_FILE_DIR_OPTION.to_string(), - stored.clone(), - ), - None => extra.remove(INDEX_FILE_IN_DATA_FILE_DIR_OPTION), - }; + for key in [ + INDEX_FILE_IN_DATA_FILE_DIR_OPTION, + "data-file.path-directory", + ] { + match self.options.get(key) { + Some(stored) => extra.insert(key.to_string(), stored.clone()), + None => extra.remove(key), + }; + } let mut new_schema = self.clone(); new_schema.options.extend(extra); new_schema @@ -1213,6 +1215,7 @@ impl Schema { Self::validate_bucket_keys(options, fields, partition_keys, primary_keys)?; Self::validate_sequence_field(options, fields, partition_keys, primary_keys)?; Self::validate_read_batch_size(options)?; + CoreOptions::new(options).validate_data_file_path_directory()?; let core_options = CoreOptions::new(options); if !core_options.is_format_table() { core_options.target_file_row_num()?; diff --git a/crates/paimon/src/table/bucket_assigner_dynamic.rs b/crates/paimon/src/table/bucket_assigner_dynamic.rs index 57fbc48e4..2919d5992 100644 --- a/crates/paimon/src/table/bucket_assigner_dynamic.rs +++ b/crates/paimon/src/table/bucket_assigner_dynamic.rs @@ -167,6 +167,7 @@ struct HashIndexLayout<'a> { table_path: &'a str, /// Partition directory, already terminated by `/`, or empty when unpartitioned. partition_path: &'a str, + data_file_path_directory: Option<&'a str>, index_file_in_data_file_dir: bool, } @@ -182,7 +183,8 @@ impl HashIndexLayout<'_> { } fn bucket_path(&self, bucket: i32) -> String { - bucket_path_under(self.table_path, self.partition_path, bucket) + let data_root = crate::spec::data_file_path(self.table_path, self.data_file_path_directory); + bucket_path_under(&data_root, self.partition_path, bucket) } /// The directory a new hash index file for `bucket` is written into. @@ -373,6 +375,7 @@ pub(crate) struct DynamicBucketAssigner { target_bucket_row_number: i64, file_io: FileIO, table_location: String, + data_file_path_directory: Option, /// Cached index manifest entries from the latest snapshot (loaded once). cached_index_entries: Option>, /// Overwrite mode: skip loading existing index entries. @@ -406,6 +409,7 @@ impl DynamicBucketAssigner { target_bucket_row_number, file_io, table_location, + data_file_path_directory: None, cached_index_entries: None, is_overwrite, partition_computer, @@ -413,6 +417,11 @@ impl DynamicBucketAssigner { } } + pub(super) fn with_data_file_path_directory(mut self, directory: Option<&str>) -> Self { + self.data_file_path_directory = directory.map(str::to_string); + self + } + pub fn set_overwrite(&mut self, is_overwrite: bool) { self.is_overwrite = is_overwrite; } @@ -466,6 +475,7 @@ impl DynamicBucketAssigner { if !partition_entries.is_empty() { let partition_path = self.partition_path(partition_bytes)?; let layout = HashIndexLayout { + data_file_path_directory: self.data_file_path_directory.as_deref(), table_path: self.table_location.trim_end_matches('/'), partition_path: &partition_path, index_file_in_data_file_dir: self.index_file_in_data_file_dir, @@ -542,6 +552,7 @@ impl BucketAssigner for DynamicBucketAssigner { } for (partition_bytes, partition_path) in partition_keys.into_iter().zip(partition_paths) { let layout = HashIndexLayout { + data_file_path_directory: self.data_file_path_directory.as_deref(), table_path: &table_path, partition_path: &partition_path, index_file_in_data_file_dir, @@ -649,6 +660,7 @@ mod tests { let table_path = format!("file://{}", tmp.path().display()); let file_io = FileIO::from_url(&table_path).unwrap().build().unwrap(); let layout = super::HashIndexLayout { + data_file_path_directory: None, table_path: &table_path, partition_path: "pt=1/", index_file_in_data_file_dir, @@ -697,6 +709,7 @@ mod tests { for index_file_in_data_file_dir in [false, true] { let layout = super::HashIndexLayout { + data_file_path_directory: None, table_path: &table_path, partition_path: "pt=1/", index_file_in_data_file_dir, diff --git a/crates/paimon/src/table/data_evolution_writer.rs b/crates/paimon/src/table/data_evolution_writer.rs index f070fe4e0..85e21762c 100644 --- a/crates/paimon/src/table/data_evolution_writer.rs +++ b/crates/paimon/src/table/data_evolution_writer.rs @@ -457,6 +457,7 @@ impl DataEvolutionWriter { /// write that follows it cannot disagree about where the file goes. struct DeletionVectorLayout { table_path: String, + relative_bucket_path: String, bucket_path: String, index_file_in_data_file_dir: bool, } @@ -722,10 +723,19 @@ impl DataEvolutionDeleteWriter { } else { EMPTY_BINARY_ROW }; + let partition_path = match &computer { + Some(computer) => computer.generate_partition_path(&partition_row)?, + None => String::new(), + }; Ok(DeletionVectorLayout { + relative_bucket_path: crate::spec::relative_bucket_path( + &partition_path, + bucket, + core_options.data_file_path_directory(), + ), table_path: self.table.location().trim_end_matches('/').to_string(), bucket_path: bucket_path( - self.table.location(), + &self.table.data_file_location(), computer.as_ref(), &partition_row, bucket, @@ -806,10 +816,7 @@ impl DataEvolutionDeleteWriter { 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"), + &layout.relative_bucket_path, &file_name, )?; let path = layout diff --git a/crates/paimon/src/table/data_file_writer.rs b/crates/paimon/src/table/data_file_writer.rs index 97037d6d2..d2c013812 100644 --- a/crates/paimon/src/table/data_file_writer.rs +++ b/crates/paimon/src/table/data_file_writer.rs @@ -378,7 +378,11 @@ impl DataFileWriter { } fn bucket_dir(&self) -> String { - bucket_path_under(&self.table_location, &self.partition_path, self.bucket) + let data_root = crate::spec::data_file_path( + &self.table_location, + CoreOptions::new(&self.format_options).data_file_path_directory(), + ); + bucket_path_under(&data_root, &self.partition_path, self.bucket) } #[allow(clippy::too_many_arguments)] diff --git a/crates/paimon/src/table/dedicated_format_file_writer.rs b/crates/paimon/src/table/dedicated_format_file_writer.rs index 4e6102d2e..fb833dfc7 100644 --- a/crates/paimon/src/table/dedicated_format_file_writer.rs +++ b/crates/paimon/src/table/dedicated_format_file_writer.rs @@ -117,7 +117,7 @@ impl AppendDedicatedFormatFileWriter { write_buffer_size, "blob".to_string(), vec![field.clone()], - HashMap::new(), + format_options.clone(), Some(0), None, Some(vec![field.name().to_string()]), diff --git a/crates/paimon/src/table/external_path.rs b/crates/paimon/src/table/external_path.rs index 2b694d15c..9d7b730f4 100644 --- a/crates/paimon/src/table/external_path.rs +++ b/crates/paimon/src/table/external_path.rs @@ -134,7 +134,7 @@ impl ExternalPathProvider { }; Ok(Some(Self { paths: roots, - bucket: bucket.trim_matches('/').into(), + bucket: bucket.to_string(), position, entropy, cumulative_weights, @@ -150,11 +150,10 @@ impl ExternalPathProvider { 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); - } + let mut path = crate::spec::data_file_path( + &self.paths[index], + (!self.bucket.is_empty()).then_some(self.bucket.as_str()), + ); if self.entropy { let hash = crate::spec::murmur_hash::hash_bytes_guava(file_name.as_bytes()) as u32 & 0xfffff; @@ -195,8 +194,8 @@ mod tests { .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!(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() @@ -210,6 +209,26 @@ mod tests { ); } + #[test] + fn configured_directories_resolve_against_external_roots() { + for strategy in ["round-robin", "weight-robin", "entropy-inject"] { + for (bucket, expected) in [ + ("/shared/bucket-0", "file:/shared/bucket-0/"), + ( + "s3://other:9000/shared/bucket-0", + "s3://other:9000/shared/bucket-0/", + ), + ] { + let mut provider = ExternalPathProvider::new(&options(strategy), bucket) + .unwrap() + .unwrap(); + let path = provider.next_path("index-1"); + assert!(path.starts_with(expected), "{strategy}: {path}"); + assert!(path.ends_with("/index-1")); + } + } + } + #[test] fn specific_fs_and_invalid_options() { let mut options = options("specific-fs"); @@ -224,7 +243,7 @@ mod tests { .unwrap() .unwrap() .next_path("index-1"), - "file:///a/bucket-0/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()); @@ -253,7 +272,7 @@ mod tests { .unwrap(); assert_eq!(provider.cumulative_weights, vec![1, 3]); for _ in 0..10 { - assert!(["file:///a/bucket-0/index", "file:///b/bucket-0/index"] + assert!(["file:/a/bucket-0/index", "file:/b/bucket-0/index"] .contains(&provider.next_path("index").as_str())); } } @@ -284,11 +303,11 @@ mod tests { // 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" + "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" + "file:/a/p=x/bucket-0/1011/1111/1010/01000111/hello" ); } } diff --git a/crates/paimon/src/table/global_index_build_common/vector.rs b/crates/paimon/src/table/global_index_build_common/vector.rs index 72a37d68a..a226597df 100644 --- a/crates/paimon/src/table/global_index_build_common/vector.rs +++ b/crates/paimon/src/table/global_index_build_common/vector.rs @@ -240,7 +240,9 @@ fn bucket_path( partition: &BinaryRow, bucket: i32, ) -> Result { - let base = table_location.trim_end_matches('/'); + let data_root = + crate::spec::data_file_path(table_location, core_options.data_file_path_directory()); + let base = data_root.trim_end_matches('/'); if partition_keys.is_empty() { return Ok(format!("{base}/{}", bucket_dir_name(bucket))); } diff --git a/crates/paimon/src/table/index_file_path.rs b/crates/paimon/src/table/index_file_path.rs index a4a3cbd4e..8523e03e2 100644 --- a/crates/paimon/src/table/index_file_path.rs +++ b/crates/paimon/src/table/index_file_path.rs @@ -150,7 +150,7 @@ pub(crate) async fn resolve_legacy_deletion_vector_entries( } let partition = BinaryRow::from_serialized_bytes(&entry.partition)?; let bucket_path = bucket_path( - table_path, + &table.data_file_location(), partition_computer.as_ref(), &partition, entry.bucket, diff --git a/crates/paimon/src/table/kv_file_writer.rs b/crates/paimon/src/table/kv_file_writer.rs index 000a60572..ebef74ab0 100644 --- a/crates/paimon/src/table/kv_file_writer.rs +++ b/crates/paimon/src/table/kv_file_writer.rs @@ -491,7 +491,10 @@ impl KeyValueFileWriter { write.file_format, ); let bucket_dir = bucket_path_under( - &self.config.table_location, + &crate::spec::data_file_path( + &self.config.table_location, + CoreOptions::new(&self.config.table_options).data_file_path_directory(), + ), &self.config.partition_path, self.config.bucket, ); @@ -1173,7 +1176,10 @@ impl KeyValueFileWriter { let _ = reservation.try_resize(0); } let bucket_path = bucket_path_under( - &self.config.table_location, + &crate::spec::data_file_path( + &self.config.table_location, + CoreOptions::new(&self.config.table_options).data_file_path_directory(), + ), &self.config.partition_path, self.config.bucket, ); diff --git a/crates/paimon/src/table/managed_blob_writer.rs b/crates/paimon/src/table/managed_blob_writer.rs index 886a58af9..d53be2322 100644 --- a/crates/paimon/src/table/managed_blob_writer.rs +++ b/crates/paimon/src/table/managed_blob_writer.rs @@ -102,7 +102,10 @@ impl ManagedBlobWriteState { let options = CoreOptions::new(&config.table_options); let writer = ManagedBlobWriter::new( file_io.clone(), - &config.table_location, + &crate::spec::data_file_path( + &config.table_location, + options.data_file_path_directory(), + ), &config.partition_path, config.bucket, &config.data_file_prefix, diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs index c195deda5..5cdf174af 100644 --- a/crates/paimon/src/table/mod.rs +++ b/crates/paimon/src/table/mod.rs @@ -318,6 +318,14 @@ impl Table { &self.location } + /// Root for bucket data. Metadata continues to use `location()`. + pub(crate) fn data_file_location(&self) -> String { + crate::spec::data_file_path( + self.location(), + self.schema().core_options().data_file_path_directory(), + ) + } + /// Get the table's schema. pub fn schema(&self) -> &TableSchema { &self.schema diff --git a/crates/paimon/src/table/read_builder.rs b/crates/paimon/src/table/read_builder.rs index 751833658..d3a308b4e 100644 --- a/crates/paimon/src/table/read_builder.rs +++ b/crates/paimon/src/table/read_builder.rs @@ -553,6 +553,7 @@ impl<'a> PaimonReadBuilder<'a> { // `to_arrow` (e.g. an empty-splits fast path) can't bypass the guard. let core_options = self.table.schema.core_options(); core_options.ensure_read_authorized()?; + core_options.validate_data_file_path_directory()?; let read_type = match self.resolve_read_type()? { None => self.table.schema.fields().to_vec(), Some(fields) => fields, diff --git a/crates/paimon/src/table/referenced_files.rs b/crates/paimon/src/table/referenced_files.rs index 812cf566a..9b5470c40 100644 --- a/crates/paimon/src/table/referenced_files.rs +++ b/crates/paimon/src/table/referenced_files.rs @@ -29,8 +29,8 @@ use std::sync::{ use crate::io::FileIO; use crate::spec::{ - bucket_dir_name, BinaryRow, DataField, DataFileMeta, IndexManifest, Manifest, ManifestEntry, - ManifestFileMeta, PartitionComputer, + bucket_dir_name, BinaryRow, CoreOptions, DataField, DataFileMeta, IndexManifest, Manifest, + ManifestEntry, ManifestFileMeta, PartitionComputer, }; use crate::table::{BranchManager, SnapshotManager, TagManager}; use futures::future::try_join_all; @@ -187,10 +187,40 @@ pub async fn collect_referenced_files_summary( table_location: &str, partition_keys: &[String], schema_fields: &[DataField], +) -> crate::Result> { + collect_referenced_files_summary_with_options( + file_io, + table_location, + partition_keys, + schema_fields, + &CoreOptions::new(&HashMap::new()), + ) + .await +} + +/// Collect referenced sizes using the table's data-directory and partition options. +pub async fn collect_referenced_files_summary_with_options( + file_io: &FileIO, + table_location: &str, + partition_keys: &[String], + schema_fields: &[DataField], + options: &CoreOptions<'_>, ) -> crate::Result> { let manifest_cache: ManifestCache = Mutex::new(HashMap::new()); let manifest_cache_ref = &manifest_cache; - let extra_resolver = ExtraFileResolver::new(table_location, partition_keys, schema_fields); + let mut extra_resolver = ExtraFileResolver::new( + &crate::spec::data_file_path(table_location, options.data_file_path_directory()), + partition_keys, + schema_fields, + ); + if !partition_keys.is_empty() { + extra_resolver.partition_computer = Some(PartitionComputer::new( + partition_keys, + schema_fields, + options.partition_default_name(), + options.legacy_partition_name(), + )?); + } let extra_resolver_ref = &extra_resolver; let sm = SnapshotManager::new(file_io.clone(), table_location.to_string()); @@ -728,62 +758,31 @@ fn is_index_file_in_bucket(segments: &[&str], partition_depth: usize) -> bool { && is_bucket_index_file_name(segments[partition_depth + 1]) } -fn is_data_file_in_data_dir( - relative_path: &str, - data_dir_relative_path: &str, - partition_depth: usize, -) -> bool { - let data_dir_relative_path = data_dir_relative_path.trim_matches('/'); - let data_relative_path = if data_dir_relative_path.is_empty() { - relative_path - } else { - let Some(rest) = relative_path - .strip_prefix(data_dir_relative_path) - .and_then(|rest| rest.strip_prefix('/')) - else { - return false; - }; - rest - }; - let segments = data_relative_path.split('/').collect::>(); - is_data_file_in_bucket(&segments, partition_depth) -} - fn classify_physical_path( table_location: &str, path: &str, partition_depth: usize, data_file_path_directory: Option<&str>, ) -> PhysicalFileKind { - let Some(relative_path) = table_relative_path(table_location, path) else { + let data_root = crate::spec::data_file_path(table_location, data_file_path_directory); + if let Some(relative) = table_relative_path(&data_root, path) { + let segments: Vec<_> = relative.split('/').collect(); + if is_index_file_in_bucket(&segments, partition_depth) { + return PhysicalFileKind::Index; + } + if is_data_file_in_bucket(&segments, partition_depth) { + return PhysicalFileKind::Data; + } + } + let Some(relative) = table_relative_path(table_location, path) else { return PhysicalFileKind::Other; }; - let relative_path = relative_path.trim_matches('/'); - if relative_path.is_empty() { - return PhysicalFileKind::Other; - } - - let segments = relative_path.split('/').collect::>(); - + let segments: Vec<_> = relative.split('/').collect(); match segments.as_slice() { ["manifest", name] if is_manifest_file_name(name) => PhysicalFileKind::Manifest, ["statistics", _] => PhysicalFileKind::Statistics, ["index", _] => PhysicalFileKind::Index, - _ if is_index_file_in_bucket(&segments, partition_depth) => PhysicalFileKind::Index, - _ => { - if let Some(data_dir) = data_file_path_directory { - let data_dir = table_relative_path(table_location, data_dir).unwrap_or(data_dir); - if is_data_file_in_data_dir(relative_path, data_dir, partition_depth) { - PhysicalFileKind::Data - } else { - PhysicalFileKind::Other - } - } else if is_data_file_in_bucket(&segments, partition_depth) { - PhysicalFileKind::Data - } else { - PhysicalFileKind::Other - } - } + _ => PhysicalFileKind::Other, } } @@ -803,13 +802,31 @@ pub async fn collect_physical_files_summary( table_location: &str, partition_depth: usize, ) -> crate::Result { + collect_physical_files_summary_with_options( + file_io, + table_location, + partition_depth, + &CoreOptions::new(&HashMap::new()), + ) + .await +} + +/// Include the configured data directory, including directories outside the table root. +pub async fn collect_physical_files_summary_with_options( + file_io: &FileIO, + table_location: &str, + partition_depth: usize, + options: &CoreOptions<'_>, +) -> crate::Result { + let data_directory = options.data_file_path_directory(); + let data_root = crate::spec::data_file_path(table_location, data_directory); // List top-level entries to discover subdirectories and top-level files let top_entries = match file_io.list_status(table_location).await { Ok(s) => s, Err(crate::Error::IoUnexpected { ref source, .. }) if source.kind() == opendal::ErrorKind::NotFound => { - return Ok(PhysicalFilesSummary::default()); + Vec::new() } Err(e) => return Err(e), }; @@ -818,16 +835,21 @@ pub async fn collect_physical_files_summary( // Classify top-level files directly let mut sub_dirs = Vec::new(); + if table_relative_path(table_location, &data_root).is_none() { + sub_dirs.push(data_root); + } + let mut seen = HashSet::new(); for entry in &top_entries { if entry.is_dir { sub_dirs.push(entry.path.clone()); - } else { + } else if seen.insert(entry.path.clone()) { accumulate_file( &mut summary, table_location, &entry.path, partition_depth, entry.size, + data_directory, ); } } @@ -852,12 +874,16 @@ pub async fn collect_physical_files_summary( for result in dir_results { let statuses = result?; for status in &statuses { + if !seen.insert(status.path.clone()) { + continue; + } accumulate_file( &mut summary, table_location, &status.path, partition_depth, status.size, + data_directory, ); } } @@ -871,8 +897,9 @@ fn accumulate_file( path: &str, partition_depth: usize, size: u64, + data_directory: Option<&str>, ) { - match classify_physical_path(table_location, path, partition_depth, None) { + match classify_physical_path(table_location, path, partition_depth, data_directory) { PhysicalFileKind::Manifest | PhysicalFileKind::Statistics => { summary.manifest_file_count += 1; summary.manifest_file_size += size as i64; @@ -891,6 +918,49 @@ fn accumulate_file( #[cfg(test)] mod tests { + #[tokio::test] + async fn configured_directory_counts_data_and_bucket_indexes_once() { + for directory in ["data/nested", "memory:/outside-data", ".."] { + let io = test_file_io(); + let table_path = "memory:/sizes/table"; + let options = + HashMap::from([("data-file.path-directory".into(), directory.to_string())]); + let root = crate::spec::data_file_path(table_path, Some(directory)); + write_test_file(&io, &format!("{table_path}/manifest/manifest-1"), "meta").await; + write_test_file( + &io, + &format!("{root}/pt=a/bucket-0/data.parquet"), + "payload", + ) + .await; + write_test_file(&io, &format!("{root}/pt=a/bucket-0/index-dv"), "dv").await; + write_test_file(&io, &format!("{table_path}/bucket-0/unrelated"), "ignore").await; + let sizes = collect_physical_files_summary_with_options( + &io, + table_path, + 1, + &CoreOptions::new(&options), + ) + .await + .unwrap(); + assert_eq!( + (sizes.manifest_file_count, sizes.manifest_file_size), + (1, 4), + "{directory}" + ); + assert_eq!( + (sizes.data_file_count, sizes.data_file_size), + (1, 7), + "{directory}" + ); + assert_eq!( + (sizes.index_file_count, sizes.index_file_size), + (1, 2), + "{directory}" + ); + } + } + use super::*; use crate::io::FileIOBuilder; use crate::spec::{CommitKind, Snapshot}; diff --git a/crates/paimon/src/table/sorted_global_index_build_builder/planning.rs b/crates/paimon/src/table/sorted_global_index_build_builder/planning.rs index 7c52fd806..7e487b113 100644 --- a/crates/paimon/src/table/sorted_global_index_build_builder/planning.rs +++ b/crates/paimon/src/table/sorted_global_index_build_builder/planning.rs @@ -227,7 +227,9 @@ fn bucket_path( partition: &BinaryRow, bucket: i32, ) -> Result { - let base = table_location.trim_end_matches('/'); + let data_root = + crate::spec::data_file_path(table_location, core_options.data_file_path_directory()); + let base = data_root.trim_end_matches('/'); if partition_keys.is_empty() { return Ok(format!("{base}/{}", bucket_dir_name(bucket))); } diff --git a/crates/paimon/src/table/table_commit.rs b/crates/paimon/src/table/table_commit.rs index ec3049dc8..e56b9477b 100644 --- a/crates/paimon/src/table/table_commit.rs +++ b/crates/paimon/src/table/table_commit.rs @@ -943,7 +943,11 @@ impl TableCommit { // An unpartitioned table's buckets sit directly under the table path, // so the partition blob is never decoded — callers are free to pass an // empty one. - return Ok(bucket_path_under(self.table.location(), "", bucket)); + return Ok(bucket_path_under( + &self.table.data_file_location(), + "", + bucket, + )); } let core_options = CoreOptions::new(self.table.schema().options()); let computer = PartitionComputer::new( @@ -953,7 +957,7 @@ impl TableCommit { core_options.legacy_partition_name(), )?; bucket_path( - self.table.location(), + &self.table.data_file_location(), Some(&computer), &BinaryRow::from_serialized_bytes(partition)?, bucket, diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index e975dda79..4165b592e 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1498,7 +1498,7 @@ impl<'a> PaimonTableScan<'a> { /// for `scan.version`; the strict selectors mirror Java's typed /// `scan.snapshot-id` / `scan.tag-name` handling. pub async fn plan(&self) -> crate::Result { - self.ensure_query_auth_allowed()?; + self.validate_read_options()?; self.validate_shard_strategy()?; let data_evolution_read_field_ids = self.projected_read_field_ids()?; let snapshot = match super::time_travel::resolve_snapshot(self.table).await? { @@ -1511,7 +1511,7 @@ impl<'a> PaimonTableScan<'a> { /// Plan the full scan and return metadata-pruning trace counters. pub async fn plan_with_trace(&self) -> crate::Result<(Plan, ScanTrace)> { - self.ensure_query_auth_allowed()?; + self.validate_read_options()?; self.validate_shard_strategy()?; let mut trace = ScanTrace { limit: self.limit, @@ -1537,8 +1537,10 @@ impl<'a> PaimonTableScan<'a> { /// Fail closed for a `query-auth.enabled` table: scan planning — including /// `with_scan_all_files`, which read-facing system tables like `files` use — /// exposes file paths, row counts, and stats the client can't authorize. - fn ensure_query_auth_allowed(&self) -> crate::Result<()> { - CoreOptions::new(self.table.schema().options()).ensure_read_authorized() + fn validate_read_options(&self) -> crate::Result<()> { + let options = CoreOptions::new(self.table.schema().options()); + options.ensure_read_authorized()?; + options.validate_data_file_path_directory() } fn validate_shard_strategy(&self) -> crate::Result<()> { @@ -1886,7 +1888,7 @@ impl<'a> PaimonTableScan<'a> { /// Reuses the same split-building path as a full snapshot plan, but only /// reads the delta manifest list and keeps ADD entries. pub(crate) async fn plan_snapshot_delta(&self, snapshot: &Snapshot) -> crate::Result { - self.ensure_query_auth_allowed()?; + self.validate_read_options()?; let data_evolution_read_field_ids = self.projected_read_field_ids()?; let mut scan = self.clone(); scan.incremental_split_mode = Some(IncrementalSplitMode::Streaming); @@ -1905,7 +1907,7 @@ impl<'a> PaimonTableScan<'a> { snapshots: &[Snapshot], end_snapshot: &Snapshot, ) -> crate::Result { - self.ensure_query_auth_allowed()?; + self.validate_read_options()?; self.validate_shard_strategy()?; let data_evolution_read_field_ids = self.projected_read_field_ids()?; let mut scan = self.clone(); @@ -1925,7 +1927,7 @@ impl<'a> PaimonTableScan<'a> { /// reads the changelog manifest list and keeps ADD entries. Snapshots /// without a changelog list yield an empty plan. pub(crate) async fn plan_snapshot_changelog(&self, snapshot: &Snapshot) -> crate::Result { - self.ensure_query_auth_allowed()?; + self.validate_read_options()?; let Some(list_name) = snapshot.changelog_manifest_list() else { return Ok(Plan::new(Vec::new()).with_snapshot_id(snapshot.id())); }; @@ -2008,7 +2010,7 @@ impl<'a> PaimonTableScan<'a> { before: &Snapshot, after: &Snapshot, ) -> crate::Result<(Plan, Plan)> { - self.ensure_query_auth_allowed()?; + self.validate_read_options()?; let core_options = CoreOptions::new(self.table.schema().options()); // Both forms: without `first_row_id` a filter stays an ordinary data // predicate instead of becoming a row range. @@ -2464,7 +2466,7 @@ impl<'a> PaimonTableScan<'a> { 'groups: for ((partition, bucket), (total_buckets, data_files)) in groups { let partition_row = BinaryRow::from_serialized_bytes(&partition)?; let bucket_path = bucket_path( - base_path, + &self.table.data_file_location(), partition_computer.as_ref(), &partition_row, bucket, diff --git a/crates/paimon/src/table/table_write.rs b/crates/paimon/src/table/table_write.rs index 1f28cc6df..4e78c42e1 100644 --- a/crates/paimon/src/table/table_write.rs +++ b/crates/paimon/src/table/table_write.rs @@ -245,6 +245,7 @@ impl TableWrite { // A dynamic-bucket write reads the persisted PK hash index; the rest are // refused too, since their commit is blocked anyway. CoreOptions::new(table.schema().options()).ensure_read_authorized()?; + CoreOptions::new(table.schema().options()).validate_data_file_path_directory()?; let is_overwrite = false; let schema = table.schema(); let write_schema = build_target_arrow_schema(schema.fields())?; @@ -405,20 +406,23 @@ impl TableWrite { merge_engine, ))) } else if is_dynamic_bucket { - BucketAssignerEnum::Dynamic(Box::new(DynamicBucketAssigner::new( - partition_field_indices, - primary_key_indices.clone(), - schema.fields().to_vec(), - target_bucket_row_number, - table.file_io().clone(), - table.location().to_string(), - is_overwrite, - // The same computer this writer already built: a hash index kept in - // the data-file directory must land in the directory the writer and - // the reader both derive, so both must agree on partition naming. - partition_computer.clone(), - core_options.index_file_in_data_file_dir(), - ))) + BucketAssignerEnum::Dynamic(Box::new( + DynamicBucketAssigner::new( + partition_field_indices, + primary_key_indices.clone(), + schema.fields().to_vec(), + target_bucket_row_number, + table.file_io().clone(), + table.location().to_string(), + is_overwrite, + // The same computer this writer already built: a hash index kept in + // the data-file directory must land in the directory the writer and + // the reader both derive, so both must agree on partition naming. + partition_computer.clone(), + core_options.index_file_in_data_file_dir(), + ) + .with_data_file_path_directory(core_options.data_file_path_directory()), + )) } else if total_buckets == POSTPONE_BUCKET { BucketAssignerEnum::Constant(ConstantBucketAssigner::new( partition_field_indices, @@ -1202,7 +1206,7 @@ impl TableWrite { PostponeFileWriter::new( self.table.file_io().clone(), PostponeWriteConfig { - table_location: self.table.location().to_string(), + table_location: self.table.data_file_location(), partition_path, bucket, schema_id: self.schema_id, diff --git a/crates/paimon/tests/data_file_directory_test.rs b/crates/paimon/tests/data_file_directory_test.rs new file mode 100644 index 000000000..e5847ab34 --- /dev/null +++ b/crates/paimon/tests/data_file_directory_test.rs @@ -0,0 +1,470 @@ +// 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::{ArrayRef, Int32Array, 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::Table; +use std::sync::Arc; + +async fn table(primary_key: bool, evolution: bool, directory: &str, bucket_indexes: bool) -> Table { + let path = "memory:/directory_test"; + let mut schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("pt", DataType::VarChar(VarCharType::string_type())) + .column("value", DataType::Int(IntType::new())) + .partition_keys(["pt"]) + .option("data-file.path-directory", directory) + .option("target-file-row-num", "2"); + if primary_key { + schema = schema + .primary_key(["id", "pt"]) + .option("bucket", "1") + .option("changelog-producer", "input"); + } + if evolution { + schema = schema + .option("row-tracking.enabled", "true") + .option("data-evolution.enabled", "true") + .option("deletion-vectors.enabled", "true") + .option("index-file-in-data-file-dir", bucket_indexes.to_string()); + } + let (io, table) = memory_table(path, TableSchema::new(0, &schema.build().unwrap())); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + table +} + +fn batch(ids: Vec, parts: Vec>, values: Vec) -> RecordBatch { + RecordBatch::try_from_iter([ + ("id", Arc::new(Int32Array::from(ids)) as ArrayRef), + ("pt", Arc::new(StringArray::from(parts)) as ArrayRef), + ("value", Arc::new(Int32Array::from(values)) as ArrayRef), + ]) + .unwrap() +} + +async fn values(table: &Table) -> Vec { + let builder = table.new_read_builder(); + let plan = builder.new_scan().plan().await.unwrap(); + for split in plan.splits() { + assert!( + split.bucket_path().contains("/data/nested/pt="), + "{}", + split.bucket_path() + ); + for file in split.data_files() { + assert!(table + .file_io() + .exists(&file.data_file_path(split.bucket_path())) + .await + .unwrap()); + } + } + let batches: Vec = builder + .new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect() + .await + .unwrap(); + let mut result: Vec = batches + .iter() + .flat_map(|batch| { + batch + .column_by_name("value") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }) + .collect(); + result.sort(); + result +} + +#[tokio::test] +async fn configured_directory_contains_append_and_primary_key_files() { + for primary_key in [false, true] { + let table = table(primary_key, false, "data/nested", false).await; + write_batch( + &table, + &batch( + vec![1, 2, 3], + vec![ + Some("a/b"), + Some("a/b"), + if primary_key { Some("c") } else { None }, + ], + vec![10, 20, 30], + ), + ) + .await; + assert_eq!(values(&table).await, vec![10, 20, 30]); + assert!(table + .file_io() + .exists("memory:/directory_test/snapshot/snapshot-1") + .await + .unwrap()); + assert!(!table + .file_io() + .exists("memory:/directory_test/data/nested/snapshot/snapshot-1") + .await + .unwrap()); + } +} + +#[tokio::test] +async fn relocated_updates_upserts_and_repeated_deletes_share_paths() { + for bucket_indexes in [false, true] { + let table = table(false, true, "data/nested", bucket_indexes).await; + write_batch( + &table, + &batch(vec![1, 2, 3], vec![Some("a"); 3], vec![10, 20, 30]), + ) + .await; + let mut update = table.new_write_builder().new_update().unwrap(); + update.with_update_type(vec!["value".into()]).unwrap(); + let input = RecordBatch::try_from_iter([ + ( + "_ROW_ID", + Arc::new(arrow_array::Int64Array::from(vec![0])) as ArrayRef, + ), + ("value", Arc::new(Int32Array::from(vec![11])) as ArrayRef), + ]) + .unwrap(); + let messages = update + .update_by_arrow_with_row_id(vec![input]) + .await + .unwrap(); + table + .new_write_builder() + .new_commit() + .commit(messages) + .await + .unwrap(); + assert_eq!(values(&table).await, vec![11, 20, 30]); + let messages = update + .upsert_by_arrow_with_key( + vec![batch(vec![2, 4], vec![Some("a"), Some("b")], vec![22, 40])], + vec!["id".into()], + ) + .await + .unwrap(); + table + .new_write_builder() + .new_commit() + .commit(messages) + .await + .unwrap(); + assert_eq!(values(&table).await, vec![11, 22, 30, 40]); + for (row_id, expected) in [(0, vec![22, 30, 40]), (1, vec![30, 40])] { + let messages = update.delete_by_row_id(vec![row_id]).await.unwrap(); + for message in &messages { + for index in &message.new_index_files { + let directory = if bucket_indexes { + "data/nested/pt=a/bucket-0" + } else { + "index" + }; + assert!(table + .file_io() + .exists(&format!( + "{}/{directory}/{}", + table.location(), + index.file_name + )) + .await + .unwrap()); + } + } + table + .new_write_builder() + .new_commit() + .commit(messages) + .await + .unwrap(); + assert_eq!(values(&table).await, expected); + } + // Abort deletes only newly staged files, leaving the committed DV live. + let messages = update.delete_by_row_id(vec![2]).await.unwrap(); + let paths: Vec = messages + .iter() + .flat_map(|message| { + message.new_index_files.iter().map(|index| { + let directory = if bucket_indexes { + "data/nested/pt=a/bucket-0" + } else { + "index" + }; + format!("{}/{directory}/{}", table.location(), index.file_name) + }) + }) + .collect(); + assert!(!paths.is_empty()); + table + .new_write_builder() + .new_commit() + .abort(&messages) + .await + .unwrap(); + for path in paths { + assert!(!table.file_io().exists(&path).await.unwrap()); + } + assert_eq!(values(&table).await, vec![30, 40]); + } +} + +#[tokio::test] +async fn abort_removes_relocated_data_and_changelog_files() { + for primary_key in [false, true] { + let table = table(primary_key, false, "data/nested", false).await; + let mut writer = table.new_write_builder().new_write().unwrap(); + writer + .write_arrow_batch(&batch(vec![1], vec![Some("a")], vec![10])) + .await + .unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + let mut paths = Vec::new(); + for message in &messages { + if primary_key { + assert!(!message.new_changelog_files.is_empty()); + } + for file in message.new_files.iter().chain(&message.new_changelog_files) { + let path = format!( + "{}/data/nested/pt=a/bucket-{}/{}", + table.location(), + message.bucket, + file.file_name + ); + assert!(table.file_io().exists(&path).await.unwrap()); + paths.push(path); + } + } + assert!(!paths.is_empty()); + table + .new_write_builder() + .new_commit() + .abort(&messages) + .await + .unwrap(); + for path in paths { + assert!(!table.file_io().exists(&path).await.unwrap()); + } + assert!(table + .snapshot_manager() + .get_latest_snapshot() + .await + .unwrap() + .is_none()); + } +} + +#[tokio::test] +async fn managed_blob_packs_and_parquet_share_relocated_bucket() { + use arrow_array::{Array, LargeBinaryArray}; + use paimon::spec::BlobType; + for primary_key in [false, true] { + let path = "memory:/blob_directory"; + let mut builder = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("payload", DataType::Blob(BlobType::new())) + .option("data-file.path-directory", "data/nested"); + if primary_key { + builder = builder.primary_key(["id"]).option("bucket", "1"); + } else { + builder = builder + .option("row-tracking.enabled", "true") + .option("data-evolution.enabled", "true"); + } + let (io, table) = memory_table(path, TableSchema::new(0, &builder.build().unwrap())); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + let input = RecordBatch::try_from_iter([ + ("id", Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef), + ( + "payload", + Arc::new(LargeBinaryArray::from(vec![ + Some(b"hello".as_slice()), + None, + Some(b"".as_slice()), + ])) as ArrayRef, + ), + ]) + .unwrap(); + write_batch(&table, &input).await; + let read = table.new_read_builder(); + let plan = read.new_scan().plan().await.unwrap(); + assert!(plan + .splits() + .iter() + .all(|split| split.bucket_path() == format!("{path}/data/nested/bucket-0"))); + let batches: Vec = read + .new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect() + .await + .unwrap(); + let mut rows = Vec::new(); + for batch in batches { + let ids = batch + .column_by_name("id") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + let payload = batch + .column_by_name("payload") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + for row in 0..batch.num_rows() { + rows.push(( + ids.value(row), + (!payload.is_null(row)).then(|| payload.value(row).to_vec()), + )); + } + } + rows.sort(); + assert_eq!( + rows, + vec![(1, Some(b"hello".to_vec())), (2, None), (3, Some(vec![]))] + ); + } +} + +#[tokio::test] +async fn dynamic_bucket_hash_indexes_follow_data_directory() { + let path = "memory:/dynamic_directory"; + let schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("pt", DataType::VarChar(VarCharType::string_type())) + .column("value", DataType::Int(IntType::new())) + .partition_keys(["pt"]) + .primary_key(["pt", "id"]) + .option("bucket", "-1") + .option("dynamic-bucket.target-row-num", "1") + .option("index-file-in-data-file-dir", "true") + .option("data-file.path-directory", "data/nested") + .build() + .unwrap(); + let (io, table) = memory_table(path, TableSchema::new(0, &schema)); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + for value in [10, 20] { + let mut writer = table.new_write_builder().new_write().unwrap(); + writer + .write_arrow_batch(&batch(vec![1, 2, 3], vec![Some("a"); 3], vec![value; 3])) + .await + .unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + for message in &messages { + for index in &message.new_index_files { + assert_eq!(index.index_type, "HASH"); + assert!(io + .exists(&format!( + "{path}/data/nested/pt=a/bucket-{}/{}", + message.bucket, index.file_name + )) + .await + .unwrap()); + } + } + table + .new_write_builder() + .new_commit() + .commit(messages) + .await + .unwrap(); + assert_eq!(values(&table).await, vec![value; 3]); + } +} + +#[tokio::test] +async fn referenced_sidecars_use_the_configured_directory() { + use paimon::table::referenced_files::collect_referenced_files_summary_with_options; + let table = table(false, false, "data/nested", false) + .await + .copy_with_options(std::collections::HashMap::from([ + ("file-index.bloom-filter.columns".into(), "value".into()), + ("file-index.bloom-filter.value.items".into(), "10".into()), + ("file-index.in-manifest-threshold".into(), "0 B".into()), + ])); + write_batch(&table, &batch(vec![1, 2], vec![Some("a"); 2], vec![10, 20])).await; + let plan = table.new_read_builder().new_scan().plan().await.unwrap(); + let mut expected_size = 0; + let mut expected_count = 0; + for split in plan.splits() { + for file in split.data_files() { + assert!(!file.extra_files.is_empty()); + for path in file.collect_files(split.bucket_path()) { + expected_size += table + .file_io() + .new_input(&path) + .unwrap() + .metadata() + .await + .unwrap() + .size as i64; + expected_count += 1; + } + } + } + let summaries = collect_referenced_files_summary_with_options( + table.file_io(), + table.location(), + table.schema().partition_keys(), + table.schema().fields(), + &table.schema().core_options(), + ) + .await + .unwrap(); + let total = summaries.iter().find(|row| row.source == "total").unwrap(); + assert_eq!(total.data_file_size, expected_size); + assert_eq!(total.data_file_count, expected_count); +} + +#[test] +fn directory_is_immutable_and_cannot_be_empty() { + use std::collections::HashMap; + for directory in [None, Some("data/nested")] { + let mut builder = Schema::builder().column("id", DataType::Int(IntType::new())); + if let Some(directory) = directory { + builder = builder.option("data-file.path-directory", directory); + } + let schema = TableSchema::new(0, &builder.build().unwrap()); + let copied = schema.copy_with_options(HashMap::from([( + "data-file.path-directory".into(), + "other".into(), + )])); + assert_eq!(copied.core_options().data_file_path_directory(), directory); + } + assert!(Schema::builder() + .column("id", DataType::Int(IntType::new())) + .option("data-file.path-directory", "") + .build() + .is_err()); +} From 2a46e49fe7433165bf066c1e4c6abac8cf3cec60 Mon Sep 17 00:00:00 2001 From: leaves12138 <41894543+leaves12138@users.noreply.github.com> Date: Mon, 28 Sep 2026 08:44:40 +0800 Subject: [PATCH 2/2] fix(core): preserve Java URI semantics when normalizing data directories --- crates/paimon/src/spec/partition_utils.rs | 37 ++++++++++++++++++----- 1 file changed, 30 insertions(+), 7 deletions(-) diff --git a/crates/paimon/src/spec/partition_utils.rs b/crates/paimon/src/spec/partition_utils.rs index 96c0d78ac..7c3dda8c1 100644 --- a/crates/paimon/src/spec/partition_utils.rs +++ b/crates/paimon/src/spec/partition_utils.rs @@ -64,8 +64,20 @@ fn resolve_data_path(parent: &str, child: &str, windows: bool) -> String { } let leading = if path.starts_with('/') { "/" } else { "" }; let normalized = format!("{leading}{}", components.join("/")); - if windows && has_drive(&normalized) && normalized.len() == 3 && path != normalized { + if windows + && normalized.starts_with('/') + && has_drive(&normalized) + && normalized.len() == 3 + && path != normalized + { format!("{normalized}/") + } else if !normalized.starts_with('/') + && normalized + .split('/') + .next() + .is_some_and(|part| part.contains(':')) + { + format!("./{normalized}") } else { normalized } @@ -133,12 +145,13 @@ fn resolve_data_path(parent: &str, child: &str, windows: bool) -> String { } else if child_path.starts_with('/') { (parent_scheme, parent_authority, child_path) } else { - let path = - if parent_path.is_empty() && parent_scheme.is_empty() && parent_authority.is_empty() { - child_path - } else { - format!("{}/{child_path}", parent_path.trim_end_matches('/')) - }; + let path = if parent_path.is_empty() + && (child_path.is_empty() || (parent_scheme.is_empty() && parent_authority.is_empty())) + { + child_path + } else { + format!("{}/{child_path}", parent_path.trim_end_matches('/')) + }; (parent_scheme, parent_authority, path) }; let path = normalize(&path, windows); @@ -830,6 +843,8 @@ mod tests { ("/t", "..", "/"), ("/t", "../../data", "/../data"), ("s3://bucket/", "data", "s3://bucket/data"), + ("s3://bucket", ".", "s3://bucket"), + ("s3://bucket", "x/..", "s3://bucket"), ("warehouse", "a/../b:c", "warehouse/b:c"), ("memory:/t", "a/../../b", "memory:/b"), ("/t", "///data", "/data"), @@ -840,6 +855,8 @@ mod tests { ("s3://bucket/t", "//other:9000/data", "s3://other:9000/data"), ("//host:8020/table", "//other/data", "//other/data"), ("relative", "../b:c", "./b:c"), + ("relative", "../a:bb", "./a:bb"), + ("relative", "../b:c/x", "./b:c/x"), ("file:/t", "file:/data", "file:/data"), ] { assert_eq!( @@ -853,6 +870,12 @@ mod tests { #[test] fn windows_data_directories_match_java() { for (parent, child, expected) in [ + ("relative", "../b:c", "./b:c"), + ("relative", "../a:bb", "./a:bb"), + ("relative", "../b:c/x", "./b:c/x"), + ("warehouse/t", "data/../../../b:c", "./b:c"), + (".", "a/../b:c", "./b:c"), + ("s3://bucket", ".", "s3://bucket"), (r"C:\warehouse\table", "../data", "C:/warehouse/data"), (r"C:\warehouse\table", r"..\data", "C:/warehouse/data"), ("C:/warehouse/table", "/data", "/data"),