mixermt commented on code in PR #3111:
URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4174790837


##########
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:
   Checked this: `Url::host_str()` returns the host slice of the serialized 
URL, so IPv6 keeps its brackets. `hdfs://[::1]:8020/a/b` already resolves to 
`hdfs://[::1]:8020` on the current code; I added 
`test_hdfs_native_parse_path_ipv6_authority_keeps_brackets` to pin it. I also 
pointed the client at a local `[::1]` listener through the path authority and 
it connects fine. `url.host()` would give the same string, so I left the code 
as is.



##########
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:
   Done: `relativize_path` now goes through `hdfs_native_effective_name_node`, 
so an authority-less path errors the same way in both unless a NameNode is 
configured; test updated. For the record it was not reachable (`delete_stream` 
only relativizes paths whose batch deleter `create_operator` already built), 
but one rule for all callers is cleaner.



##########
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:
   The representation is already private: `public-api.txt` only shows 
`HdfsNativeOperatorCache(_)` with Clone/Default/Debug, so swapping the 
internals is not a breaking change. Enum variant fields are always public in 
Rust, so a cached variant has to expose some opaque type either way; 
`HdfsNative(Arc<HdfsNativeState>)` would just rename it and hide `config`, 
which every other variant exposes. `Memory(Operator)` and 
`S3::customized_credential_load` already carry non-config state the same way. I 
added the lifecycle note to the cache doc as suggested.



##########
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:
   The duplicate-cache part doesn't happen: the configured list is the only key 
inside a storage, so paths with different authorities still share one operator 
(test extended to pin that), and opendal strips trailing slashes per entry 
anyway. But trimming per entry is still the right call, for a different reason: 
a space after the comma (`a, b`) was passed through as-is and silently broke 
failover to that NameNode; I verified that against a local listener. Now each 
entry is trimmed and empties dropped, with tests.



##########
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:
   I'd keep `HdfsNative`. blackmwk's request was actually to rename the module 
to `hdfs_native` (I renamed the variants to match and offered to revert; nobody 
asked). opendal has two HDFS services, `services-hdfs` (libhdfs/JNI) and 
`services-hdfs-native`, and #1130's design names them `Hdfs` and `HdfsNative` 
respectively, with the maintainers wanting both selectable eventually. Calling 
this one `Hdfs` would take the libhdfs backend's name. The crate already names 
variants after opendal services (`Azdls`, `Hf`), so this follows the convention.



##########
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:
   Added both. `test_file_io_hdfs_delete_stream_two_name_nodes` writes under 
`localhost` and `127.0.0.1`, checks each spelling sees the other's file, 
deletes through one stream and checks all four spellings are gone. Since the 
fixture is a single cluster it exercises the two-batch path rather than 
cross-cluster routing, which the batch-key unit tests pin. 
`test_hdfs_native_operator_cache_keeps_first_insert` inserts two 
distinguishable operators for one key, asserts the first is kept and returned 
to both callers, then races eight threads through `create_operator`.



##########
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:
   Renamed to `DEFAULT_HDFS_PORT`.



##########
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:
   Reworded: the doc now says `hadoop.dfs.client.failover.random.order` is 
forwarded as `dfs.client.failover.random.order`, and mentions 
`hadoop.fs.defaultFS`.



-- 
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]

Reply via email to