laskoviymishka commented on code in PR #3111:
URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4187742576
##########
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:
This is @comphead's public-API point from round 1, still open: `operators`
is a public field on a public variant, so `HdfsNativeOperatorCache` is now
permanent semver surface and callers can construct the variant with an
arbitrary cache — the siblings only expose `config: Arc<...>`. Folding both
fields behind one opaque type (`HdfsNative(HdfsNativeStorage)` with private
fields), or `#[non_exhaustive]` on the variant plus `pub(crate)` on the cache,
keeps it off the stable surface. Worth settling before this is published rather
than after.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,635 @@
+// 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_HOST, HDFS_NAME_NODE,
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+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;
+
+/// 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: Vec<&str> = name_node
+ .split(',')
+ .map(|entry| entry.trim().trim_end_matches('/'))
+ .filter(|entry| !entry.is_empty())
+ .collect();
+ // hdfs-native dials each entry as `host:port` (`hdfs://` optional);
+ // without a port it fails only at the first I/O.
+ if let Some(entry) = entries.iter().find(|entry| {
+ entry
+ .rsplit_once(':')
+ .is_none_or(|(_, port)| port.parse::<u16>().is_err())
+ }) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_NAME_NODE}` entry: {entry}, expected
host:port"),
+ ));
+ }
+ if !entries.is_empty() {
+ cfg.name_node = Some(entries.join(","));
+ }
+ }
+
+ let host = m.remove(HDFS_HOST).map(|s| s.trim().to_string());
+ let port = m.remove(HDFS_PORT).map(|s| s.trim().to_string());
+
+ 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();
+ // 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.
+ if let Some(host) = host.filter(|s| !s.is_empty()) {
+ let port = match port.filter(|s| !s.is_empty()) {
+ Some(port) => port.parse::<u16>().map_err(|e| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_PORT}`: {port}: {e}"),
+ )
+ })?,
+ None => HDFS_DEFAULT_PORT,
+ };
+ // An IPv6 literal needs brackets in a URI authority.
+ let host = if host.contains(':') && !host.starts_with('[') {
+ format!("[{host}]")
+ } else {
+ host
+ };
+ options
+ .entry(FS_DEFAULT_FS.to_string())
+ .or_insert_with(|| format!("hdfs://{host}:{port}"));
+ }
+ 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, else `fs.defaultFS` —
+/// plus the relative path. 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_name_node, relative_path) = hdfs_native_parse_path(path)?;
+ let name_node = config
+ .name_node
+ .clone()
+ .or(authority_name_node)
Review Comment:
`config.name_node.clone().or(authority_name_node)` means a configured
`hdfs.name-node` wins over the authority in *every* path, so
`hdfs://other-cluster:8020/...` silently routes to the configured cluster — a
wrong-cluster read/write with no error. That's exactly right for an HA logical
name, but it also swallows a real second cluster. I'd apply the configured
name-node only when the path authority is absent or isn't a concrete
`host:port` (i.e. a logical nameservice), and otherwise let the authority
stand. If we'd rather keep the simple precedence, the `HDFS_NAME_NODE` doc
needs to call out the all-paths-go-to-one-cluster behavior explicitly, with a
test pinning the mixed-cluster case.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,635 @@
+// 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_HOST, HDFS_NAME_NODE,
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+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;
+
+/// 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: Vec<&str> = name_node
+ .split(',')
+ .map(|entry| entry.trim().trim_end_matches('/'))
+ .filter(|entry| !entry.is_empty())
+ .collect();
+ // hdfs-native dials each entry as `host:port` (`hdfs://` optional);
+ // without a port it fails only at the first I/O.
+ if let Some(entry) = entries.iter().find(|entry| {
+ entry
+ .rsplit_once(':')
+ .is_none_or(|(_, port)| port.parse::<u16>().is_err())
+ }) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_NAME_NODE}` entry: {entry}, expected
host:port"),
+ ));
+ }
+ if !entries.is_empty() {
+ cfg.name_node = Some(entries.join(","));
+ }
+ }
+
+ let host = m.remove(HDFS_HOST).map(|s| s.trim().to_string());
+ let port = m.remove(HDFS_PORT).map(|s| s.trim().to_string());
+
+ 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();
+ // 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.
+ if let Some(host) = host.filter(|s| !s.is_empty()) {
+ let port = match port.filter(|s| !s.is_empty()) {
+ Some(port) => port.parse::<u16>().map_err(|e| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_PORT}`: {port}: {e}"),
+ )
+ })?,
+ None => HDFS_DEFAULT_PORT,
+ };
+ // An IPv6 literal needs brackets in a URI authority.
+ let host = if host.contains(':') && !host.starts_with('[') {
+ format!("[{host}]")
+ } else {
+ host
+ };
+ options
+ .entry(FS_DEFAULT_FS.to_string())
+ .or_insert_with(|| format!("hdfs://{host}:{port}"));
+ }
+ 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, else `fs.defaultFS` —
+/// plus the relative path. 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_name_node, relative_path) = hdfs_native_parse_path(path)?;
+ let name_node = config
+ .name_node
+ .clone()
+ .or(authority_name_node)
+ .or_else(|| hdfs_native_default_fs(config))
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid hdfs path: {path}, authority-less paths require
`{HDFS_NAME_NODE}` or `{HDFS_HOST}`"
+ ),
+ )
+ })?;
+ Ok((name_node, relative_path))
+}
+
+/// `fs.defaultFS` from the forwarded options, when it is an HDFS URI.
+fn hdfs_native_default_fs(config: &HdfsNativeConfig) -> Option<String> {
+ config
+ .options
+ .as_ref()?
+ .get(FS_DEFAULT_FS)
+ .map(|s| s.trim().trim_end_matches('/'))
+ .filter(|s| {
+ s.strip_prefix("hdfs://")
+ .is_some_and(|rest| !rest.is_empty())
+ })
+ .map(str::to_string)
+}
+
+/// 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 (a private one when built outside any). The
+/// cache lives as long as the storage that owns it (clones share it) and
+/// never evicts, so the storage must not be used from another runtime once
+/// the building one is dropped: `hdfs-native` panics on a dead runtime.
+#[derive(Clone, Debug, Default)]
+pub struct HdfsNativeOperatorCache(Arc<RwLock<HashMap<String, Operator>>>);
Review Comment:
The doc here is honest about the hazard but leaves it as a landmine. The
cache never evicts, and each cached Operator pins its RPC tasks to whichever
runtime was current when it was first built — so a FileIO/catalog that outlives
that runtime (reused across `#[tokio::test]`s, rebuilt after a worker shuts
down, shared across runtimes in a server) aborts the process the next time it
touches HDFS, because hdfs-native panics on a dead runtime. That's reachable
from a plain library call, which is the never-panic line I don't want us to
cross.
Keying the cache on the runtime id and rebuilding on a miss, or just
building per call like Gcs/Oss/Azdls, both close it. If it's genuinely out of
scope for an experimental backend, I'd want this as a loud warning on the
variant doc and the README rather than a comment on a private struct — plus a
test that actually reproduces the cross-runtime abort.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,635 @@
+// 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_HOST, HDFS_NAME_NODE,
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+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;
+
+/// 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: Vec<&str> = name_node
+ .split(',')
+ .map(|entry| entry.trim().trim_end_matches('/'))
+ .filter(|entry| !entry.is_empty())
+ .collect();
+ // hdfs-native dials each entry as `host:port` (`hdfs://` optional);
+ // without a port it fails only at the first I/O.
+ if let Some(entry) = entries.iter().find(|entry| {
Review Comment:
We reject a configured `hdfs.name-node` without a port here, but a path
authority (`hdfs://ns1/x`) and `hadoop.fs.defaultFS` (`hdfs://ns1`) both sail
through portless, and `operator_build` then feeds that portless string straight
to hdfs-native. So the two inputs disagree: either portless is a valid logical
nameservice — in which case this check is too strict and rejects the natural HA
spelling `hdfs://ns1` — or hdfs-native dials `host:port` and the portless
authority fails at first I/O with a confusing error, the very thing this check
exists to prevent. Worth confirming what hdfs-native actually does with a
portless NameNode, then applying one rule to all three sources and pinning it
with a test.
##########
crates/iceberg/src/io/storage/config/hdfs.rs:
##########
@@ -0,0 +1,37 @@
+// 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. Takes precedence over the path
+/// authority; when unset, the NameNode is derived from the path authority.
+pub const HDFS_NAME_NODE: &str = "hdfs.name-node";
+/// NameNode host for authority-less paths, as in PyIceberg; paths that carry
+/// an authority ignore it, as they do there. Combined with [`HDFS_PORT`] into
+/// Hadoop's `fs.defaultFS`. PyIceberg's `hdfs.user` and `hdfs.kerberos_ticket`
+/// have no equivalent: the client reads `HADOOP_USER_NAME` and the default
+/// Kerberos credential cache.
+pub const HDFS_HOST: &str = "hdfs.host";
+/// NameNode port for [`HDFS_HOST`]; defaults to `8020`.
Review Comment:
The doc is upfront that `hdfs.user` and `hdfs.kerberos_ticket` have no
equivalent, which is good — but a config copied over from PyIceberg will then
silently run as a different identity (`HADOOP_USER_NAME` / OS user / default
ccache) with no signal. I'd lean toward a `log::warn!` when either key is
present so the drop is at least visible, or an outright `FeatureUnsupported` if
we'd rather fail loud. Do you see a clean path to mapping `hdfs.user` through
hdfs-native's user hook, or is warn-and-drop the pragmatic line for now?
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,635 @@
+// 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_HOST, HDFS_NAME_NODE,
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+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;
+
+/// 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: Vec<&str> = name_node
+ .split(',')
+ .map(|entry| entry.trim().trim_end_matches('/'))
+ .filter(|entry| !entry.is_empty())
+ .collect();
+ // hdfs-native dials each entry as `host:port` (`hdfs://` optional);
+ // without a port it fails only at the first I/O.
+ if let Some(entry) = entries.iter().find(|entry| {
+ entry
+ .rsplit_once(':')
+ .is_none_or(|(_, port)| port.parse::<u16>().is_err())
+ }) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_NAME_NODE}` entry: {entry}, expected
host:port"),
+ ));
+ }
+ if !entries.is_empty() {
+ cfg.name_node = Some(entries.join(","));
+ }
+ }
+
+ let host = m.remove(HDFS_HOST).map(|s| s.trim().to_string());
+ let port = m.remove(HDFS_PORT).map(|s| s.trim().to_string());
+
+ 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();
+ // 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.
+ if let Some(host) = host.filter(|s| !s.is_empty()) {
+ let port = match port.filter(|s| !s.is_empty()) {
+ Some(port) => port.parse::<u16>().map_err(|e| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_PORT}`: {port}: {e}"),
+ )
+ })?,
+ None => HDFS_DEFAULT_PORT,
+ };
+ // An IPv6 literal needs brackets in a URI authority.
+ let host = if host.contains(':') && !host.starts_with('[') {
+ format!("[{host}]")
+ } else {
+ host
+ };
+ options
+ .entry(FS_DEFAULT_FS.to_string())
+ .or_insert_with(|| format!("hdfs://{host}:{port}"));
+ }
+ 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, else `fs.defaultFS` —
+/// plus the relative path. 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_name_node, relative_path) = hdfs_native_parse_path(path)?;
+ let name_node = config
+ .name_node
+ .clone()
+ .or(authority_name_node)
+ .or_else(|| hdfs_native_default_fs(config))
+ .ok_or_else(|| {
+ Error::new(
+ ErrorKind::DataInvalid,
+ format!(
+ "Invalid hdfs path: {path}, authority-less paths require
`{HDFS_NAME_NODE}` or `{HDFS_HOST}`"
+ ),
+ )
+ })?;
+ Ok((name_node, relative_path))
+}
+
+/// `fs.defaultFS` from the forwarded options, when it is an HDFS URI.
+fn hdfs_native_default_fs(config: &HdfsNativeConfig) -> Option<String> {
+ config
+ .options
+ .as_ref()?
+ .get(FS_DEFAULT_FS)
+ .map(|s| s.trim().trim_end_matches('/'))
+ .filter(|s| {
+ s.strip_prefix("hdfs://")
+ .is_some_and(|rest| !rest.is_empty())
+ })
+ .map(str::to_string)
+}
+
+/// 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 (a private one when built outside any). The
+/// cache lives as long as the storage that owns it (clones share it) and
+/// never evicts, so the storage must not be used from another runtime once
+/// the building one is dropped: `hdfs-native` panics on a dead runtime.
+#[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));
+ }
+
+ // The build reads the Hadoop XML config synchronously (~0.1 ms, once
+ // per NameNode), so it runs outside the lock and is not worth a
+ // blocking-thread hop. A racing first caller may build too; the loser
+ // is dropped before opening any connection.
+ let op = hdfs_native_operator_build(config, &name_node)?;
Review Comment:
This is the synchronous Hadoop XML read I flagged last round — it's
documented now, but still running on the executor thread. The `~0.1 ms` holds
for a local `$HADOOP_CONF_DIR`; on a slow or networked mount it's a blocking
call on an async worker, and when no `hdfs.name-node` is configured it fires
once per distinct NameNode seen in metadata. I'd still rather wrap the build in
`spawn_blocking`. If we're consciously accepting it for the experimental
backend, that's a fine call — just say so at the call site rather than only
describing the timing.
##########
.github/workflows/ci.yml:
##########
@@ -241,6 +241,9 @@ jobs:
- name: Start Docker containers
if: matrix.test-suite.name == 'default'
+ env:
+ # The HDFS fixture needs host networking; opt in here (Linux only).
+ COMPOSE_PROFILES: hdfs
Review Comment:
Opting the HDFS fixture in on the default suite means every default job now
pulls the Hadoop image and waits on two container health checks (up to 30×5s)
before any test runs — a heavy fixed cost on every PR for an experimental
backend, and the whole workspace suite goes red if the fixture flakes. I'd gate
this to a separate job (or to `crates/storage/opendal/**` changes) so an HDFS
hiccup can't block unrelated work. Worth double-checking `make docker-up`
passes `--wait` and honors the profile so a slow start fails cleanly rather
than racing the tests.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,635 @@
+// 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_HOST, HDFS_NAME_NODE,
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+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;
+
+/// 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: Vec<&str> = name_node
+ .split(',')
+ .map(|entry| entry.trim().trim_end_matches('/'))
+ .filter(|entry| !entry.is_empty())
+ .collect();
+ // hdfs-native dials each entry as `host:port` (`hdfs://` optional);
+ // without a port it fails only at the first I/O.
+ if let Some(entry) = entries.iter().find(|entry| {
+ entry
+ .rsplit_once(':')
+ .is_none_or(|(_, port)| port.parse::<u16>().is_err())
+ }) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_NAME_NODE}` entry: {entry}, expected
host:port"),
+ ));
+ }
+ if !entries.is_empty() {
+ cfg.name_node = Some(entries.join(","));
+ }
+ }
+
+ let host = m.remove(HDFS_HOST).map(|s| s.trim().to_string());
+ let port = m.remove(HDFS_PORT).map(|s| s.trim().to_string());
+
+ 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();
+ // 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.
+ if let Some(host) = host.filter(|s| !s.is_empty()) {
Review Comment:
Port handling is gated behind `hdfs.host` being set, so `hdfs.port` on its
own is silently dropped and a bad value like `hdfs.port=x` is only caught when
a host is also present. A `tracing::warn!` on a port with no host (or
validating it regardless) would avoid a config that looks applied but isn't.
##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,635 @@
+// 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_HOST, HDFS_NAME_NODE,
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use url::Url;
+
+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;
+
+/// 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: Vec<&str> = name_node
+ .split(',')
+ .map(|entry| entry.trim().trim_end_matches('/'))
+ .filter(|entry| !entry.is_empty())
+ .collect();
+ // hdfs-native dials each entry as `host:port` (`hdfs://` optional);
+ // without a port it fails only at the first I/O.
+ if let Some(entry) = entries.iter().find(|entry| {
+ entry
+ .rsplit_once(':')
+ .is_none_or(|(_, port)| port.parse::<u16>().is_err())
+ }) {
+ return Err(Error::new(
+ ErrorKind::DataInvalid,
+ format!("Invalid `{HDFS_NAME_NODE}` entry: {entry}, expected
host:port"),
+ ));
+ }
+ if !entries.is_empty() {
+ cfg.name_node = Some(entries.join(","));
+ }
+ }
+
+ let host = m.remove(HDFS_HOST).map(|s| s.trim().to_string());
+ let port = m.remove(HDFS_PORT).map(|s| s.trim().to_string());
+
+ let mut options: HashMap<String, String> = m
+ .into_iter()
+ .filter_map(|(key, value)| {
+ key.strip_prefix(HDFS_HADOOP_CONF_PREFIX)
Review Comment:
A bare `hadoop.` key strips to `""` and gets forwarded to opendal as an
empty option key. Worth filtering out the empty stripped key here.
##########
crates/storage/opendal/tests/file_io_hdfs_test.rs:
##########
@@ -0,0 +1,353 @@
+// 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() {
+ require_hdfs!();
+ let file_io = get_file_io();
+
+ let paths: Vec<String> = (0..5)
+ .map(|i| format!("{}/file-{i}",
test_path("test_file_io_hdfs_delete_stream")))
+ .collect();
+ for path in &paths {
+ let _ = file_io.delete(path).await;
+ file_io
+ .new_output(path)
+ .unwrap()
+ .write("delete-me".into())
+ .await
+ .unwrap();
+ assert!(file_io.exists(path).await.unwrap());
+ }
+
+ let stream = futures::stream::iter(paths.clone()).boxed();
+ file_io.delete_stream(stream).await.unwrap();
+
+ for path in &paths {
+ assert!(!file_io.exists(path).await.unwrap());
+ }
+ }
+
+ /// Paths are batched per effective NameNode; two spellings of the
fixture's
+ /// NameNode drive two batches through one stream. The fixture is one
+ /// cluster, so the keying itself is pinned by the unit tests.
+ #[tokio::test]
+ async fn test_file_io_hdfs_delete_stream_two_name_nodes() {
+ require_hdfs!();
+ let file_io = get_file_io();
+ let dir = test_path("test_file_io_hdfs_delete_stream_two_name_nodes");
+ let alt_dir = dir.replacen("localhost", "127.0.0.1", 1);
Review Comment:
`dir.replacen("localhost", "127.0.0.1", 1)` is a no-op unless the endpoint
literally contains `localhost` — against `hdfs://namenode:8020` `alt_dir ==
dir`, both paths key to one batch, and the test passes having exercised
nothing. An `assert_ne!(dir, alt_dir)` (or skipping when the endpoint doesn't
contain `localhost`) would keep it honest. Worth confirming it's stable on the
GH runners too — `localhost` vs `127.0.0.1` can resolve to different stacks.
--
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]