Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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)?;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
})
Expand Down
60 changes: 52 additions & 8 deletions crates/paimon/src/catalog/filesystem.rs
Original file line number Diff line number Diff line change
Expand Up @@ -704,8 +704,8 @@ fn local_datetime_to_millis<Tz: TimeZone>(
/// 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.
Expand Down Expand Up @@ -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"
),
});
Expand Down Expand Up @@ -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;
Expand Down
17 changes: 17 additions & 0 deletions crates/paimon/src/spec/core_options.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion crates/paimon/src/spec/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading