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


##########
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:
   Warn-and-drop it is: opendal's `HdfsNativeConfig` has no user field and 
hdfs-native has no config key for it (`with_user` exists on its builder, but 
opendal never calls it), so mapping would need an opendal change first. An 
error would reject every PyIceberg config that just carries `hdfs.user`. Both 
keys now log a warning naming the key, and a test pins that they are never 
forwarded as Hadoop options. Lands with the next push.



##########
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:
   `hdfs.port` is now validated whenever it is present, host or not, and a 
valid port without `hdfs.host` logs a warning instead of vanishing. Tests cover 
both. Lands with the next push.



##########
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:
   Fixed: a bare `hadoop.` key is rejected with a `DataInvalid` naming the 
prefix, with a test. Lands with the next push.



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