use std::collections::HashMap;
use std::fmt::Display;
use std::str::FromStr;
use crate::compression::CompressionCodec;
use crate::error::{Error, ErrorKind, Result};
fn parse_property<T: FromStr>(
properties: &HashMap<String, String>,
key: &str,
default: T,
) -> Result<T>
where
<T as FromStr>::Err: Display,
{
properties.get(key).map_or(Ok(default), |value| {
value.parse::<T>().map_err(|e| {
Error::new(
ErrorKind::DataInvalid,
format!("Invalid value for {key}: {e}"),
)
})
})
}
pub(crate) fn parse_metadata_file_compression(
properties: &HashMap<String, String>,
) -> Result<CompressionCodec> {
let value = properties
.get(TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC)
.map(|s| s.as_str())
.unwrap_or(TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC_DEFAULT);
if value.is_empty() {
return Ok(CompressionCodec::None);
}
let lowercase_value = value.to_lowercase();
let codec: CompressionCodec = serde_json::from_value(serde_json::Value::String(
lowercase_value,
))
.map_err(|_| {
Error::new(
ErrorKind::DataInvalid,
format!(
"Invalid metadata compression codec: {value}. Only '{}' and '{}' are supported.",
CompressionCodec::None.name(),
CompressionCodec::gzip_default().name()
),
)
})?;
match codec {
CompressionCodec::None | CompressionCodec::Gzip(_) => Ok(codec),
_ => Err(Error::new(
ErrorKind::DataInvalid,
format!(
"Invalid metadata compression codec: {value}. Only '{}' and '{}' are supported for metadata files.",
CompressionCodec::None.name(),
CompressionCodec::gzip_default().name()
),
)),
}
}
#[derive(Debug)]
pub struct TableProperties {
pub commit_num_retries: usize,
pub commit_min_retry_wait_ms: u64,
pub commit_max_retry_wait_ms: u64,
pub commit_total_retry_timeout_ms: u64,
pub write_format_default: String,
pub write_target_file_size_bytes: usize,
pub metadata_compression_codec: CompressionCodec,
pub write_datafusion_fanout_enabled: bool,
pub gc_enabled: bool,
pub max_snapshot_age_ms: i64,
pub min_snapshots_to_keep: usize,
pub max_ref_age_ms: i64,
pub cdc_enabled: bool,
pub cdc_min_chunk_size: usize,
pub cdc_max_chunk_size: usize,
pub cdc_norm_level: i32,
pub encryption_key_id: Option<String>,
pub encryption_data_key_length: usize,
}
impl TableProperties {
pub const PROPERTY_FORMAT_VERSION: &str = "format-version";
pub const PROPERTY_UUID: &str = "uuid";
pub const PROPERTY_SNAPSHOT_COUNT: &str = "snapshot-count";
pub const PROPERTY_CURRENT_SNAPSHOT_SUMMARY: &str = "current-snapshot-summary";
pub const PROPERTY_CURRENT_SNAPSHOT_ID: &str = "current-snapshot-id";
pub const PROPERTY_CURRENT_SNAPSHOT_TIMESTAMP: &str = "current-snapshot-timestamp-ms";
pub const PROPERTY_CURRENT_SCHEMA: &str = "current-schema";
pub const PROPERTY_DEFAULT_PARTITION_SPEC: &str = "default-partition-spec";
pub const PROPERTY_DEFAULT_SORT_ORDER: &str = "default-sort-order";
pub const PROPERTY_METADATA_PREVIOUS_VERSIONS_MAX: &str =
"write.metadata.previous-versions-max";
pub const PROPERTY_METADATA_PREVIOUS_VERSIONS_MAX_DEFAULT: usize = 100;
pub const PROPERTY_WRITE_PARTITION_SUMMARY_LIMIT: &str = "write.summary.partition-limit";
pub const PROPERTY_WRITE_PARTITION_SUMMARY_LIMIT_DEFAULT: u64 = 0;
pub const RESERVED_PROPERTIES: [&str; 9] = [
Self::PROPERTY_FORMAT_VERSION,
Self::PROPERTY_UUID,
Self::PROPERTY_SNAPSHOT_COUNT,
Self::PROPERTY_CURRENT_SNAPSHOT_ID,
Self::PROPERTY_CURRENT_SNAPSHOT_SUMMARY,
Self::PROPERTY_CURRENT_SNAPSHOT_TIMESTAMP,
Self::PROPERTY_CURRENT_SCHEMA,
Self::PROPERTY_DEFAULT_PARTITION_SPEC,
Self::PROPERTY_DEFAULT_SORT_ORDER,
];
pub const PROPERTY_COMMIT_NUM_RETRIES: &str = "commit.retry.num-retries";
pub const PROPERTY_COMMIT_NUM_RETRIES_DEFAULT: usize = 4;
pub const PROPERTY_COMMIT_MIN_RETRY_WAIT_MS: &str = "commit.retry.min-wait-ms";
pub const PROPERTY_COMMIT_MIN_RETRY_WAIT_MS_DEFAULT: u64 = 100;
pub const PROPERTY_COMMIT_MAX_RETRY_WAIT_MS: &str = "commit.retry.max-wait-ms";
pub const PROPERTY_COMMIT_MAX_RETRY_WAIT_MS_DEFAULT: u64 = 60 * 1000;
pub const PROPERTY_COMMIT_TOTAL_RETRY_TIME_MS: &str = "commit.retry.total-timeout-ms";
pub const PROPERTY_COMMIT_TOTAL_RETRY_TIME_MS_DEFAULT: u64 = 30 * 60 * 1000;
pub const PROPERTY_DEFAULT_FILE_FORMAT: &str = "write.format.default";
pub const PROPERTY_DELETE_DEFAULT_FILE_FORMAT: &str = "write.delete.format.default";
pub const PROPERTY_DEFAULT_FILE_FORMAT_DEFAULT: &str = "parquet";
pub const PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES: &str = "write.target-file-size-bytes";
pub const PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT: usize = 512 * 1024 * 1024;
pub const PROPERTY_METADATA_COMPRESSION_CODEC: &str = "write.metadata.compression-codec";
pub const PROPERTY_METADATA_COMPRESSION_CODEC_DEFAULT: &str = "none";
pub const PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED: &str = "write.datafusion.fanout.enabled";
pub const PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED_DEFAULT: bool = true;
pub const PROPERTY_GC_ENABLED: &str = "gc.enabled";
pub const PROPERTY_GC_ENABLED_DEFAULT: bool = true;
pub const PROPERTY_MAX_SNAPSHOT_AGE_MS: &str = "history.expire.max-snapshot-age-ms";
pub const PROPERTY_MAX_SNAPSHOT_AGE_MS_DEFAULT: i64 = 5 * 24 * 60 * 60 * 1000;
pub const PROPERTY_MIN_SNAPSHOTS_TO_KEEP: &str = "history.expire.min-snapshots-to-keep";
pub const PROPERTY_MIN_SNAPSHOTS_TO_KEEP_DEFAULT: usize = 1;
pub const PROPERTY_MAX_REF_AGE_MS: &str = "history.expire.max-ref-age-ms";
pub const PROPERTY_MAX_REF_AGE_MS_DEFAULT: i64 = i64::MAX;
pub const PROPERTY_PARQUET_CDC_ENABLED: &str = "write.parquet.content-defined-chunking.enabled";
pub const PROPERTY_PARQUET_CDC_ENABLED_DEFAULT: bool = false;
pub const PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE: &str =
"write.parquet.content-defined-chunking.min-chunk-size";
pub const PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE_DEFAULT: usize = 256 * 1024;
pub const PROPERTY_PARQUET_CDC_MAX_CHUNK_SIZE: &str =
"write.parquet.content-defined-chunking.max-chunk-size";
pub const PROPERTY_PARQUET_CDC_MAX_CHUNK_SIZE_DEFAULT: usize = 1024 * 1024;
pub const PROPERTY_PARQUET_CDC_NORM_LEVEL: &str =
"write.parquet.content-defined-chunking.norm-level";
pub const PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT: i32 = 0;
pub const PROPERTY_ENCRYPTION_KEY_ID: &str = "encryption.key-id";
pub const PROPERTY_ENCRYPTION_DATA_KEY_LENGTH: &str = "encryption.data-key-length";
pub const PROPERTY_ENCRYPTION_DATA_KEY_LENGTH_DEFAULT: usize = 16;
}
impl TryFrom<&HashMap<String, String>> for TableProperties {
type Error = Error;
fn try_from(props: &HashMap<String, String>) -> Result<Self> {
Ok(TableProperties {
commit_num_retries: parse_property(
props,
TableProperties::PROPERTY_COMMIT_NUM_RETRIES,
TableProperties::PROPERTY_COMMIT_NUM_RETRIES_DEFAULT,
)?,
commit_min_retry_wait_ms: parse_property(
props,
TableProperties::PROPERTY_COMMIT_MIN_RETRY_WAIT_MS,
TableProperties::PROPERTY_COMMIT_MIN_RETRY_WAIT_MS_DEFAULT,
)?,
commit_max_retry_wait_ms: parse_property(
props,
TableProperties::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS,
TableProperties::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS_DEFAULT,
)?,
commit_total_retry_timeout_ms: parse_property(
props,
TableProperties::PROPERTY_COMMIT_TOTAL_RETRY_TIME_MS,
TableProperties::PROPERTY_COMMIT_TOTAL_RETRY_TIME_MS_DEFAULT,
)?,
write_format_default: parse_property(
props,
TableProperties::PROPERTY_DEFAULT_FILE_FORMAT,
TableProperties::PROPERTY_DEFAULT_FILE_FORMAT_DEFAULT.to_string(),
)?,
write_target_file_size_bytes: parse_property(
props,
TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES,
TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT,
)?,
metadata_compression_codec: parse_metadata_file_compression(props)?,
write_datafusion_fanout_enabled: parse_property(
props,
TableProperties::PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED,
TableProperties::PROPERTY_DATAFUSION_WRITE_FANOUT_ENABLED_DEFAULT,
)?,
gc_enabled: parse_property(
props,
TableProperties::PROPERTY_GC_ENABLED,
TableProperties::PROPERTY_GC_ENABLED_DEFAULT,
)?,
max_snapshot_age_ms: parse_property(
props,
TableProperties::PROPERTY_MAX_SNAPSHOT_AGE_MS,
TableProperties::PROPERTY_MAX_SNAPSHOT_AGE_MS_DEFAULT,
)?,
min_snapshots_to_keep: parse_property(
props,
TableProperties::PROPERTY_MIN_SNAPSHOTS_TO_KEEP,
TableProperties::PROPERTY_MIN_SNAPSHOTS_TO_KEEP_DEFAULT,
)?,
max_ref_age_ms: parse_property(
props,
TableProperties::PROPERTY_MAX_REF_AGE_MS,
TableProperties::PROPERTY_MAX_REF_AGE_MS_DEFAULT,
)?,
cdc_enabled: parse_property(
props,
TableProperties::PROPERTY_PARQUET_CDC_ENABLED,
TableProperties::PROPERTY_PARQUET_CDC_ENABLED_DEFAULT,
)?,
cdc_min_chunk_size: parse_property(
props,
TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE,
TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE_DEFAULT,
)?,
cdc_max_chunk_size: parse_property(
props,
TableProperties::PROPERTY_PARQUET_CDC_MAX_CHUNK_SIZE,
TableProperties::PROPERTY_PARQUET_CDC_MAX_CHUNK_SIZE_DEFAULT,
)?,
cdc_norm_level: parse_property(
props,
TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL,
TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL_DEFAULT,
)?,
encryption_key_id: props
.get(TableProperties::PROPERTY_ENCRYPTION_KEY_ID)
.cloned(),
encryption_data_key_length: parse_property(
props,
TableProperties::PROPERTY_ENCRYPTION_DATA_KEY_LENGTH,
TableProperties::PROPERTY_ENCRYPTION_DATA_KEY_LENGTH_DEFAULT,
)?,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::compression::CompressionCodec;
#[test]
fn test_table_properties_default() {
let props = HashMap::new();
let table_properties = TableProperties::try_from(&props).unwrap();
assert_eq!(
table_properties.commit_num_retries,
TableProperties::PROPERTY_COMMIT_NUM_RETRIES_DEFAULT
);
assert_eq!(
table_properties.commit_min_retry_wait_ms,
TableProperties::PROPERTY_COMMIT_MIN_RETRY_WAIT_MS_DEFAULT
);
assert_eq!(
table_properties.commit_max_retry_wait_ms,
TableProperties::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS_DEFAULT
);
assert_eq!(
table_properties.write_format_default,
TableProperties::PROPERTY_DEFAULT_FILE_FORMAT_DEFAULT.to_string()
);
assert_eq!(
table_properties.write_target_file_size_bytes,
TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT
);
assert_eq!(
table_properties.metadata_compression_codec,
CompressionCodec::None
);
assert_eq!(
table_properties.gc_enabled,
TableProperties::PROPERTY_GC_ENABLED_DEFAULT
);
assert_eq!(
table_properties.max_snapshot_age_ms,
TableProperties::PROPERTY_MAX_SNAPSHOT_AGE_MS_DEFAULT
);
assert_eq!(
table_properties.min_snapshots_to_keep,
TableProperties::PROPERTY_MIN_SNAPSHOTS_TO_KEEP_DEFAULT
);
assert_eq!(
table_properties.max_ref_age_ms,
TableProperties::PROPERTY_MAX_REF_AGE_MS_DEFAULT
);
}
#[test]
fn test_table_properties_history_expire_overrides() {
let props = HashMap::from([
(
TableProperties::PROPERTY_MAX_SNAPSHOT_AGE_MS.to_string(),
"1234".to_string(),
),
(
TableProperties::PROPERTY_MIN_SNAPSHOTS_TO_KEEP.to_string(),
"7".to_string(),
),
(
TableProperties::PROPERTY_MAX_REF_AGE_MS.to_string(),
"5678".to_string(),
),
]);
let table_properties = TableProperties::try_from(&props).unwrap();
assert_eq!(table_properties.max_snapshot_age_ms, 1234);
assert_eq!(table_properties.min_snapshots_to_keep, 7);
assert_eq!(table_properties.max_ref_age_ms, 5678);
}
#[test]
fn test_table_properties_compression() {
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"gzip".to_string(),
)]);
let table_properties = TableProperties::try_from(&props).unwrap();
assert_eq!(
table_properties.metadata_compression_codec,
CompressionCodec::gzip_default()
);
}
#[test]
fn test_table_properties_compression_none() {
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"none".to_string(),
)]);
let table_properties = TableProperties::try_from(&props).unwrap();
assert_eq!(
table_properties.metadata_compression_codec,
CompressionCodec::None
);
}
#[test]
fn test_table_properties_compression_case_insensitive() {
let props_upper = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"GZIP".to_string(),
)]);
let table_properties = TableProperties::try_from(&props_upper).unwrap();
assert_eq!(
table_properties.metadata_compression_codec,
CompressionCodec::gzip_default()
);
let props_mixed = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"GzIp".to_string(),
)]);
let table_properties = TableProperties::try_from(&props_mixed).unwrap();
assert_eq!(
table_properties.metadata_compression_codec,
CompressionCodec::gzip_default()
);
let props_none_upper = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"NONE".to_string(),
)]);
let table_properties = TableProperties::try_from(&props_none_upper).unwrap();
assert_eq!(
table_properties.metadata_compression_codec,
CompressionCodec::None
);
}
#[test]
fn test_table_properties_valid() {
let props = HashMap::from([
(
TableProperties::PROPERTY_COMMIT_NUM_RETRIES.to_string(),
"10".to_string(),
),
(
TableProperties::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS.to_string(),
"20".to_string(),
),
(
TableProperties::PROPERTY_DEFAULT_FILE_FORMAT.to_string(),
"avro".to_string(),
),
(
TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES.to_string(),
"512".to_string(),
),
(
TableProperties::PROPERTY_GC_ENABLED.to_string(),
"false".to_string(),
),
]);
let table_properties = TableProperties::try_from(&props).unwrap();
assert_eq!(table_properties.commit_num_retries, 10);
assert_eq!(table_properties.commit_max_retry_wait_ms, 20);
assert_eq!(table_properties.write_format_default, "avro".to_string());
assert_eq!(table_properties.write_target_file_size_bytes, 512);
assert!(!table_properties.gc_enabled);
}
#[test]
fn test_table_properties_invalid() {
let invalid_retries = HashMap::from([(
TableProperties::PROPERTY_COMMIT_NUM_RETRIES.to_string(),
"abc".to_string(),
)]);
let table_properties = TableProperties::try_from(&invalid_retries).unwrap_err();
assert!(
table_properties.to_string().contains(
"Invalid value for commit.retry.num-retries: invalid digit found in string"
)
);
let invalid_min_wait = HashMap::from([(
TableProperties::PROPERTY_COMMIT_MIN_RETRY_WAIT_MS.to_string(),
"abc".to_string(),
)]);
let table_properties = TableProperties::try_from(&invalid_min_wait).unwrap_err();
assert!(
table_properties.to_string().contains(
"Invalid value for commit.retry.min-wait-ms: invalid digit found in string"
)
);
let invalid_max_wait = HashMap::from([(
TableProperties::PROPERTY_COMMIT_MAX_RETRY_WAIT_MS.to_string(),
"abc".to_string(),
)]);
let table_properties = TableProperties::try_from(&invalid_max_wait).unwrap_err();
assert!(
table_properties.to_string().contains(
"Invalid value for commit.retry.max-wait-ms: invalid digit found in string"
)
);
let invalid_target_size = HashMap::from([(
TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES.to_string(),
"abc".to_string(),
)]);
let table_properties = TableProperties::try_from(&invalid_target_size).unwrap_err();
assert!(table_properties.to_string().contains(
"Invalid value for write.target-file-size-bytes: invalid digit found in string"
));
let invalid_gc_enabled = HashMap::from([(
TableProperties::PROPERTY_GC_ENABLED.to_string(),
"notabool".to_string(),
)]);
let table_properties = TableProperties::try_from(&invalid_gc_enabled).unwrap_err();
assert!(
table_properties
.to_string()
.contains("Invalid value for gc.enabled")
);
}
#[test]
fn test_table_properties_compression_invalid_rejected() {
let invalid_codecs = ["lz4", "zstd", "snappy"];
for codec in invalid_codecs {
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
codec.to_string(),
)]);
let err = TableProperties::try_from(&props).unwrap_err();
let err_msg = err.to_string();
assert!(
err_msg.contains(&format!("Invalid metadata compression codec: {codec}")),
"Expected error message to contain codec '{codec}', got: {err_msg}"
);
assert!(
err_msg.contains("Only 'none' and 'gzip' are supported"),
"Expected error message to contain supported codecs, got: {err_msg}"
);
}
}
#[test]
fn test_parse_metadata_file_compression_valid() {
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"none".to_string(),
)]);
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::None
);
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"".to_string(),
)]);
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::None
);
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"gzip".to_string(),
)]);
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::gzip_default()
);
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"NONE".to_string(),
)]);
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::None
);
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"GZIP".to_string(),
)]);
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::gzip_default()
);
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
"GzIp".to_string(),
)]);
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::gzip_default()
);
let props = HashMap::new();
assert_eq!(
parse_metadata_file_compression(&props).unwrap(),
CompressionCodec::None
);
}
#[test]
fn test_parse_metadata_file_compression_invalid() {
let invalid_codecs = ["lz4", "zstd", "snappy"];
for codec in invalid_codecs {
let props = HashMap::from([(
TableProperties::PROPERTY_METADATA_COMPRESSION_CODEC.to_string(),
codec.to_string(),
)]);
let err = parse_metadata_file_compression(&props).unwrap_err();
let err_msg = err.to_string();
assert!(
err_msg.contains("Invalid metadata compression codec"),
"Expected error message to contain 'Invalid metadata compression codec', got: {err_msg}"
);
assert!(
err_msg.contains("Only 'none' and 'gzip' are supported"),
"Expected error message to contain supported codecs, got: {err_msg}"
);
}
}
#[test]
fn test_cdc_disabled_by_default() {
let props = HashMap::new();
let tp = TableProperties::try_from(&props).unwrap();
assert!(!tp.cdc_enabled);
}
#[test]
fn test_cdc_enabled_via_flag() {
let props = HashMap::from([(
TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
"true".to_string(),
)]);
let tp = TableProperties::try_from(&props).unwrap();
assert!(tp.cdc_enabled);
assert_eq!(tp.cdc_min_chunk_size, 256 * 1024);
assert_eq!(tp.cdc_max_chunk_size, 1024 * 1024);
assert_eq!(tp.cdc_norm_level, 0);
}
#[test]
fn test_cdc_size_props_alone_do_not_enable() {
let props = HashMap::from([(
TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE.to_string(),
"262144".to_string(),
)]);
let tp = TableProperties::try_from(&props).unwrap();
assert!(!tp.cdc_enabled);
}
#[test]
fn test_cdc_custom_values() {
let props = HashMap::from([
(
TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
"true".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE.to_string(),
"200000".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_MAX_CHUNK_SIZE.to_string(),
"900000".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL.to_string(),
"1".to_string(),
),
]);
let tp = TableProperties::try_from(&props).unwrap();
assert!(tp.cdc_enabled);
assert_eq!(tp.cdc_min_chunk_size, 200000);
assert_eq!(tp.cdc_max_chunk_size, 900000);
assert_eq!(tp.cdc_norm_level, 1);
}
#[test]
fn test_cdc_partial_override() {
let props = HashMap::from([
(
TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
"true".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL.to_string(),
"2".to_string(),
),
]);
let tp = TableProperties::try_from(&props).unwrap();
assert!(tp.cdc_enabled);
assert_eq!(tp.cdc_min_chunk_size, 256 * 1024);
assert_eq!(tp.cdc_max_chunk_size, 1024 * 1024);
assert_eq!(tp.cdc_norm_level, 2);
}
#[test]
fn test_cdc_negative_norm_level() {
let props = HashMap::from([
(
TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
"true".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL.to_string(),
"-2".to_string(),
),
]);
let tp = TableProperties::try_from(&props).unwrap();
assert_eq!(tp.cdc_norm_level, -2);
}
#[test]
fn test_cdc_invalid_min_chunk_size() {
let props = HashMap::from([
(
TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
"true".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE.to_string(),
"not_a_number".to_string(),
),
]);
let err = TableProperties::try_from(&props).unwrap_err();
assert!(
err.to_string().contains(
"Invalid value for write.parquet.content-defined-chunking.min-chunk-size"
)
);
}
#[test]
fn test_cdc_invalid_norm_level() {
let props = HashMap::from([
(
TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
"true".to_string(),
),
(
TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL.to_string(),
"not_a_number".to_string(),
),
]);
let err = TableProperties::try_from(&props).unwrap_err();
assert!(
err.to_string()
.contains("Invalid value for write.parquet.content-defined-chunking.norm-level")
);
}
#[test]
fn test_cdc_no_properties() {
let props = HashMap::from([("some.other.property".to_string(), "value".to_string())]);
let tp = TableProperties::try_from(&props).unwrap();
assert!(!tp.cdc_enabled);
}
}