laskoviymishka commented on code in PR #3111:
URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4169088910
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -460,6 +499,11 @@ impl OpenDalStorage {
))
}
}
+ #[cfg(feature = "opendal-hdfs-native")]
+ OpenDalStorage::HdfsNative { .. } => {
+ let (_, relative_path) = hdfs_native_parse_path(path)?;
Review Comment:
`relativize_path` only calls `hdfs_native_parse_path`, so `hdfs:///a/b`
returns `Ok("a/b")` here while `create_operator` on the same unconfigured path
errors — the two disagree on whether an authority-less path is usable, and
every other backend fails fast in both. I'd route this through
`hdfs_native_effective_name_node` so a path that can't resolve a NameNode
errors the same way on both.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,411 @@
+// 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::sync::{Arc, RwLock};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_NAME_NODE};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+use crate::utils::from_opendal_error;
+
+/// Parse iceberg properties to [`HdfsNativeConfig`].
+pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ let mut cfg = HdfsNativeConfig::default();
+
+ // `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)
+ .map(|s| s.trim().trim_end_matches('/').to_string())
+ .filter(|s| !s.is_empty())
+ {
+ cfg.name_node = Some(name_node);
+ }
+
+ let 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();
+ 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:
`url.host_str()` strips the brackets off an IPv6 literal, so
`hdfs://[::1]:8020` reconstructs here as `hdfs://::1:8020` — not a parseable
URL, and both `url` and `hdfs-native` reject it. Any IPv6 NameNode then fails
every operator build and cache lookup.
`url.host()` keeps the bracketed form via its `Display`:
```rust
let name_node = url.host().filter(|h| !matches!(h,
url::Host::Domain(""))).map(|host| {
url.port().map(|port| format!("hdfs://{host}:{port}")).unwrap_or_else(||
format!("hdfs://{host}"))
});
```
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -243,6 +261,18 @@ pub enum OpenDalStorage {
/// GCS configuration.
config: Arc<GcsConfig>,
},
+ /// HDFS storage variant.
+ ///
+ /// The NameNode is taken from the `hdfs.name-node` property when set
+ /// (comma-separated endpoints enable HA failover), else the path
authority.
+ #[cfg(feature = "opendal-hdfs-native")]
+ HdfsNative {
+ /// HDFS configuration.
+ config: Arc<HdfsNativeConfig>,
+ /// Operator cache keyed by effective NameNode.
+ #[serde(skip, default)]
+ operators: HdfsNativeOperatorCache,
Review Comment:
@comphead already flagged this — seconding it. `operators` as a `pub` field
drags `HdfsNativeOperatorCache` into the crate's public API (it's in
`public-api.txt` now), which no other variant does — they expose only `config`.
That locks a connection-pooling detail into semver and lets a caller build the
variant with an isolated cache, silently losing reuse. I'd wrap the state in a
newtype, `HdfsNative(Arc<HdfsNativeState>)` with private internals (mirrors
`Memory(Operator)`), and drop the `pub use`. While we're here, a line
documenting the cache lifecycle — one entry per effective NameNode, lives with
the storage instance, no eviction — would help too.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,411 @@
+// 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::sync::{Arc, RwLock};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_NAME_NODE};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+use crate::utils::from_opendal_error;
+
+/// Parse iceberg properties to [`HdfsNativeConfig`].
+pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ let mut cfg = HdfsNativeConfig::default();
+
+ // `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)
+ .map(|s| s.trim().trim_end_matches('/').to_string())
+ .filter(|s| !s.is_empty())
+ {
+ cfg.name_node = Some(name_node);
+ }
+
+ let 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();
+ 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 — the configured
+/// `hdfs.name-node` when set, else the path authority — plus the relative
+/// path. Both the operator cache and `delete_stream` batching key on 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_name_node, relative_path) = hdfs_native_parse_path(path)?;
+ let name_node = config
+ .name_node
+ .clone()
+ .or(authority_name_node)
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid hdfs path: {path}, authority-less paths require
the `{HDFS_NAME_NODE}` property"
+ ),
+ )
+ })?;
+ Ok((name_node, relative_path))
+}
+
+/// 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.
+#[derive(Clone, Debug, Default)]
+pub struct HdfsNativeOperatorCache(Arc<RwLock<HashMap<String, Operator>>>);
+
+impl HdfsNativeOperatorCache {
+ fn get(&self, name_node: &str) -> Result<Option<Operator>> {
+ Ok(self.0.read().map_err(poisoned)?.get(name_node).cloned())
+ }
+
+ /// Inserts `op` unless a concurrent caller got there first, returning
+ /// whichever operator the cache now holds.
+ fn insert(&self, name_node: String, op: Operator) -> Result<Operator> {
+ Ok(self
+ .0
+ .write()
+ .map_err(poisoned)?
+ .entry(name_node)
+ .or_insert(op)
+ .clone())
+ }
+
+ #[cfg(test)]
+ fn len(&self) -> usize {
+ self.0.read().unwrap().len()
+ }
+}
+
+fn poisoned<T>(_: T) -> Error {
+ Error::new(ErrorKind::Unexpected, "HDFS operator cache lock poisoned")
+}
+
+/// Creates an operator for the path, reusing the cached one for its
+/// effective NameNode.
+pub(crate) fn hdfs_native_create_operator<'a>(
+ path: &'a str,
+ config: &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));
+ }
+
+ // Built outside the lock: the build reads the Hadoop XML config
+ // synchronously. A racing first caller may build too; the loser is
+ // dropped before opening any connection.
+ let op = hdfs_native_operator_build(config, &name_node)?;
+ Ok((operators.insert(name_node, op)?, relative_path))
+}
+
+/// Returns the `delete_stream` grouping key for a path: the effective
+/// NameNode, so paths that resolve to different operators never share a
+/// deleter. Unresolvable paths key on themselves (as `hf_batch_key` does);
+/// `create_operator` then reports the real error.
+pub(crate) fn hdfs_native_batch_key(config: &HdfsNativeConfig, path: &str) ->
String {
+ hdfs_native_effective_name_node(config, path)
+ .map(|(name_node, _)| name_node)
+ .unwrap_or_else(|_| path.to_string())
+}
+
+/// Build a new OpenDAL [`Operator`]: OpenDAL splits `name_node` on commas
+/// into a synthetic HA name service; `$HADOOP_CONF_DIR` XML still merges in.
+fn hdfs_native_operator_build(config: &HdfsNativeConfig, name_node: &str) ->
Result<Operator> {
+ let mut cfg = config.clone();
+ cfg.name_node = Some(name_node.to_string());
+ Operator::from_config(cfg).map_err(from_opendal_error)
Review Comment:
`Operator::from_config` reads the Hadoop XML synchronously, and this runs on
the tokio executor thread (sync fn reached from the async `Storage` methods)
with no `spawn_blocking`. @comphead's "blocking under the write lock" point is
handled — the build is outside the lock now — but the cold build still blocks
the executor thread, so first-access-per-NameNode under load stalls other
tasks. I'd make `hdfs_native_create_operator` async and wrap the build in
`spawn_blocking`, or document it as a known experimental constraint with a
tracking issue.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -135,6 +144,9 @@ pub enum OpenDalStorageFactory {
/// GCS storage factory.
#[cfg(feature = "opendal-gcs")]
Gcs,
+ /// HDFS storage factory.
+ #[cfg(feature = "opendal-hdfs-native")]
+ HdfsNative,
Review Comment:
@blackmwk already raised this and I agree — every other variant is named for
the storage system (`Gcs`, `Oss`, `Azdls`, `Hf`), so `HdfsNative` is the odd
one out leaking the opendal implementation choice. I'd name the public variant
`Hdfs` and keep `services-hdfs-native` as the internal/feature detail.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,411 @@
+// 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::sync::{Arc, RwLock};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_NAME_NODE};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+use crate::utils::from_opendal_error;
+
+/// Parse iceberg properties to [`HdfsNativeConfig`].
+pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ let mut cfg = HdfsNativeConfig::default();
+
+ // `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)
+ .map(|s| s.trim().trim_end_matches('/').to_string())
+ .filter(|s| !s.is_empty())
+ {
+ cfg.name_node = Some(name_node);
Review Comment:
A bare `hdfs.name-node = "namenode:8020"` (no `hdfs://`) is accepted and
passed straight through, and `hdfs-native` then fails with an opaque error
rather than something that points at the misconfigured property. I'd validate
the value starts with `hdfs://` here and return `DataInvalid` otherwise.
##########
crates/iceberg/src/io/storage/config/hdfs.rs:
##########
@@ -0,0 +1,27 @@
+// 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 configuration.
+
+/// HDFS NameNode RPC endpoint(s), e.g. `hdfs://namenode:8020`; a
+/// comma-separated list enables HA failover. When unset, the NameNode is
+/// derived from the path authority.
+pub const HDFS_NAME_NODE: &str = "hdfs.name-node";
Review Comment:
PyIceberg already has a de-facto convention here — `hdfs.host` + `hdfs.port`
(plus `hdfs.user` / `hdfs.kerberos_ticket`) — and this introduces
`hdfs.name-node` as a single combined key, so catalog properties carried over
from PyIceberg get silently ignored. It's not a metadata-interop issue (these
keys aren't stored in table metadata), but it's a real usability gap. At
minimum I'd note in the doc-comment that this key is
iceberg-rust/hdfs-native-specific and differs from PyIceberg; accepting
`hdfs.host`+`hdfs.port` as a fallback would be better still.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,411 @@
+// 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::sync::{Arc, RwLock};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_NAME_NODE};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+use crate::utils::from_opendal_error;
+
+/// Parse iceberg properties to [`HdfsNativeConfig`].
+pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ let mut cfg = HdfsNativeConfig::default();
+
+ // `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)
+ .map(|s| s.trim().trim_end_matches('/').to_string())
+ .filter(|s| !s.is_empty())
+ {
+ cfg.name_node = Some(name_node);
+ }
+
+ let 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();
+ 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 — the configured
+/// `hdfs.name-node` when set, else the path authority — plus the relative
+/// path. Both the operator cache and `delete_stream` batching key on 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_name_node, relative_path) = hdfs_native_parse_path(path)?;
+ let name_node = config
+ .name_node
+ .clone()
+ .or(authority_name_node)
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid hdfs path: {path}, authority-less paths require
the `{HDFS_NAME_NODE}` property"
+ ),
+ )
+ })?;
+ Ok((name_node, relative_path))
+}
+
+/// 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.
+#[derive(Clone, Debug, Default)]
+pub struct HdfsNativeOperatorCache(Arc<RwLock<HashMap<String, Operator>>>);
+
+impl HdfsNativeOperatorCache {
+ fn get(&self, name_node: &str) -> Result<Option<Operator>> {
+ Ok(self.0.read().map_err(poisoned)?.get(name_node).cloned())
+ }
+
+ /// Inserts `op` unless a concurrent caller got there first, returning
+ /// whichever operator the cache now holds.
+ fn insert(&self, name_node: String, op: Operator) -> Result<Operator> {
+ Ok(self
+ .0
+ .write()
+ .map_err(poisoned)?
+ .entry(name_node)
+ .or_insert(op)
+ .clone())
+ }
+
+ #[cfg(test)]
+ fn len(&self) -> usize {
+ self.0.read().unwrap().len()
+ }
+}
+
+fn poisoned<T>(_: T) -> Error {
+ Error::new(ErrorKind::Unexpected, "HDFS operator cache lock poisoned")
+}
+
+/// Creates an operator for the path, reusing the cached one for its
+/// effective NameNode.
+pub(crate) fn hdfs_native_create_operator<'a>(
+ path: &'a str,
+ config: &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));
+ }
+
+ // Built outside the lock: the build reads the Hadoop XML config
+ // synchronously. A racing first caller may build too; the loser is
+ // dropped before opening any connection.
+ let op = hdfs_native_operator_build(config, &name_node)?;
+ Ok((operators.insert(name_node, op)?, relative_path))
+}
+
+/// Returns the `delete_stream` grouping key for a path: the effective
+/// NameNode, so paths that resolve to different operators never share a
+/// deleter. Unresolvable paths key on themselves (as `hf_batch_key` does);
+/// `create_operator` then reports the real error.
+pub(crate) fn hdfs_native_batch_key(config: &HdfsNativeConfig, path: &str) ->
String {
+ hdfs_native_effective_name_node(config, path)
+ .map(|(name_node, _)| name_node)
+ .unwrap_or_else(|_| path.to_string())
+}
+
+/// Build a new OpenDAL [`Operator`]: OpenDAL splits `name_node` on commas
+/// into a synthetic HA name service; `$HADOOP_CONF_DIR` XML still merges in.
+fn hdfs_native_operator_build(config: &HdfsNativeConfig, name_node: &str) ->
Result<Operator> {
+ let mut cfg = config.clone();
+ cfg.name_node = Some(name_node.to_string());
+ Operator::from_config(cfg).map_err(from_opendal_error)
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn test_hdfs_native_config_parse_name_node_and_options() {
+ let props = HashMap::from([
+ (
+ HDFS_NAME_NODE.to_string(),
+ "hdfs://nn1:8020,hdfs://nn2:8020".to_string(),
+ ),
+ (
+ "hadoop.dfs.client.failover.random.order".to_string(),
+ "true".to_string(),
+ ),
+ ("unrelated.key".to_string(), "ignored".to_string()),
+ ]);
+
+ let cfg = hdfs_native_config_parse(props).unwrap();
+
+ assert_eq!(
+ cfg.name_node.as_deref(),
+ Some("hdfs://nn1:8020,hdfs://nn2:8020")
+ );
+ let options = cfg.options.unwrap();
+ assert_eq!(
+ options.get("dfs.client.failover.random.order"),
+ Some(&"true".to_string())
+ );
+ assert!(!options.contains_key("unrelated.key"));
+ }
+
+ #[test]
+ fn test_hdfs_native_config_parse_empty() {
+ let cfg = hdfs_native_config_parse(HashMap::new()).unwrap();
+
+ assert_eq!(cfg.name_node, None);
+ assert_eq!(cfg.options, None);
+ }
+
+ #[test]
+ fn test_hdfs_native_config_parse_normalizes_name_node() {
+ let parse = |value: &str| {
+ hdfs_native_config_parse(HashMap::from([(
+ HDFS_NAME_NODE.to_string(),
+ value.to_string(),
+ )]))
+ .unwrap()
+ .name_node
+ };
+
+ // Empty must not shadow the path-authority fallback.
+ assert_eq!(parse(""), None);
+ assert_eq!(parse(" "), None);
+ // Trailing `/` would otherwise yield a second cache entry for one
cluster.
+ assert_eq!(
+ parse(" hdfs://nn:8020/ ").as_deref(),
+ Some("hdfs://nn:8020")
+ );
+ }
+
+ #[test]
+ fn test_hdfs_native_effective_name_node_precedence() {
+ let configured = hdfs_native_config_parse(HashMap::from([(
+ HDFS_NAME_NODE.to_string(),
+ "hdfs://nn1:8020,hdfs://nn2:8020".to_string(),
+ )]))
+ .unwrap();
+ let unconfigured = HdfsNativeConfig::default();
+
+ // Configured wins over the authority, including for authority-less
paths.
+ for path in ["hdfs://ns-a/x", "hdfs:///y"] {
+ let (nn, _) = hdfs_native_effective_name_node(&configured,
path).unwrap();
+ assert_eq!(nn, "hdfs://nn1:8020,hdfs://nn2:8020");
+ }
+ // Otherwise the authority, including its port.
+ let (nn, rel) =
+ hdfs_native_effective_name_node(&unconfigured,
"hdfs://nn:9000/a/b").unwrap();
+ assert_eq!((nn.as_str(), rel), ("hdfs://nn:9000", "a/b"));
+ // Neither: a pointed error.
+ let err = hdfs_native_effective_name_node(&unconfigured,
"hdfs:///a").unwrap_err();
+ assert!(err.to_string().contains(HDFS_NAME_NODE));
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_with_authority_and_rel() {
+ let (nn, rel) =
hdfs_native_parse_path("hdfs://nameservice1/a/b").unwrap();
+
+ assert_eq!(nn.as_deref(), Some("hdfs://nameservice1"));
+ assert_eq!(rel, "a/b");
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_with_authority_and_port() {
+ let (nn, rel) = hdfs_native_parse_path("hdfs://nn:8020/foo").unwrap();
+
+ assert_eq!(nn.as_deref(), Some("hdfs://nn:8020"));
+ assert_eq!(rel, "foo");
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_with_authority_no_path() {
+ let (nn, rel) = hdfs_native_parse_path("hdfs://nameservice1").unwrap();
+
+ assert_eq!(nn.as_deref(), Some("hdfs://nameservice1"));
+ assert_eq!(rel, "");
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_with_authority_trailing_slash() {
+ let (nn, rel) =
hdfs_native_parse_path("hdfs://nameservice1/").unwrap();
+
+ assert_eq!(nn.as_deref(), Some("hdfs://nameservice1"));
+ assert_eq!(rel, "");
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_authority_less_returns_none() {
+ let (nn, rel) = hdfs_native_parse_path("hdfs:///a/b").unwrap();
+
+ assert_eq!(nn, None);
+ assert_eq!(rel, "a/b");
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_wrong_scheme_errors() {
+ let err = hdfs_native_parse_path("file:///tmp/x").unwrap_err();
+
+ assert!(err.to_string().contains("expected scheme `hdfs://`"));
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_invalid_url_errors() {
+ let err = hdfs_native_parse_path("not-a-url").unwrap_err();
+
+ assert!(err.to_string().contains("Invalid hdfs path"));
+ }
+
+ #[test]
+ fn test_hdfs_native_parse_path_non_hierarchical_errors() {
+ // `hdfs:x` parses as a valid non-hierarchical URL; it must be
+ // rejected rather than panic on slicing.
+ for path in ["hdfs:x", "hdfs:/x", "hdfs:"] {
+ let err = hdfs_native_parse_path(path).unwrap_err();
+ assert!(err.to_string().contains("expected scheme `hdfs://`"));
+ }
+ }
+
+ #[test]
+ fn test_hdfs_native_batch_key_distinguishes_ports() {
+ let config = HdfsNativeConfig::default();
+
+ assert_eq!(
+ hdfs_native_batch_key(&config, "hdfs://namenode:8020/a"),
+ "hdfs://namenode:8020"
+ );
+ assert_eq!(
+ hdfs_native_batch_key(&config, "hdfs://namenode:9000/b"),
+ "hdfs://namenode:9000"
+ );
+ }
+
+ #[test]
+ fn test_hdfs_native_batch_key_invalid_path_keys_on_itself() {
+ let config = HdfsNativeConfig::default();
+
+ // Unresolvable paths must not collapse onto a shared "" key.
+ assert_eq!(hdfs_native_batch_key(&config, "not-a-url"), "not-a-url");
+ assert_eq!(hdfs_native_batch_key(&config, "hdfs:///a"), "hdfs:///a");
+ }
+
+ #[test]
Review Comment:
These are plain `#[test]`s with no runtime, relying on
`Operator::from_config` being runtime-free today. If a future `hdfs-native`
builds the `Client` eagerly (it captures `Handle::try_current()` to spawn
tasks), these panic with "no reactor running" — the same eager-construction
path behind @comphead's runtime-pinning question. I'd switch them to
`#[tokio::test]` and document whether `from_config` is actually runtime-free,
so the contract is pinned rather than incidental.
##########
crates/storage/opendal/tests/file_io_hdfs_test.rs:
##########
@@ -0,0 +1,315 @@
+// 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.
+
+//! Integration tests for HDFS FileIO via OpenDAL `services-hdfs-native`.
+//!
+//! These tests need the HDFS fixture in `dev/docker-compose.yaml` and are
+//! skipped when `ICEBERG_TEST_HDFS_ENDPOINT` is not set. The fixture uses
+//! host networking (Linux, or a Docker runtime that supports it), so it sits
+//! behind a compose profile:
+//!
+//! ```text
+//! COMPOSE_PROFILES=hdfs make docker-up
+//! ICEBERG_TEST_HDFS_ENDPOINT=hdfs://localhost:8020 cargo test -p
iceberg-storage-opendal \
+//! --features opendal-hdfs-native --test file_io_hdfs_test
+//! ```
+
+#[cfg(feature = "opendal-hdfs-native")]
+mod tests {
+ use std::sync::Arc;
+
+ use bytes::Bytes;
+ use futures::StreamExt;
+ use iceberg::io::{FileIO, FileIOBuilder, HDFS_NAME_NODE};
+ use iceberg_storage_opendal::{OpenDalResolvingStorageFactory,
OpenDalStorageFactory};
+ use iceberg_test_utils::{
+ ENV_HDFS_ENDPOINT, get_hdfs_endpoint, normalize_test_name_with_parts,
set_up,
+ };
+
+ /// Skips the calling test unless the HDFS fixture endpoint is configured;
+ /// an unset *or* empty variable means "not provided" (see the HF tests).
+ macro_rules! require_hdfs {
+ () => {
+ match std::env::var(ENV_HDFS_ENDPOINT) {
+ Ok(v) if !v.is_empty() => {}
+ _ => {
+ eprintln!("Skipping HDFS test: {} not set",
ENV_HDFS_ENDPOINT);
+ return;
+ }
+ }
+ };
+ }
+
+ fn get_file_io() -> FileIO {
+ set_up();
+ FileIOBuilder::new(Arc::new(OpenDalStorageFactory::HdfsNative)).build()
+ }
+
+ fn test_path(suffix: &str) -> String {
+ format!(
+ "{}/{}",
+ get_hdfs_endpoint(),
+ normalize_test_name_with_parts!(suffix)
+ )
+ }
+
+ #[tokio::test]
+ async fn test_file_io_hdfs_exists() {
+ require_hdfs!();
+ let file_io = get_file_io();
+
+ let absent = test_path("test_file_io_hdfs_exists_absent");
+ assert!(!file_io.exists(&absent).await.unwrap());
+ }
+
+ #[tokio::test]
+ async fn test_file_io_hdfs_write_and_read() {
+ require_hdfs!();
+ let file_io = get_file_io();
+ let path = test_path("test_file_io_hdfs_write_and_read");
+ let _ = file_io.delete(&path).await;
+
+ let output = file_io.new_output(&path).unwrap();
+ output
+ .write(Bytes::from_static(b"hello hdfs"))
+ .await
+ .unwrap();
+
+ assert!(file_io.exists(&path).await.unwrap());
+ let input = file_io.new_input(&path).unwrap();
+ assert_eq!(
+ input.read().await.unwrap(),
+ Bytes::from_static(b"hello hdfs")
+ );
+ }
+
+ /// The HA flow: table locations carry a logical authority while
+ /// `hdfs.name-node` carries the (comma-separated) endpoints; it wins.
+ #[tokio::test]
+ async fn test_file_io_hdfs_configured_name_node() {
+ require_hdfs!();
+ set_up();
+ let file_io =
FileIOBuilder::new(Arc::new(OpenDalStorageFactory::HdfsNative))
+ .with_prop(HDFS_NAME_NODE, get_hdfs_endpoint())
+ .build();
+
+ // The path authority is a logical name; the configured NameNode wins.
+ let path = format!(
+ "hdfs://logical-nameservice/{}",
+
normalize_test_name_with_parts!("test_file_io_hdfs_configured_name_node")
+ );
+ let _ = file_io.delete(&path).await;
+
+ file_io
+ .new_output(&path)
+ .unwrap()
+ .write(Bytes::from_static(b"via configured name node"))
+ .await
+ .unwrap();
+
+ assert!(file_io.exists(&path).await.unwrap());
+ assert_eq!(
+ file_io.new_input(&path).unwrap().read().await.unwrap(),
+ Bytes::from_static(b"via configured name node")
+ );
+ }
+
+ #[tokio::test]
+ async fn test_file_io_hdfs_overwrite() {
+ require_hdfs!();
+ let file_io = get_file_io();
+ let path = test_path("test_file_io_hdfs_overwrite");
+ let _ = file_io.delete(&path).await;
+
+ for content in [b"first".as_slice(), b"second, longer".as_slice()] {
+ file_io
+ .new_output(&path)
+ .unwrap()
+ .write(Bytes::from_static(content))
+ .await
+ .unwrap();
+ }
+
+ assert_eq!(
+ file_io.new_input(&path).unwrap().read().await.unwrap(),
+ Bytes::from_static(b"second, longer")
+ );
+ }
+
+ #[tokio::test]
+ async fn test_file_io_hdfs_delete_stream() {
Review Comment:
The batch-key-by-effective-NameNode logic is the one real behavioral
difference from the other backends, and it's never exercised end-to-end with
two distinct NameNodes — this test uses one. I'd add a case with paths across
two effective NameNodes (on the single-node fixture, `localhost` vs `127.0.0.1`
gives two keys) asserting both groups delete. A concurrent cache-insert unit
test racing two callers for the same NameNode (asserting `len() == 1`) would
also lock in the `or_insert` race handling, and that one needs no real HDFS.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,411 @@
+// 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::sync::{Arc, RwLock};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_NAME_NODE};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+use crate::utils::from_opendal_error;
+
+/// Parse iceberg properties to [`HdfsNativeConfig`].
+pub(crate) fn hdfs_native_config_parse(mut m: HashMap<String, String>) ->
Result<HdfsNativeConfig> {
+ let mut cfg = HdfsNativeConfig::default();
+
+ // `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)
+ .map(|s| s.trim().trim_end_matches('/').to_string())
Review Comment:
`trim_end_matches('/')` only trims the end of the whole string, so for a
comma-separated HA list like `hdfs://nn1:8020/,hdfs://nn2:8020/` the first
entry keeps its slash and we cache two operators for one cluster. I'd split on
`,` and trim each entry before rejoining.
##########
crates/iceberg/src/io/storage/config/hdfs.rs:
##########
@@ -0,0 +1,27 @@
+// 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 configuration.
+
+/// HDFS NameNode RPC endpoint(s), e.g. `hdfs://namenode:8020`; a
+/// comma-separated list enables HA failover. When unset, the NameNode is
+/// derived from the path authority.
+pub const HDFS_NAME_NODE: &str = "hdfs.name-node";
+/// Prefix for properties forwarded to the HDFS client configuration, e.g.
+/// `hadoop.dfs.client.failover.random.order`. Forwarded values (prefix
+/// stripped) override those loaded from `$HADOOP_CONF_DIR`.
+pub const HDFS_HADOOP_CONF_PREFIX: &str = "hadoop.";
Review Comment:
The doc doesn't quite say the `hadoop.` prefix is stripped before
forwarding. Worth spelling out — something like "property `hadoop.<key>` is
forwarded to the HDFS client config as `<key>`".
##########
crates/test_utils/src/lib.rs:
##########
@@ -47,13 +47,15 @@ mod common {
pub const ENV_HMS_ENDPOINT: &str = "ICEBERG_TEST_HMS_ENDPOINT";
pub const ENV_GLUE_ENDPOINT: &str = "ICEBERG_TEST_GLUE_ENDPOINT";
pub const ENV_GCS_ENDPOINT: &str = "ICEBERG_TEST_GCS_ENDPOINT";
+ pub const ENV_HDFS_ENDPOINT: &str = "ICEBERG_TEST_HDFS_ENDPOINT";
// Default ports matching dev/docker-compose.yaml
pub const DEFAULT_OBJECT_STORE_PORT: u16 = 9000;
pub const DEFAULT_REST_CATALOG_PORT: u16 = 8181;
pub const DEFAULT_HMS_PORT: u16 = 9083;
pub const DEFAULT_GLUE_PORT: u16 = 5001;
pub const DEFAULT_GCS_PORT: u16 = 4443;
+ pub const DEFAULT_HDFS_NN_PORT: u16 = 8020;
Review Comment:
The siblings are all `DEFAULT_<SERVICE>_PORT` — this one slips in an `_NN_`
infix. I'd rename to `DEFAULT_HDFS_PORT` for consistency.
--
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]