laskoviymishka commented on code in PR #3111: URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4198203938
########## crates/storage/opendal/src/hdfs_native.rs: ########## @@ -0,0 +1,1244 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! HDFS storage backend via OpenDAL's `services-hdfs-native` (pure Rust, no JNI). + +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::sync::{Arc, RwLock, Weak}; + +use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_HOST, HDFS_NAME_NODE, HDFS_PORT}; +use iceberg::{Error, ErrorKind, Result}; +use opendal::Operator; +use opendal::services::HdfsNativeConfig; +use serde::{Deserialize, Serialize}; +use tokio::runtime::Handle; +use tokio::task::JoinHandle; +use url::Url; + +use crate::OpenDalClientConfig; +use crate::utils::from_opendal_error; + +/// Hadoop's default filesystem, which serves authority-less paths. +const FS_DEFAULT_FS: &str = "fs.defaultFS"; +const HDFS_DEFAULT_PORT: u16 = 8020; +/// PyIceberg keys with no equivalent in opendal's config. +const HDFS_UNSUPPORTED_KEYS: [&str; 2] = ["hdfs.user", "hdfs.kerberos_ticket"]; +/// Hadoop's keys declaring an HA nameservice: `hdfs.name-node.<nameservice>` +/// expands to them, and they are honored as well when passed through `hadoop.`. +const HA_NAMENODES_PREFIX: &str = "dfs.ha.namenodes"; +const HA_NAMENODE_RPC_ADDRESS_PREFIX: &str = "dfs.namenode.rpc-address"; + +/// Normalizes one NameNode spelling to `hdfs://host:port`, the form path +/// authorities take, so every source shares cache keys. `hdfs-native` dials +/// a socket address and has no default port, so anything else is `None`: +/// another scheme, a logical name, port 0, userinfo, a path. +fn hdfs_native_name_node(entry: &str) -> Option<String> { + let rest = entry.trim().trim_end_matches('/'); + let rest = rest.strip_prefix("hdfs://").unwrap_or(rest); + if rest.is_empty() || rest.contains("://") { + return None; + } + let url = Url::parse(&format!("hdfs://{rest}")).ok()?; + let (host, port) = (url.host_str()?, url.port()?); + let plain = host.is_empty() + || port == 0 + || !url.username().is_empty() + || url.password().is_some() + || !matches!(url.path(), "" | "/") + || url.query().is_some() + || url.fragment().is_some(); + (!plain).then(|| format!("hdfs://{host}:{port}")) +} + +/// Parses a comma-separated NameNode list, each entry normalized; an empty +/// list is `Ok` and empty. +fn hdfs_native_name_node_list(property: &str, value: &str) -> Result<Vec<String>> { + value + .split(',') + .map(str::trim) + .filter(|entry| !entry.is_empty()) + .map(|entry| { + hdfs_native_name_node(entry).ok_or_else(|| { + Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid `{property}` entry: {entry}, expected host:port (hdfs:// optional)" + ), + ) + }) + }) + .collect() +} + +/// Parse iceberg properties to [`HdfsNativeConfig`]. +pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) -> Result<HdfsNativeConfig> { + let mut cfg = HdfsNativeConfig::default(); + + // Entries are trimmed one by one: opendal splits the list on `,` as is, + // so a space after a comma would break failover to that NameNode. An + // empty result is dropped because `Operator::from_config` bypasses the + // builder's empty-string guard and `Some("")` would shadow the + // path-authority fallback below. + if let Some(name_node) = m.remove(HDFS_NAME_NODE) { + let entries = hdfs_native_name_node_list(HDFS_NAME_NODE, &name_node)?; + if !entries.is_empty() { + cfg.name_node = Some(entries.join(",")); + } + } + + // `hdfs.name-node.<nameservice>` is sugar for Hadoop's own declaration of + // an HA nameservice, which the resolver reads back from the options. + let nameservice_prefix = format!("{HDFS_NAME_NODE}."); + let declared_keys: Vec<String> = m + .keys() + .filter(|key| key.starts_with(&nameservice_prefix)) + .cloned() + .collect(); + let mut declared = Vec::new(); + for key in declared_keys { + let value = m.remove(&key).unwrap_or_default(); + let nameservice = key[nameservice_prefix.len()..].to_string(); + if nameservice.is_empty() || nameservice.chars().any(char::is_whitespace) { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid property `{key}`: a nameservice name must follow `{nameservice_prefix}`" + ), + )); + } + let entries = hdfs_native_name_node_list(&key, &value)?; + if entries.is_empty() { + return Err(Error::new( + ErrorKind::DataInvalid, + format!("Invalid property `{key}`: no NameNodes"), + )); + } + declared.push((nameservice, entries)); + } + + // A config carried over from PyIceberg would otherwise change identity + // silently; the client reads `HADOOP_USER_NAME` and the Kerberos cache. + for key in HDFS_UNSUPPORTED_KEYS { + if m.remove(key).is_some() { + tracing::warn!("`{key}` is not supported by the hdfs-native backend and is ignored"); + } + } + if m.contains_key(HDFS_HADOOP_CONF_PREFIX) { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid property `{HDFS_HADOOP_CONF_PREFIX}`: a Hadoop key must follow the prefix" + ), + )); + } + + let host = m + .remove(HDFS_HOST) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()); + let port = m + .remove(HDFS_PORT) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()) + .map(|port| { + port.parse::<u16>().map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid `{HDFS_PORT}`: {port}: {e}"), + ) + }) + }) + .transpose()?; + + let mut options: HashMap<String, String> = m + .into_iter() + .filter_map(|(key, value)| { + key.strip_prefix(HDFS_HADOOP_CONF_PREFIX) + .map(|stripped| (stripped.to_string(), value)) + }) + .collect(); + // Explicit `hadoop.` keys win over the sugar. + for (nameservice, entries) in declared { + let ids: Vec<String> = (0..entries.len()).map(|i| format!("nn{i}")).collect(); + options + .entry(format!("{HA_NAMENODES_PREFIX}.{nameservice}")) + .or_insert_with(|| ids.join(",")); + for (id, entry) in ids.iter().zip(&entries) { + options + .entry(format!( + "{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}" + )) + .or_insert_with(|| entry.trim_start_matches("hdfs://").to_string()); + } + } + // PyIceberg's `hdfs.host`/`hdfs.port` name the filesystem for + // authority-less paths, which is what Hadoop's `fs.defaultFS` means; an + // explicit `hadoop.fs.defaultFS` wins. + match host { + Some(host) => { + // An IPv6 literal needs brackets in a URI authority. + let host = if host.contains(':') && !host.starts_with('[') { + format!("[{host}]") + } else { + host + }; + let port = port.unwrap_or(HDFS_DEFAULT_PORT); + options + .entry(FS_DEFAULT_FS.to_string()) + .or_insert_with(|| format!("hdfs://{host}:{port}")); + } + None if port.is_some() => { + tracing::warn!("`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is ignored"); + } + None => {} + } + if !options.is_empty() { + cfg.options = Some(options); + } + + Ok(cfg) +} + +/// Parse an HDFS path into `Some("hdfs://<authority>")` (`None` when +/// authority-less) and the relative path (no leading `/`, opendal style). +pub(crate) fn hdfs_native_parse_path(path: &str) -> Result<(Option<String>, &str)> { + let url = Url::parse(path).map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid hdfs path: {path}: {e}"), + ) + })?; + // Non-special schemes parse even without `//` (e.g. `hdfs:x` is a valid + // non-hierarchical URL), so require the literal prefix before slicing. + let (Some(after_scheme), "hdfs") = (path.strip_prefix("hdfs://"), url.scheme()) else { + return Err(Error::new( + ErrorKind::DataInvalid, + format!("Invalid hdfs path: {path}, expected scheme `hdfs://`"), + )); + }; + + let name_node = url.host_str().filter(|h| !h.is_empty()).map(|host| { + url.port() + .map(|port| format!("hdfs://{host}:{port}")) + .unwrap_or_else(|| format!("hdfs://{host}")) + }); + + // `url.path()` borrows from `url` and can't be returned with the input's + // lifetime. Slice the path component out of the original input instead; + // it starts after the first `/` following the `hdfs://` prefix. Opendal + // paths must not start with `/` (`Deleter::delete` rejects them). + let rel = match after_scheme.find('/') { + Some(i) => after_scheme[i..].trim_start_matches('/'), + None => "", + }; + + Ok((name_node, rel)) +} + +/// Resolves the effective NameNode for a path, plus the relative path. As in +/// Hadoop, an authority with a port is used as is; a logical nameservice +/// authority (no port) resolves through `hdfs.name-node`, and an +/// authority-less path through `hdfs.name-node`, else `fs.defaultFS`. The +/// operator cache, `delete_stream` batching and `relativize_path` all go +/// through this, so they cannot drift apart. +pub(crate) fn hdfs_native_effective_name_node<'a>( + config: &HdfsNativeConfig, + path: &'a str, +) -> Result<(String, &'a str)> { + let (authority, relative_path) = hdfs_native_parse_path(path)?; + let invalid = |reason: String| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid hdfs path: {path}, {reason}"), + ) + }; + let name_node = match authority { + Some(authority) if hdfs_native_name_node(&authority).is_some() => authority, + Some(logical) => { + let nameservice = logical.trim_start_matches("hdfs://"); + match hdfs_native_nameservice(config, nameservice)? { + Some(name_node) => name_node, + None => config.name_node.clone().ok_or_else(|| { Review Comment: With the plain `hdfs.name-node` set, any portless authority that isn't a declared nameservice lands here and resolves to that one cluster — so `hdfs://ns-typo/...` or a bare second-cluster hostname gets read, written, and `delete_prefix`'d against the wrong NameNode, no error. On a single cluster the fallback is accidentally harmless; on anything multi-cluster it's a silent wrong-cluster delete. I'd require the declaration (`hdfs.name-node.<ns>`) for logical authorities and keep the plain key for authority-less paths only — or, if the single-cluster convenience matters, fall back only when the authority matches the `fs.defaultFS` nameservice. Either way an unknown name should error rather than guess. ########## crates/storage/opendal/src/hdfs_native.rs: ########## @@ -0,0 +1,1244 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! HDFS storage backend via OpenDAL's `services-hdfs-native` (pure Rust, no JNI). + +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::sync::{Arc, RwLock, Weak}; + +use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_HOST, HDFS_NAME_NODE, HDFS_PORT}; +use iceberg::{Error, ErrorKind, Result}; +use opendal::Operator; +use opendal::services::HdfsNativeConfig; +use serde::{Deserialize, Serialize}; +use tokio::runtime::Handle; +use tokio::task::JoinHandle; +use url::Url; + +use crate::OpenDalClientConfig; +use crate::utils::from_opendal_error; + +/// Hadoop's default filesystem, which serves authority-less paths. +const FS_DEFAULT_FS: &str = "fs.defaultFS"; +const HDFS_DEFAULT_PORT: u16 = 8020; +/// PyIceberg keys with no equivalent in opendal's config. +const HDFS_UNSUPPORTED_KEYS: [&str; 2] = ["hdfs.user", "hdfs.kerberos_ticket"]; +/// Hadoop's keys declaring an HA nameservice: `hdfs.name-node.<nameservice>` +/// expands to them, and they are honored as well when passed through `hadoop.`. +const HA_NAMENODES_PREFIX: &str = "dfs.ha.namenodes"; +const HA_NAMENODE_RPC_ADDRESS_PREFIX: &str = "dfs.namenode.rpc-address"; + +/// Normalizes one NameNode spelling to `hdfs://host:port`, the form path +/// authorities take, so every source shares cache keys. `hdfs-native` dials +/// a socket address and has no default port, so anything else is `None`: +/// another scheme, a logical name, port 0, userinfo, a path. +fn hdfs_native_name_node(entry: &str) -> Option<String> { + let rest = entry.trim().trim_end_matches('/'); + let rest = rest.strip_prefix("hdfs://").unwrap_or(rest); + if rest.is_empty() || rest.contains("://") { + return None; + } + let url = Url::parse(&format!("hdfs://{rest}")).ok()?; + let (host, port) = (url.host_str()?, url.port()?); + let plain = host.is_empty() + || port == 0 + || !url.username().is_empty() + || url.password().is_some() + || !matches!(url.path(), "" | "/") + || url.query().is_some() + || url.fragment().is_some(); + (!plain).then(|| format!("hdfs://{host}:{port}")) +} + +/// Parses a comma-separated NameNode list, each entry normalized; an empty +/// list is `Ok` and empty. +fn hdfs_native_name_node_list(property: &str, value: &str) -> Result<Vec<String>> { + value + .split(',') + .map(str::trim) + .filter(|entry| !entry.is_empty()) + .map(|entry| { + hdfs_native_name_node(entry).ok_or_else(|| { + Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid `{property}` entry: {entry}, expected host:port (hdfs:// optional)" + ), + ) + }) + }) + .collect() +} + +/// Parse iceberg properties to [`HdfsNativeConfig`]. +pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) -> Result<HdfsNativeConfig> { + let mut cfg = HdfsNativeConfig::default(); + + // Entries are trimmed one by one: opendal splits the list on `,` as is, + // so a space after a comma would break failover to that NameNode. An + // empty result is dropped because `Operator::from_config` bypasses the + // builder's empty-string guard and `Some("")` would shadow the + // path-authority fallback below. + if let Some(name_node) = m.remove(HDFS_NAME_NODE) { + let entries = hdfs_native_name_node_list(HDFS_NAME_NODE, &name_node)?; + if !entries.is_empty() { + cfg.name_node = Some(entries.join(",")); + } + } + + // `hdfs.name-node.<nameservice>` is sugar for Hadoop's own declaration of + // an HA nameservice, which the resolver reads back from the options. + let nameservice_prefix = format!("{HDFS_NAME_NODE}."); + let declared_keys: Vec<String> = m + .keys() + .filter(|key| key.starts_with(&nameservice_prefix)) + .cloned() + .collect(); + let mut declared = Vec::new(); + for key in declared_keys { + let value = m.remove(&key).unwrap_or_default(); + let nameservice = key[nameservice_prefix.len()..].to_string(); + if nameservice.is_empty() || nameservice.chars().any(char::is_whitespace) { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid property `{key}`: a nameservice name must follow `{nameservice_prefix}`" + ), + )); + } + let entries = hdfs_native_name_node_list(&key, &value)?; + if entries.is_empty() { + return Err(Error::new( + ErrorKind::DataInvalid, + format!("Invalid property `{key}`: no NameNodes"), + )); + } + declared.push((nameservice, entries)); + } + + // A config carried over from PyIceberg would otherwise change identity + // silently; the client reads `HADOOP_USER_NAME` and the Kerberos cache. + for key in HDFS_UNSUPPORTED_KEYS { + if m.remove(key).is_some() { + tracing::warn!("`{key}` is not supported by the hdfs-native backend and is ignored"); + } + } + if m.contains_key(HDFS_HADOOP_CONF_PREFIX) { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid property `{HDFS_HADOOP_CONF_PREFIX}`: a Hadoop key must follow the prefix" + ), + )); + } + + let host = m + .remove(HDFS_HOST) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()); + let port = m + .remove(HDFS_PORT) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()) + .map(|port| { + port.parse::<u16>().map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid `{HDFS_PORT}`: {port}: {e}"), + ) + }) + }) + .transpose()?; + + let mut options: HashMap<String, String> = m + .into_iter() + .filter_map(|(key, value)| { + key.strip_prefix(HDFS_HADOOP_CONF_PREFIX) + .map(|stripped| (stripped.to_string(), value)) + }) + .collect(); + // Explicit `hadoop.` keys win over the sugar. + for (nameservice, entries) in declared { + let ids: Vec<String> = (0..entries.len()).map(|i| format!("nn{i}")).collect(); + options + .entry(format!("{HA_NAMENODES_PREFIX}.{nameservice}")) + .or_insert_with(|| ids.join(",")); + for (id, entry) in ids.iter().zip(&entries) { + options + .entry(format!( + "{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}" + )) + .or_insert_with(|| entry.trim_start_matches("hdfs://").to_string()); + } + } + // PyIceberg's `hdfs.host`/`hdfs.port` name the filesystem for + // authority-less paths, which is what Hadoop's `fs.defaultFS` means; an + // explicit `hadoop.fs.defaultFS` wins. + match host { + Some(host) => { + // An IPv6 literal needs brackets in a URI authority. + let host = if host.contains(':') && !host.starts_with('[') { + format!("[{host}]") + } else { + host + }; + let port = port.unwrap_or(HDFS_DEFAULT_PORT); + options + .entry(FS_DEFAULT_FS.to_string()) + .or_insert_with(|| format!("hdfs://{host}:{port}")); Review Comment: small one, not blocking — `hdfs.host`/`hdfs.port` aren't validated before composing, so `hdfs.host=nn:8020` becomes `hdfs://[nn:8020]:8020` and `hdfs.port=0` becomes `hdfs://host:0`, both failing later with a confusing error. Running the composed value through `hdfs_native_name_node` here and erroring on `hdfs.host`/`hdfs.port` would catch it at config time. ########## crates/storage/opendal/src/hdfs_native.rs: ########## @@ -0,0 +1,1244 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! HDFS storage backend via OpenDAL's `services-hdfs-native` (pure Rust, no JNI). + +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::sync::{Arc, RwLock, Weak}; + +use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_HOST, HDFS_NAME_NODE, HDFS_PORT}; +use iceberg::{Error, ErrorKind, Result}; +use opendal::Operator; +use opendal::services::HdfsNativeConfig; +use serde::{Deserialize, Serialize}; +use tokio::runtime::Handle; +use tokio::task::JoinHandle; +use url::Url; + +use crate::OpenDalClientConfig; +use crate::utils::from_opendal_error; + +/// Hadoop's default filesystem, which serves authority-less paths. +const FS_DEFAULT_FS: &str = "fs.defaultFS"; +const HDFS_DEFAULT_PORT: u16 = 8020; +/// PyIceberg keys with no equivalent in opendal's config. +const HDFS_UNSUPPORTED_KEYS: [&str; 2] = ["hdfs.user", "hdfs.kerberos_ticket"]; +/// Hadoop's keys declaring an HA nameservice: `hdfs.name-node.<nameservice>` +/// expands to them, and they are honored as well when passed through `hadoop.`. +const HA_NAMENODES_PREFIX: &str = "dfs.ha.namenodes"; +const HA_NAMENODE_RPC_ADDRESS_PREFIX: &str = "dfs.namenode.rpc-address"; + +/// Normalizes one NameNode spelling to `hdfs://host:port`, the form path +/// authorities take, so every source shares cache keys. `hdfs-native` dials +/// a socket address and has no default port, so anything else is `None`: +/// another scheme, a logical name, port 0, userinfo, a path. +fn hdfs_native_name_node(entry: &str) -> Option<String> { + let rest = entry.trim().trim_end_matches('/'); + let rest = rest.strip_prefix("hdfs://").unwrap_or(rest); + if rest.is_empty() || rest.contains("://") { + return None; + } + let url = Url::parse(&format!("hdfs://{rest}")).ok()?; + let (host, port) = (url.host_str()?, url.port()?); + let plain = host.is_empty() + || port == 0 + || !url.username().is_empty() + || url.password().is_some() + || !matches!(url.path(), "" | "/") + || url.query().is_some() + || url.fragment().is_some(); + (!plain).then(|| format!("hdfs://{host}:{port}")) +} + +/// Parses a comma-separated NameNode list, each entry normalized; an empty +/// list is `Ok` and empty. +fn hdfs_native_name_node_list(property: &str, value: &str) -> Result<Vec<String>> { + value + .split(',') + .map(str::trim) + .filter(|entry| !entry.is_empty()) + .map(|entry| { + hdfs_native_name_node(entry).ok_or_else(|| { + Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid `{property}` entry: {entry}, expected host:port (hdfs:// optional)" + ), + ) + }) + }) + .collect() +} + +/// Parse iceberg properties to [`HdfsNativeConfig`]. +pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) -> Result<HdfsNativeConfig> { + let mut cfg = HdfsNativeConfig::default(); + + // Entries are trimmed one by one: opendal splits the list on `,` as is, + // so a space after a comma would break failover to that NameNode. An + // empty result is dropped because `Operator::from_config` bypasses the + // builder's empty-string guard and `Some("")` would shadow the + // path-authority fallback below. + if let Some(name_node) = m.remove(HDFS_NAME_NODE) { + let entries = hdfs_native_name_node_list(HDFS_NAME_NODE, &name_node)?; + if !entries.is_empty() { + cfg.name_node = Some(entries.join(",")); + } + } + + // `hdfs.name-node.<nameservice>` is sugar for Hadoop's own declaration of + // an HA nameservice, which the resolver reads back from the options. + let nameservice_prefix = format!("{HDFS_NAME_NODE}."); + let declared_keys: Vec<String> = m + .keys() + .filter(|key| key.starts_with(&nameservice_prefix)) + .cloned() + .collect(); + let mut declared = Vec::new(); + for key in declared_keys { + let value = m.remove(&key).unwrap_or_default(); + let nameservice = key[nameservice_prefix.len()..].to_string(); + if nameservice.is_empty() || nameservice.chars().any(char::is_whitespace) { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid property `{key}`: a nameservice name must follow `{nameservice_prefix}`" + ), + )); + } + let entries = hdfs_native_name_node_list(&key, &value)?; + if entries.is_empty() { + return Err(Error::new( + ErrorKind::DataInvalid, + format!("Invalid property `{key}`: no NameNodes"), + )); + } + declared.push((nameservice, entries)); + } + + // A config carried over from PyIceberg would otherwise change identity + // silently; the client reads `HADOOP_USER_NAME` and the Kerberos cache. + for key in HDFS_UNSUPPORTED_KEYS { + if m.remove(key).is_some() { + tracing::warn!("`{key}` is not supported by the hdfs-native backend and is ignored"); + } + } + if m.contains_key(HDFS_HADOOP_CONF_PREFIX) { + return Err(Error::new( + ErrorKind::DataInvalid, + format!( + "Invalid property `{HDFS_HADOOP_CONF_PREFIX}`: a Hadoop key must follow the prefix" + ), + )); + } + + let host = m + .remove(HDFS_HOST) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()); + let port = m + .remove(HDFS_PORT) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()) + .map(|port| { + port.parse::<u16>().map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid `{HDFS_PORT}`: {port}: {e}"), + ) + }) + }) + .transpose()?; + + let mut options: HashMap<String, String> = m + .into_iter() + .filter_map(|(key, value)| { + key.strip_prefix(HDFS_HADOOP_CONF_PREFIX) + .map(|stripped| (stripped.to_string(), value)) + }) + .collect(); + // Explicit `hadoop.` keys win over the sugar. + for (nameservice, entries) in declared { + let ids: Vec<String> = (0..entries.len()).map(|i| format!("nn{i}")).collect(); + options + .entry(format!("{HA_NAMENODES_PREFIX}.{nameservice}")) + .or_insert_with(|| ids.join(",")); + for (id, entry) in ids.iter().zip(&entries) { + options + .entry(format!( + "{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}" + )) + .or_insert_with(|| entry.trim_start_matches("hdfs://").to_string()); + } + } + // PyIceberg's `hdfs.host`/`hdfs.port` name the filesystem for + // authority-less paths, which is what Hadoop's `fs.defaultFS` means; an + // explicit `hadoop.fs.defaultFS` wins. + match host { + Some(host) => { + // An IPv6 literal needs brackets in a URI authority. + let host = if host.contains(':') && !host.starts_with('[') { + format!("[{host}]") + } else { + host + }; + let port = port.unwrap_or(HDFS_DEFAULT_PORT); + options + .entry(FS_DEFAULT_FS.to_string()) + .or_insert_with(|| format!("hdfs://{host}:{port}")); + } + None if port.is_some() => { + tracing::warn!("`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is ignored"); + } + None => {} + } + if !options.is_empty() { + cfg.options = Some(options); + } + + Ok(cfg) +} + +/// Parse an HDFS path into `Some("hdfs://<authority>")` (`None` when +/// authority-less) and the relative path (no leading `/`, opendal style). +pub(crate) fn hdfs_native_parse_path(path: &str) -> Result<(Option<String>, &str)> { + let url = Url::parse(path).map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid hdfs path: {path}: {e}"), + ) + })?; + // Non-special schemes parse even without `//` (e.g. `hdfs:x` is a valid + // non-hierarchical URL), so require the literal prefix before slicing. + let (Some(after_scheme), "hdfs") = (path.strip_prefix("hdfs://"), url.scheme()) else { + return Err(Error::new( + ErrorKind::DataInvalid, + format!("Invalid hdfs path: {path}, expected scheme `hdfs://`"), + )); + }; + + let name_node = url.host_str().filter(|h| !h.is_empty()).map(|host| { Review Comment: Two gaps here that the config side already closes: userinfo is silently dropped (`hdfs://alice@nn:8020/a` resolves as `nn:8020`, while `hdfs_native_name_node` rejects userinfo for configured values), and port 0 passes through as `hdfs://nn:0` — which then fails `hdfs_native_name_node`, gets re-read as a logical nameservice `nn:0`, and takes the fallback below. I'd reject both here with `DataInvalid`; an authority that carries a port but fails validation should be an error, not a nameservice. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
