mixermt commented on code in PR #3111:
URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4234494910
##########
crates/iceberg/src/io/storage/config/hdfs_native.rs:
##########
Review Comment:
Done: the module is now `hdfs_native`, matching its siblings (`azdls`,
`gcs`, `hf`, `oss`, `s3`). It is private and re-exported, so the public
`iceberg::io::HDFS_*` paths are unchanged.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -1159,6 +1201,102 @@ mod tests {
);
}
+ #[cfg(feature = "opendal-hdfs-native")]
+ fn hdfs_native_test_storage() -> OpenDalStorage {
+ OpenDalStorage::HdfsNative(HdfsNativeStorage::new(
+ opendal::services::HdfsNativeConfig::default(),
+ OpenDalClientConfig::default(),
+ ))
+ }
+
+ /// The configuration round-trips through serde; the operator cache does
+ /// not and starts empty.
+ #[cfg(feature = "opendal-hdfs-native")]
+ #[tokio::test]
+ async fn test_hdfs_native_storage_serde_round_trip() {
Review Comment:
Done: both tests moved to `hdfs_native` with the module's
`test_hdfs_native_` prefix. Opened #3385 for splitting the remaining `lib.rs`
tests.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1347 @@
+// 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, PoisonError, 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: {}, expected host:port
(hdfs:// optional)",
+ hdfs_native_redact(entry)
+ ),
+ )
+ })
+ })
+ .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
Review Comment:
Done: the parser now reads `hdfs.name-node`, the
`hdfs.name-node.<nameservice>` and `hadoop.` prefixes, `hdfs.host` and
`hdfs.port` through `#[derive(Properties)]`. The validation on top is unchanged.
##########
.github/workflows/ci.yml:
##########
@@ -265,6 +265,46 @@ jobs:
if: always() && matrix.test-suite.name == 'default'
run: make docker-down
+ tests-hdfs:
Review Comment:
Done: the HDFS fixture now follows the other fixtures. It runs on the shared
bridge network with published ports, so `make docker-up` starts it on every
platform (the DataNode advertises `127.0.0.1` and the tests set
`dfs.client.use.datanode.hostname`), the tests always run, and the dedicated
job is gone. Trade-off: the default job pays the Hadoop startup again (~40 s),
which @laskoviymishka had asked to move out of it.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1347 @@
+// 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, PoisonError, 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: {}, expected host:port
(hdfs:// optional)",
+ hdfs_native_redact(entry)
+ ),
+ )
+ })
+ })
+ .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
+ mut m: HashMap<String, String>,
+ mut warn: impl FnMut(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() {
+ warn(format!(
+ "`{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 bracketed = if host.contains(':') && !host.starts_with('[') {
+ format!("[{host}]")
+ } else {
+ host.clone()
+ };
+ let port = port.unwrap_or(HDFS_DEFAULT_PORT);
+ // Validated like every other NameNode spelling, so a host that
+ // carries a scheme or port fails here and not at the first I/O.
+ let default_fs =
hdfs_native_name_node(&format!("hdfs://{bracketed}:{port}"))
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid `{HDFS_HOST}`/`{HDFS_PORT}`: {}, expected
a host name or IP and a port",
+ hdfs_native_redact(&format!("{host}:{port}"))
+ ),
+ )
+ })?;
+ options
+ .entry(FS_DEFAULT_FS.to_string())
+ .or_insert(default_fs);
+ }
+ None if port.is_some() => {
+ warn(format!(
+ "`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is
ignored"
+ ));
+ }
+ None => {}
+ }
+ if !options.is_empty() {
+ cfg.options = Some(options);
+ }
+
+ Ok(cfg)
+}
+
+/// `s` with everything between its scheme and its last `@` masked, so an
+/// error message never echoes userinfo, even a password holding a raw `/`,
+/// `?`, `#` or `://` that a URL parser would split elsewhere. An `@` further
+/// along the path is masked the same way, which only costs context.
+fn hdfs_native_redact(s: &str) -> String {
+ let start = s
+ .find("://")
+ .filter(|&i| {
+ s[..i]
+ .chars()
+ .all(|c| c.is_ascii_alphanumeric() || "+-.".contains(c))
+ })
+ .map_or(0, |i| i + 3);
+ match s[start..].rfind('@') {
+ Some(at) => format!("{}***{}", &s[..start], &s[start + at..]),
+ None => s.to_string(),
+ }
+}
+
+/// A `DataInvalid` error for an HDFS path, its userinfo masked.
+fn hdfs_native_invalid_path(path: &str, reason: impl std::fmt::Display) ->
Error {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid hdfs path: {}, {reason}", hdfs_native_redact(path)),
+ )
+}
+
+/// 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| hdfs_native_invalid_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(hdfs_native_invalid_path(path, "expected scheme
`hdfs://`"));
+ };
+ // Userinfo has nowhere to go and port 0 cannot be dialed; silently
+ // dropping the one or reading the other as a logical name would mislead.
+ if !url.username().is_empty() || url.password().is_some() {
+ return Err(hdfs_native_invalid_path(path, "userinfo is not
supported"));
+ }
+ if url.port() == Some(0) {
+ return Err(hdfs_native_invalid_path(path, "port 0 cannot be dialed"));
+ }
+
+ 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. An
+/// authority with a port is used as is, a portless one must be a declared
+/// nameservice, and an authority-less path uses `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| hdfs_native_invalid_path(path, reason);
+ let name_node = match authority {
+ Some(authority) if hdfs_native_name_node(&authority).is_some() =>
authority,
+ Some(portless) => {
+ let nameservice = portless.trim_start_matches("hdfs://");
+ hdfs_native_nameservice(config, nameservice)?.ok_or_else(|| {
+ invalid(format!(
+ "`{nameservice}` has no port and is not a declared
nameservice; add the port or set `{HDFS_NAME_NODE}.{nameservice}`"
+ ))
+ })?
+ }
+ None => match (&config.name_node, hdfs_native_default_fs(config)) {
+ (Some(name_node), _) => name_node.clone(),
+ (None, Some(default_fs)) =>
hdfs_native_name_node(default_fs).ok_or_else(|| {
+ invalid(format!(
+ "`{FS_DEFAULT_FS}` {} is not an HDFS host:port, a logical
nameservice requires `{HDFS_NAME_NODE}`",
+ hdfs_native_redact(default_fs)
+ ))
+ })?,
+ (None, None) => {
+ return Err(invalid(format!(
+ "authority-less paths require `{HDFS_NAME_NODE}` or
`{HDFS_HOST}`"
+ )));
+ }
+ },
+ };
+ Ok((name_node, relative_path))
+}
+
+/// `fs.defaultFS` from the forwarded options, as written.
+fn hdfs_native_default_fs(config: &HdfsNativeConfig) -> Option<&str> {
+ config
+ .options
+ .as_ref()?
+ .get(FS_DEFAULT_FS)
+ .map(|s| s.trim())
+ .filter(|s| !s.is_empty())
+}
+
+/// The NameNodes that Hadoop's keys in the forwarded options declare for a
+/// nameservice, if any. A declaration with a missing or malformed address is
+/// an error rather than a silent fallback.
+fn hdfs_native_nameservice(config: &HdfsNativeConfig, nameservice: &str) ->
Result<Option<String>> {
+ let Some(options) = config.options.as_ref() else {
+ return Ok(None);
+ };
+ let Some(ids) =
options.get(&format!("{HA_NAMENODES_PREFIX}.{nameservice}")) else {
+ return Ok(None);
+ };
+ let name_nodes = ids
+ .split(',')
+ .map(str::trim)
+ .filter(|id| !id.is_empty())
+ .map(|id| {
+ let key =
format!("{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}");
+ options
+ .get(&key)
+ .and_then(|value| hdfs_native_name_node(value))
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Nameservice `{nameservice}` declares NameNode
`{id}` but `{key}` is missing or not host:port"
+ ),
+ )
+ })
+ })
+ .collect::<Result<Vec<_>>>()?;
+ if name_nodes.is_empty() {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!("Nameservice `{nameservice}` declares no NameNodes"),
+ ));
+ }
+ Ok(Some(name_nodes.join(",")))
+}
+
+/// State of [`OpenDalStorage::HdfsNative`](crate::OpenDalStorage::HdfsNative):
+/// the parsed configuration and the per-NameNode operator cache. Only the
+/// storage factories build it.
+#[derive(Clone, Debug, Serialize, Deserialize)]
+pub struct HdfsNativeStorage {
+ pub(crate) config: Arc<HdfsNativeConfig>,
+ #[serde(skip, default)]
+ pub(crate) operators: HdfsNativeOperatorCache,
+ #[serde(default)]
+ pub(crate) client_config: OpenDalClientConfig,
+}
+
+impl HdfsNativeStorage {
+ pub(crate) fn new(config: HdfsNativeConfig, client_config:
OpenDalClientConfig) -> Self {
+ Self {
+ config: Arc::new(config),
+ operators: HdfsNativeOperatorCache::default(),
+ client_config,
+ }
+ }
+}
+
+/// Operators cached per effective NameNode: each holds an `hdfs-native`
+/// client with live RPC connections, whose tasks run on the tokio runtime
+/// current when it was built. An entry is rebuilt once that runtime is
+/// gone, as `hdfs-native` panics when
+/// it spawns onto a dead one. The cache lives as long as the storage that
+/// owns it (clones share it).
+#[derive(Clone, Debug, Default)]
+pub(crate) struct HdfsNativeOperatorCache(Arc<RwLock<HashMap<String,
CachedOperator>>>);
+
+#[derive(Debug)]
+struct CachedOperator {
+ operator: Operator,
+ sentinel: RuntimeSentinel,
+}
+
+/// A task parked on the building runtime that owns the token, so the token
+/// outlives it only while that runtime is alive. Aborted on drop so entries
+/// do not leave parked tasks behind.
+#[derive(Debug)]
+struct RuntimeSentinel {
+ alive: Weak<()>,
+ task: JoinHandle<()>,
+}
+
+impl RuntimeSentinel {
+ fn spawn(handle: &Handle) -> Self {
+ let token = Arc::new(());
+ let alive = Arc::downgrade(&token);
+ let task = handle.spawn(async move {
+ let _token = token;
+ std::future::pending::<()>().await
+ });
+ Self { alive, task }
+ }
+}
+
+impl Drop for RuntimeSentinel {
+ fn drop(&mut self) {
+ self.task.abort();
+ }
+}
+
+impl CachedOperator {
+ fn new(operator: Operator, handle: &Handle) -> Self {
+ Self {
+ operator,
+ sentinel: RuntimeSentinel::spawn(handle),
+ }
+ }
+
+ fn runtime_alive(&self) -> bool {
+ self.sentinel.alive.strong_count() > 0
+ }
+}
+
+// A poisoned lock is recovered rather than reported: writes only ever
+// replace whole entries, so a panic elsewhere cannot leave the map broken.
+impl HdfsNativeOperatorCache {
+ pub(crate) fn get(&self, name_node: &str) -> Option<Operator> {
+ self.0
+ .read()
+ .unwrap_or_else(PoisonError::into_inner)
+ .get(name_node)
+ .filter(|cached| cached.runtime_alive())
+ .map(|cached| cached.operator.clone())
+ }
+
+ /// Inserts `op` unless a concurrent caller got there first, returning
+ /// whichever operator the cache now holds; a stale entry is replaced.
+ fn insert(&self, name_node: String, op: Operator, handle: &Handle) ->
Operator {
+ let mut operators =
self.0.write().unwrap_or_else(PoisonError::into_inner);
+ match operators.entry(name_node) {
+ Entry::Occupied(entry) if entry.get().runtime_alive() =>
entry.get().operator.clone(),
+ Entry::Occupied(mut entry) => {
+ entry.insert(CachedOperator::new(op.clone(), handle));
+ op
+ }
+ Entry::Vacant(entry) => {
+ entry.insert(CachedOperator::new(op.clone(), handle));
+ op
+ }
+ }
+ }
+
+ #[cfg(test)]
+ fn len(&self) -> usize {
+ self.0.read().unwrap().len()
+ }
+}
+
+/// Creates an operator for the path, reusing the cached one for its
+/// effective NameNode.
+pub(crate) async fn hdfs_native_create_operator<'a>(
+ path: &'a str,
+ config: &Arc<HdfsNativeConfig>,
+ operators: &HdfsNativeOperatorCache,
+) -> Result<(Operator, &'a str)> {
+ let (name_node, relative_path) = hdfs_native_effective_name_node(config,
path)?;
+
+ if let Some(op) = operators.get(&name_node) {
+ return Ok((op, relative_path));
+ }
+
+ // Every operator in this crate needs a tokio runtime for its I/O (the
+ // timeout layer), so say so instead of panicking in `spawn_blocking`.
+ let handle = Handle::try_current().map_err(|_| {
+ Error::new(
+ ErrorKind::FeatureUnsupported,
+ "HDFS storage requires a tokio runtime",
+ )
+ })?;
+
+ // The build reads the Hadoop XML config synchronously, so it runs on a
+ // blocking thread and outside the lock. A racing first caller may build
+ // too; the loser is dropped before opening any connection.
+ let build_config = Arc::clone(config);
+ let build_name_node = name_node.clone();
+ let op = handle
Review Comment:
Done, and the runtime handle is gone too. An hdfs-native client spawns its
NameNode connections on the runtime current when it is built and panics once
that runtime is dropped, so clients are now built inside a crate-owned runtime
and work from any caller's runtime: no `Handle::try_current`, sentinel or
`spawn_blocking` remain, and `create_operator` is sync again. The Hadoop XML
read (~0.1 ms) runs inline, as @laskoviymishka allowed for, with a note at the
call site. `FileIO` has no runtime to pass today; once it carries
`iceberg::runtime::Runtime`, the build can use its io handle instead.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -372,7 +394,7 @@ impl OpenDalStorage {
/// * An [`opendal::Operator`] instance used to operate on file.
/// * Relative path to the root uri of [`opendal::Operator`].
#[allow(unreachable_code, unused_variables)]
- pub(crate) fn create_operator<'a>(
+ pub(crate) async fn create_operator<'a>(
Review Comment:
Done: `create_operator` is sync again for every backend, as on main; the
HDFS build no longer needs a runtime (see the reply to blackmwk above).
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1347 @@
+// 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, PoisonError, 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: {}, expected host:port
(hdfs:// optional)",
+ hdfs_native_redact(entry)
+ ),
+ )
+ })
+ })
+ .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
+ mut m: HashMap<String, String>,
+ mut warn: impl FnMut(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() {
+ warn(format!(
+ "`{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 bracketed = if host.contains(':') && !host.starts_with('[') {
+ format!("[{host}]")
+ } else {
+ host.clone()
+ };
+ let port = port.unwrap_or(HDFS_DEFAULT_PORT);
+ // Validated like every other NameNode spelling, so a host that
+ // carries a scheme or port fails here and not at the first I/O.
+ let default_fs =
hdfs_native_name_node(&format!("hdfs://{bracketed}:{port}"))
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid `{HDFS_HOST}`/`{HDFS_PORT}`: {}, expected
a host name or IP and a port",
+ hdfs_native_redact(&format!("{host}:{port}"))
+ ),
+ )
+ })?;
+ options
+ .entry(FS_DEFAULT_FS.to_string())
+ .or_insert(default_fs);
+ }
+ None if port.is_some() => {
+ warn(format!(
+ "`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is
ignored"
+ ));
+ }
+ None => {}
+ }
+ if !options.is_empty() {
+ cfg.options = Some(options);
+ }
+
+ Ok(cfg)
+}
+
+/// `s` with everything between its scheme and its last `@` masked, so an
+/// error message never echoes userinfo, even a password holding a raw `/`,
+/// `?`, `#` or `://` that a URL parser would split elsewhere. An `@` further
+/// along the path is masked the same way, which only costs context.
+fn hdfs_native_redact(s: &str) -> String {
+ let start = s
+ .find("://")
+ .filter(|&i| {
+ s[..i]
+ .chars()
+ .all(|c| c.is_ascii_alphanumeric() || "+-.".contains(c))
+ })
+ .map_or(0, |i| i + 3);
+ match s[start..].rfind('@') {
+ Some(at) => format!("{}***{}", &s[..start], &s[start + at..]),
+ None => s.to_string(),
+ }
+}
+
+/// A `DataInvalid` error for an HDFS path, its userinfo masked.
+fn hdfs_native_invalid_path(path: &str, reason: impl std::fmt::Display) ->
Error {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid hdfs path: {}, {reason}", hdfs_native_redact(path)),
+ )
+}
+
+/// 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| hdfs_native_invalid_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(hdfs_native_invalid_path(path, "expected scheme
`hdfs://`"));
+ };
+ // Userinfo has nowhere to go and port 0 cannot be dialed; silently
+ // dropping the one or reading the other as a logical name would mislead.
+ if !url.username().is_empty() || url.password().is_some() {
+ return Err(hdfs_native_invalid_path(path, "userinfo is not
supported"));
+ }
+ if url.port() == Some(0) {
+ return Err(hdfs_native_invalid_path(path, "port 0 cannot be dialed"));
+ }
+
+ 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. An
+/// authority with a port is used as is, a portless one must be a declared
+/// nameservice, and an authority-less path uses `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>(
Review Comment:
Done in part. Paths now parse once into an authority enum (`NameNode`,
`Nameservice`, `Default`), so the resolver no longer re-parses the authority,
and every declaration, including passed-through `hadoop.dfs.ha.namenodes.*` and
the `fs.defaultFS` that serves authority-less paths, is validated when the
storage is built. I kept the declarations in opendal's options rather than a
second, derived copy in the serialized storage, since a lookup is a few map
reads next to an RPC, and kept one file per backend like the other backends.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1347 @@
+// 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, PoisonError, 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: {}, expected host:port
(hdfs:// optional)",
+ hdfs_native_redact(entry)
+ ),
+ )
+ })
+ })
+ .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
+ mut m: HashMap<String, String>,
+ mut warn: impl FnMut(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() {
+ warn(format!(
+ "`{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 bracketed = if host.contains(':') && !host.starts_with('[') {
+ format!("[{host}]")
+ } else {
+ host.clone()
+ };
+ let port = port.unwrap_or(HDFS_DEFAULT_PORT);
+ // Validated like every other NameNode spelling, so a host that
+ // carries a scheme or port fails here and not at the first I/O.
+ let default_fs =
hdfs_native_name_node(&format!("hdfs://{bracketed}:{port}"))
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid `{HDFS_HOST}`/`{HDFS_PORT}`: {}, expected
a host name or IP and a port",
+ hdfs_native_redact(&format!("{host}:{port}"))
+ ),
+ )
+ })?;
+ options
+ .entry(FS_DEFAULT_FS.to_string())
+ .or_insert(default_fs);
+ }
+ None if port.is_some() => {
+ warn(format!(
+ "`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is
ignored"
+ ));
+ }
+ None => {}
+ }
+ if !options.is_empty() {
+ cfg.options = Some(options);
+ }
+
+ Ok(cfg)
+}
+
+/// `s` with everything between its scheme and its last `@` masked, so an
+/// error message never echoes userinfo, even a password holding a raw `/`,
+/// `?`, `#` or `://` that a URL parser would split elsewhere. An `@` further
+/// along the path is masked the same way, which only costs context.
+fn hdfs_native_redact(s: &str) -> String {
+ let start = s
+ .find("://")
+ .filter(|&i| {
+ s[..i]
+ .chars()
+ .all(|c| c.is_ascii_alphanumeric() || "+-.".contains(c))
+ })
+ .map_or(0, |i| i + 3);
+ match s[start..].rfind('@') {
+ Some(at) => format!("{}***{}", &s[..start], &s[start + at..]),
+ None => s.to_string(),
+ }
+}
+
+/// A `DataInvalid` error for an HDFS path, its userinfo masked.
+fn hdfs_native_invalid_path(path: &str, reason: impl std::fmt::Display) ->
Error {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid hdfs path: {}, {reason}", hdfs_native_redact(path)),
+ )
+}
+
+/// 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| hdfs_native_invalid_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(hdfs_native_invalid_path(path, "expected scheme
`hdfs://`"));
+ };
+ // Userinfo has nowhere to go and port 0 cannot be dialed; silently
+ // dropping the one or reading the other as a logical name would mislead.
+ if !url.username().is_empty() || url.password().is_some() {
+ return Err(hdfs_native_invalid_path(path, "userinfo is not
supported"));
+ }
+ if url.port() == Some(0) {
+ return Err(hdfs_native_invalid_path(path, "port 0 cannot be dialed"));
+ }
+
+ 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. An
+/// authority with a port is used as is, a portless one must be a declared
+/// nameservice, and an authority-less path uses `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>(
Review Comment:
The authority is an enum now (see above). I kept the `hdfs_native_` free
functions rather than methods, matching the other backends and @blackmwk's
earlier request.
--
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]