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 e5a90870 fix(spec): validate bucket-key against table schema (#620)
e5a90870 is described below
commit e5a908705483a45cbe04b222178959b1f2bedeb3
Author: jackylee <[email protected]>
AuthorDate: Tue Jul 28 23:11:32 2026 +0800
fix(spec): validate bucket-key against table schema (#620)
---
crates/paimon/src/spec/core_options.rs | 16 ++-
crates/paimon/src/spec/schema.rs | 192 +++++++++++++++++++++++++++++++++
2 files changed, 205 insertions(+), 3 deletions(-)
diff --git a/crates/paimon/src/spec/core_options.rs
b/crates/paimon/src/spec/core_options.rs
index fce04169..ec23ad69 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -783,10 +783,20 @@ impl<'a> CoreOptions<'a> {
}
/// Explicit bucket key columns. If not set, defaults to primary keys for
PK tables.
+ ///
+ /// Blank entries are dropped and an all-blank option resolves to `None`,
+ /// mirroring Java `TableSchema#originalBucketKeys`, so callers fall back
to
+ /// the primary keys instead of treating `""` as a column name.
pub fn bucket_key(&self) -> Option<Vec<String>> {
- self.options
- .get(BUCKET_KEY_OPTION)
- .map(|v| v.split(',').map(|s| s.trim().to_string()).collect())
+ let keys: Vec<String> = self
+ .options
+ .get(BUCKET_KEY_OPTION)?
+ .split(',')
+ .map(str::trim)
+ .filter(|key| !key.is_empty())
+ .map(str::to_string)
+ .collect();
+ (!keys.is_empty()).then_some(keys)
}
pub fn commit_max_retries(&self) -> u32 {
diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs
index 0245cf41..f96bc0d7 100644
--- a/crates/paimon/src/spec/schema.rs
+++ b/crates/paimon/src/spec/schema.rs
@@ -489,6 +489,12 @@ impl TableSchema {
&new_schema.primary_keys,
&new_schema.fields,
)?;
+ Schema::validate_bucket_keys(
+ &new_schema.options,
+ &new_schema.fields,
+ &new_schema.partition_keys,
+ &new_schema.primary_keys,
+ )?;
Schema::validate_read_batch_size(&new_schema.options)?;
Ok(new_schema)
}
@@ -936,6 +942,7 @@ impl Schema {
AggregationConfig::new(&options).validate_create_mode(&primary_keys,
&fields)?;
Self::validate_first_row_changelog_producer(&options)?;
Self::validate_rowkind_field(&options, &primary_keys, &fields)?;
+ Self::validate_bucket_keys(&options, &fields, &partition_keys,
&primary_keys)?;
Self::validate_read_batch_size(&options)?;
Ok(Self {
@@ -1351,6 +1358,57 @@ impl Schema {
Ok(())
}
+ /// Validate the explicit `bucket-key` option against the schema, mirroring
+ /// Java `TableSchema#originalBucketKeys`. A bucket key that is missing,
+ /// partitioned, or outside the primary key otherwise resolves to no field
+ /// index in `TableWrite`, which silently degrades to a constant bucket 0
+ /// assigner instead of hashing.
+ fn validate_bucket_keys(
+ options: &HashMap<String, String>,
+ fields: &[DataField],
+ partition_keys: &[String],
+ primary_keys: &[String],
+ ) -> crate::Result<()> {
+ let Some(bucket_keys) = CoreOptions::new(options).bucket_key() else {
+ return Ok(());
+ };
+
+ let mut seen: HashSet<&str> = HashSet::new();
+ for key in &bucket_keys {
+ if fields.iter().all(|f| f.name() != key) {
+ return Err(crate::Error::ConfigInvalid {
+ message: format!(
+ "Field names should contain all bucket keys, but
bucket key '{key}' \
+ can not be found in table schema."
+ ),
+ });
+ }
+ if !seen.insert(key.as_str()) {
+ return Err(crate::Error::ConfigInvalid {
+ message: format!("Bucket key '{key}' is defined
repeatedly."),
+ });
+ }
+ if partition_keys.contains(key) {
+ return Err(crate::Error::ConfigInvalid {
+ message: format!(
+ "Bucket keys should not be in partition keys, but
bucket key '{key}' \
+ is a partition field."
+ ),
+ });
+ }
+ if !primary_keys.is_empty() && !primary_keys.contains(key) {
+ return Err(crate::Error::ConfigInvalid {
+ message: format!(
+ "Primary keys {primary_keys:?} should contain all
bucket keys, but \
+ bucket key '{key}' is not a primary key field."
+ ),
+ });
+ }
+ }
+
+ Ok(())
+ }
+
fn options_error_to_config_invalid(error: crate::Error) -> crate::Error {
match error {
crate::Error::Unsupported { message } =>
crate::Error::ConfigInvalid { message },
@@ -2923,6 +2981,140 @@ mod tests {
);
}
+ #[test]
+ fn test_create_schema_rejects_unknown_bucket_key() {
+ let err = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::VarChar(VarCharType::string_type()))
+ .option("bucket", "4")
+ // typo: `nmae` instead of `name`
+ .option("bucket-key", "nmae")
+ .build()
+ .unwrap_err();
+
+ assert!(
+ matches!(err, crate::Error::ConfigInvalid { ref message }
+ if message.contains("nmae") && message.contains("can not be
found")),
+ "bucket key missing from the schema should be rejected, got
{err:?}"
+ );
+ }
+
+ #[test]
+ fn test_create_schema_rejects_repeated_bucket_key() {
+ let err = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::VarChar(VarCharType::string_type()))
+ .option("bucket", "4")
+ .option("bucket-key", "name,name")
+ .build()
+ .unwrap_err();
+
+ assert!(
+ matches!(err, crate::Error::ConfigInvalid { ref message }
+ if message.contains("name") && message.contains("repeatedly")),
+ "repeated bucket key should be rejected, got {err:?}"
+ );
+ }
+
+ #[test]
+ fn test_create_schema_rejects_partitioned_bucket_key() {
+ let err = Schema::builder()
+ .column("pt", DataType::Int(IntType::new()))
+ .column("id", DataType::Int(IntType::new()))
+ .partition_keys(["pt"])
+ .option("bucket", "4")
+ .option("bucket-key", "pt")
+ .build()
+ .unwrap_err();
+
+ assert!(
+ matches!(err, crate::Error::ConfigInvalid { ref message }
+ if message.contains("pt") && message.contains("partition")),
+ "partition field used as bucket key should be rejected, got
{err:?}"
+ );
+ }
+
+ #[test]
+ fn test_create_schema_rejects_bucket_key_outside_primary_key() {
+ let err = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::VarChar(VarCharType::string_type()))
+ .primary_key(["id"])
+ .option("bucket", "4")
+ .option("bucket-key", "name")
+ .build()
+ .unwrap_err();
+
+ assert!(
+ matches!(err, crate::Error::ConfigInvalid { ref message }
+ if message.contains("name") && message.contains("primary
key")),
+ "non-primary-key bucket key on a PK table should be rejected, got
{err:?}"
+ );
+ }
+
+ #[test]
+ fn test_create_schema_accepts_bucket_key_subset_of_primary_key() {
+ let schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::VarChar(VarCharType::string_type()))
+ .column("v", DataType::Int(IntType::new()))
+ .primary_key(["id", "name"])
+ .option("bucket", "4")
+ .option("bucket-key", "id")
+ .build();
+
+ assert!(
+ schema.is_ok(),
+ "a bucket key that is a primary key field should be accepted, got
{schema:?}"
+ );
+ }
+
+ #[test]
+ fn test_blank_bucket_key_falls_back_to_primary_keys() {
+ // A blank option must not resolve to a `""` column: `TableWrite` would
+ // find no field index for it and silently write every row to bucket 0.
+ let table_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::VarChar(VarCharType::string_type()))
+ .primary_key(["id"])
+ .option("bucket", "4")
+ .option("bucket-key", " ")
+ .build()
+ .unwrap(),
+ );
+
+ assert_eq!(table_schema.bucket_keys(), vec!["id".to_string()]);
+ }
+
+ #[test]
+ fn test_alter_set_unknown_bucket_key_rejected() {
+ let table_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::VarChar(VarCharType::string_type()))
+ .option("bucket", "4")
+ .option("bucket-key", "name")
+ .build()
+ .unwrap(),
+ );
+
+ let err = table_schema
+ .apply_changes(vec![crate::spec::SchemaChange::set_option(
+ BUCKET_KEY_OPTION.to_string(),
+ "nmae".to_string(),
+ )])
+ .unwrap_err();
+
+ assert!(
+ matches!(err, crate::Error::ConfigInvalid { ref message }
+ if message.contains("nmae") && message.contains("can not be
found")),
+ "alter setting an unknown bucket key should be rejected, got
{err:?}"
+ );
+ }
+
#[test]
fn test_drop_column_referenced_by_bucket_key_rejected() {
let table_schema = TableSchema::new(