Skip to content
Open
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
104 changes: 104 additions & 0 deletions crates/paimon/src/spec/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1198,6 +1198,7 @@ impl Schema {
CoreOptions::new(options).table_type()?;
validate_no_reserved_field_names(fields)?;
Self::validate_key_field_types(fields, primary_keys, options)?;
Self::validate_bucket_count(options)?;
Self::validate_row_tracking(primary_keys, options)?;
Self::validate_blob_fields(fields, partition_keys, primary_keys, options)?;
Self::validate_primary_key_blob_configuration(fields, primary_keys, options)?;
Expand Down Expand Up @@ -1225,6 +1226,38 @@ impl Schema {
Ok(())
}

/// Reject an out-of-range or non-integer `bucket` option at create/alter time.
///
/// `bucket()` parses the option with `unwrap_or(DEFAULT_BUCKET)`, so `bucket='0'`
/// is taken verbatim (silently disabling read-side bucket pruning) and a
/// non-integer such as `bucket='abc'` silently falls back to dynamic (-1) rather
/// than the fixed bucketing the user asked for. Mirror Java `validateBucket`'s
/// "number of buckets needs to be greater than 0" check: allow -1 (dynamic),
/// -2 (postpone), and any value >= 1. (The Java `bucket=-1` + `bucket-key` arm is
/// intentionally not ported here; it is enforced separately by bucket-key
/// validation.)
fn validate_bucket_count(options: &HashMap<String, String>) -> crate::Result<()> {
let Some(raw) = options.get("bucket") else {
return Ok(());
};
// Parse the raw value exactly as `CoreOptions::bucket()` does (no trim):
// it uses `v.parse().ok().unwrap_or(-1)`, so a padded value like `" 4 "`
// silently runs as dynamic (-1). Validating a trimmed copy would accept
// `" 4 "` here yet persist a value the runtime reads as -1, so reject what
// the runtime cannot parse.
let bucket: i32 = raw.parse().map_err(|_| crate::Error::ConfigInvalid {
message: format!("Option 'bucket' must be an integer, got: '{raw}'."),
})?;
if bucket < 1 && bucket != -1 && bucket != -2 {
return Err(crate::Error::ConfigInvalid {
message: format!(
"The number of buckets needs to be greater than 0, got: {bucket}."
),
});
}
Ok(())
}

fn validate_row_tracking(
primary_keys: &[String],
options: &HashMap<String, String>,
Expand Down Expand Up @@ -5820,6 +5853,77 @@ mod tests {
}
}

#[test]
fn bucket_zero_is_rejected() {
// bucket='0' is taken verbatim and silently disables read-side bucket
// pruning; Java rejects it as "needs to be greater than 0".
let err = Schema::builder()
.column("id", DataType::Int(IntType::new()))
.option("bucket", "0")
.build()
.unwrap_err();
assert!(
matches!(err, crate::Error::ConfigInvalid { ref message }
if message.contains("greater than 0")),
"got {err:?}"
);
}

#[test]
fn bucket_non_integer_is_rejected() {
// A non-integer would silently fall back to dynamic (-1) rather than the
// fixed bucketing the user intended.
let err = Schema::builder()
.column("id", DataType::Int(IntType::new()))
.option("bucket", "abc")
.build()
.unwrap_err();
assert!(
matches!(err, crate::Error::ConfigInvalid { ref message }
if message.contains("must be an integer")),
"got {err:?}"
);
}

#[test]
fn bucket_dynamic_and_fixed_values_are_accepted() {
// -1 (dynamic), -2 (postpone) and any value >= 1 stay valid.
for value in ["-1", "-2", "1", "16"] {
Schema::builder()
.column("id", DataType::Int(IntType::new()))
.option("bucket", value)
.build()
.unwrap_or_else(|e| panic!("bucket={value} should be accepted, got {e:?}"));
}
}

#[test]
fn bucket_padded_value_is_rejected_matching_runtime() {
// `CoreOptions::bucket()` parses the raw value without trimming and falls
// back to -1, so a padded `" 4 "` silently runs as dynamic. Validation
// must use the same grammar and reject it, not persist a value the runtime
// reads as -1.
let padded = HashMap::from([("bucket".to_string(), " 4 ".to_string())]);
assert_eq!(
CoreOptions::new(&padded).bucket(),
-1,
"runtime reads a padded bucket value as dynamic"
);
let err = Schema::builder()
.column("id", DataType::Int(IntType::new()))
.option("bucket", " 4 ")
.build()
.unwrap_err();
assert!(
matches!(err, crate::Error::ConfigInvalid { ref message }
if message.contains("must be an integer")),
"got {err:?}"
);
// The accepted spelling runs as the fixed count the user asked for.
let fixed = HashMap::from([("bucket".to_string(), "4".to_string())]);
assert_eq!(CoreOptions::new(&fixed).bucket(), 4);
}

#[test]
fn rowkind_field_rejects_changelog_producer_input() {
let err = Schema::builder()
Expand Down
Loading