mixermt commented on code in PR #3111: URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4194129589
########## 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: Done this time. Since the last round the crate depends on tokio anyway (the operator cache tracks runtime liveness with it), so there was no longer a reason to keep the build inline: `create_operator` is async and the Hadoop XML read runs under `spawn_blocking`, outside the cache lock as before. The builder itself stays sync and runtime-free, which the plain unit test still pins. 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]
