From 52c318b83b0ca63d1a103398916695348fd6cce6 Mon Sep 17 00:00:00 2001 From: Ray Liu <257669749+blackmwk@users.noreply.github.com> Date: Fri, 28 Aug 2026 19:19:34 +0800 Subject: [PATCH] refactor(hms): derive catalog configuration from properties Use the shared property derive for HMS settings while retaining address and warehouse validation and transport fallback behavior. Tracks [apache/iceberg-rust#3098](https://github.com/apache/iceberg-rust/issues/3098). Generated-by: Codex --- Cargo.lock | 1 + crates/catalog/hms/Cargo.toml | 1 + crates/catalog/hms/src/catalog.rs | 113 ++++++++++++++++++++++-------- 3 files changed, 87 insertions(+), 28 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 3802cc03ff..a583e119b5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3855,6 +3855,7 @@ dependencies = [ "faststr", "hive_metastore", "iceberg", + "iceberg-property-macro", "iceberg-storage-opendal", "iceberg_test_utils", "linkedbytes", diff --git a/crates/catalog/hms/Cargo.toml b/crates/catalog/hms/Cargo.toml index ecd3a89cbe..7804f8874a 100644 --- a/crates/catalog/hms/Cargo.toml +++ b/crates/catalog/hms/Cargo.toml @@ -35,6 +35,7 @@ async-trait = { workspace = true } chrono = { workspace = true } hive_metastore = { workspace = true } iceberg = { workspace = true } +iceberg-property-macro = { workspace = true } pilota = { workspace = true } serde_json = { workspace = true } tokio = { workspace = true } diff --git a/crates/catalog/hms/src/catalog.rs b/crates/catalog/hms/src/catalog.rs index 7e5a2c391c..89e9f0d40d 100644 --- a/crates/catalog/hms/src/catalog.rs +++ b/crates/catalog/hms/src/catalog.rs @@ -34,6 +34,7 @@ use iceberg::{ Catalog, CatalogBuilder, Error, ErrorKind, MetadataLocation, Namespace, NamespaceIdent, Result, Runtime, TableCommit, TableCreation, TableIdent, }; +use iceberg_property_macro::Properties; use volo_thrift::MaybeException; use super::utils::*; @@ -103,35 +104,20 @@ impl CatalogBuilder for HmsCatalogBuilder { ) -> impl Future> + Send { self.config.name = Some(name.into()); - if props.contains_key(HMS_CATALOG_PROP_URI) { - self.config.address = props.get(HMS_CATALOG_PROP_URI).cloned().unwrap_or_default(); - } - - if let Some(tt) = props.get(HMS_CATALOG_PROP_THRIFT_TRANSPORT) { - self.config.thrift_transport = match tt.to_lowercase().as_str() { - THRIFT_TRANSPORT_FRAMED => HmsThriftTransport::Framed, - THRIFT_TRANSPORT_BUFFERED => HmsThriftTransport::Buffered, - _ => HmsThriftTransport::default(), - }; - } - - if props.contains_key(HMS_CATALOG_PROP_WAREHOUSE) { - self.config.warehouse = props - .get(HMS_CATALOG_PROP_WAREHOUSE) - .cloned() - .unwrap_or_default(); - } - - self.config.props = props - .into_iter() - .filter(|(k, _)| { - k != HMS_CATALOG_PROP_URI - && k != HMS_CATALOG_PROP_THRIFT_TRANSPORT - && k != HMS_CATALOG_PROP_WAREHOUSE - }) - .collect(); - async move { + let mut catalog_properties = HmsCatalogProperties::from_properties(&props)?; + for property in [ + HMS_CATALOG_PROP_URI, + HMS_CATALOG_PROP_THRIFT_TRANSPORT, + HMS_CATALOG_PROP_WAREHOUSE, + ] { + catalog_properties.props.remove(property); + } + self.config.address = catalog_properties.address; + self.config.thrift_transport = catalog_properties.thrift_transport; + self.config.warehouse = catalog_properties.warehouse; + self.config.props = catalog_properties.props; + let kms_client = match self.kms_client_factory { Some(factory) => Some(factory.create_kms_client(&self.config.props).await?), None => None, @@ -164,6 +150,77 @@ impl CatalogBuilder for HmsCatalogBuilder { } } +fn parse_thrift_transport(value: &str) -> Result { + Ok(match value.to_lowercase().as_str() { + THRIFT_TRANSPORT_FRAMED => HmsThriftTransport::Framed, + THRIFT_TRANSPORT_BUFFERED => HmsThriftTransport::Buffered, + _ => HmsThriftTransport::default(), + }) +} + +#[derive(Properties)] +struct HmsCatalogProperties { + #[property(key = HMS_CATALOG_PROP_URI, default = "")] + address: String, + #[property( + key = HMS_CATALOG_PROP_THRIFT_TRANSPORT, + default = HmsThriftTransport::default(), + parse_with = parse_thrift_transport + )] + thrift_transport: HmsThriftTransport, + #[property(key = HMS_CATALOG_PROP_WAREHOUSE, default = "")] + warehouse: String, + #[property(prefix = "")] + props: HashMap, +} + +#[cfg(test)] +mod catalog_properties_tests { + use super::*; + + #[test] + fn test_catalog_properties() { + let properties = HmsCatalogProperties::from_properties(&HashMap::from([ + ( + HMS_CATALOG_PROP_URI.to_string(), + "localhost:9083".to_string(), + ), + ( + HMS_CATALOG_PROP_THRIFT_TRANSPORT.to_string(), + "FRAMED".to_string(), + ), + ( + HMS_CATALOG_PROP_WAREHOUSE.to_string(), + "s3://warehouse".to_string(), + ), + ("custom.property".to_string(), "value".to_string()), + ])) + .unwrap(); + + assert_eq!(properties.address, "localhost:9083"); + assert!(matches!( + properties.thrift_transport, + HmsThriftTransport::Framed + )); + assert_eq!(properties.warehouse, "s3://warehouse"); + assert_eq!(properties.props["custom.property"], "value"); + } + + #[test] + fn test_invalid_thrift_transport_uses_default() { + let properties = HmsCatalogProperties::from_properties(&HashMap::from([( + HMS_CATALOG_PROP_THRIFT_TRANSPORT.to_string(), + "unknown".to_string(), + )])) + .unwrap(); + + assert!(matches!( + properties.thrift_transport, + HmsThriftTransport::Buffered + )); + } +} + /// Which variant of the thrift transport to communicate with HMS /// See: #[derive(Debug, Default)]