mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-03 18:45:35 +00:00
feat: add HDFS object storage backend (#8701)
* feat: add HDFS object storage backend Signed-off-by: Minghan2005 <cambrianocean@gmail.com> * fix: make HDFS storage operations durable Gate the native HDFS backend behind an explicit feature. Publish writes through same-directory temporary files and atomic HDFS Rename2 replacement, and provide streaming copy fallback for COPY_REGION. Add regression coverage for interrupted writes and the region-copy path. Signed-off-by: Minghan2005 <cambrianocean@gmail.com> * ci: run HDFS object store tests Signed-off-by: jeremyhi <fengjiachun@gmail.com> * docs: note HDFS temporary file cleanup follow-up Signed-off-by: jeremyhi <fengjiachun@gmail.com> * feat: enable HDFS object storage by default Signed-off-by: jeremyhi <fengjiachun@gmail.com> * docs: remove redundant HDFS build feature notes Signed-off-by: jeremyhi <fengjiachun@gmail.com> --------- Signed-off-by: Minghan2005 <cambrianocean@gmail.com> Signed-off-by: jeremyhi <fengjiachun@gmail.com> Co-authored-by: Minghan2005 <cambrianocean@gmail.com> Co-authored-by: jeremyhi <fengjiachun@gmail.com>
This commit is contained in:
co-authored by
Minghan2005
jeremyhi
parent
ffdd6d09a6
commit
75bd8e9ce6
@@ -22,6 +22,7 @@ required-features = ["dev-tools"]
|
||||
[features]
|
||||
default = [
|
||||
"ai_functions",
|
||||
"hdfs-object-store",
|
||||
"servers/pprof",
|
||||
"servers/mem-prof",
|
||||
"meta-srv/pg_kvbackend",
|
||||
@@ -32,6 +33,7 @@ enterprise = ["common-meta/enterprise", "frontend/enterprise", "meta-srv/enterpr
|
||||
# Developer-only helper binaries and diagnostic datanode commands.
|
||||
# Kept out of `default` so normal/release builds don't compile them.
|
||||
dev-tools = []
|
||||
hdfs-object-store = ["mito2/hdfs-object-store", "object-store/hdfs-object-store"]
|
||||
mysql-object-store = ["object-store/mysql-object-store"]
|
||||
tokio-console = ["common-telemetry/tokio-console"]
|
||||
|
||||
|
||||
@@ -125,6 +125,44 @@ fn test_load_runtime_options_without_max_blocking_threads() {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
fn test_load_datanode_hdfs_config() {
|
||||
let config = tempfile::NamedTempFile::new().unwrap();
|
||||
std::fs::write(
|
||||
config.path(),
|
||||
r#"
|
||||
[storage]
|
||||
type = "Hdfs"
|
||||
name = "local-hdfs"
|
||||
root = "/greptimedb"
|
||||
name_node = "hdfs://127.0.0.1:9000"
|
||||
enable_read_cache = false
|
||||
options = { "dfs.client.block.write.replace-datanode-on-failure.enable" = "true" }
|
||||
"#,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let options =
|
||||
GreptimeOptions::<DatanodeOptions>::load_layered_options(config.path().to_str(), "")
|
||||
.unwrap();
|
||||
let object_store::config::ObjectStoreConfig::Hdfs(hdfs) = options.component.storage.store
|
||||
else {
|
||||
unreachable!()
|
||||
};
|
||||
|
||||
assert_eq!("local-hdfs", hdfs.name);
|
||||
assert_eq!("/greptimedb", hdfs.connection.root);
|
||||
assert_eq!("hdfs://127.0.0.1:9000", hdfs.connection.name_node);
|
||||
assert_eq!(
|
||||
Some(&"true".to_string()),
|
||||
hdfs.connection
|
||||
.options
|
||||
.get("dfs.client.block.write.replace-datanode-on-failure.enable")
|
||||
);
|
||||
assert!(!hdfs.cache.enable_read_cache);
|
||||
}
|
||||
|
||||
#[allow(deprecated)]
|
||||
#[test]
|
||||
fn test_load_datanode_example_config() {
|
||||
|
||||
@@ -10,6 +10,7 @@ test = ["common-test-util", "rstest", "rstest_reuse", "rskafka"]
|
||||
testing = ["test"]
|
||||
test-shared-fs-region-migration = []
|
||||
enterprise = []
|
||||
hdfs-object-store = ["object-store/hdfs-object-store"]
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
|
||||
@@ -265,7 +265,20 @@ impl RegionFileCopier {
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use object_store::ObjectStore;
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use object_store::layers::HdfsCompatibilityLayer;
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use object_store::services::Fs;
|
||||
|
||||
use super::*;
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use crate::access_layer::AccessLayer;
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use crate::sst::index::intermediate::IntermediateManager;
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use crate::sst::index::puffin_manager::PuffinManagerFactory;
|
||||
|
||||
#[test]
|
||||
fn test_build_copy_file_paths() {
|
||||
@@ -344,4 +357,54 @@ mod tests {
|
||||
format!("/table_dir/1_0000000002/index/{}.1.puffin", file_id)
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
#[tokio::test]
|
||||
async fn test_copy_region_files_with_hdfs_fallback() {
|
||||
let (temp_dir, puffin_manager) =
|
||||
PuffinManagerFactory::new_for_test_async("hdfs-copy-region").await;
|
||||
let intermediate_manager = IntermediateManager::init_fs(temp_dir.path().to_string_lossy())
|
||||
.await
|
||||
.unwrap();
|
||||
let storage_dir = temp_dir.path().join("storage");
|
||||
std::fs::create_dir(&storage_dir).unwrap();
|
||||
let object_store = ObjectStore::new(Fs::default().root(storage_dir.to_str().unwrap()))
|
||||
.unwrap()
|
||||
.layer(HdfsCompatibilityLayer::new_for_test());
|
||||
let access_layer = Arc::new(AccessLayer::new(
|
||||
"table_dir",
|
||||
PathType::Bare,
|
||||
object_store.clone(),
|
||||
puffin_manager,
|
||||
intermediate_manager,
|
||||
));
|
||||
let copier = RegionFileCopier::new(access_layer);
|
||||
let source_region_id = RegionId::new(1, 1);
|
||||
let target_region_id = RegionId::new(1, 2);
|
||||
let file_id = FileId::random();
|
||||
let descriptor = FileDescriptor::Data { file_id, size: 8 };
|
||||
let (source_path, target_path) = build_copy_file_paths(
|
||||
source_region_id,
|
||||
target_region_id,
|
||||
descriptor,
|
||||
"table_dir",
|
||||
PathType::Bare,
|
||||
);
|
||||
object_store.write(&source_path, "contents").await.unwrap();
|
||||
|
||||
copier
|
||||
.copy_files(source_region_id, target_region_id, vec![descriptor], 1)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
b"contents",
|
||||
object_store
|
||||
.read(&target_path)
|
||||
.await
|
||||
.unwrap()
|
||||
.to_bytes()
|
||||
.as_ref()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -8,6 +8,7 @@ license.workspace = true
|
||||
workspace = true
|
||||
|
||||
[features]
|
||||
hdfs-object-store = ["dep:hdfs-native", "opendal/services-hdfs-native", "uuid"]
|
||||
mysql-object-store = ["opendal/services-mysql"]
|
||||
services-memory = ["opendal/services-memory"]
|
||||
testing = ["derive_builder", "uuid"]
|
||||
@@ -21,6 +22,7 @@ common-macro.workspace = true
|
||||
common-runtime.workspace = true
|
||||
common-telemetry.workspace = true
|
||||
derive_builder = { workspace = true, optional = true }
|
||||
hdfs-native = { workspace = true, optional = true }
|
||||
humantime-serde.workspace = true
|
||||
opendal = { version = "0.59.2", features = [
|
||||
"layers-tracing",
|
||||
|
||||
@@ -12,10 +12,14 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use std::collections::HashMap;
|
||||
use std::time::Duration;
|
||||
|
||||
use common_base::readable_size::ReadableSize;
|
||||
use common_base::secrets::{ExposeSecret, SecretString};
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use opendal::services::HdfsNative;
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
use opendal::services::Mysql;
|
||||
use opendal::services::{Azblob, Gcs, Oss, S3};
|
||||
@@ -34,6 +38,8 @@ pub enum ObjectStoreConfig {
|
||||
Oss(OssConfig),
|
||||
Azblob(AzblobConfig),
|
||||
Gcs(GcsConfig),
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
Hdfs(HdfsConfig),
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
Mysql(MysqlConfig),
|
||||
}
|
||||
@@ -53,6 +59,8 @@ impl ObjectStoreConfig {
|
||||
Self::Oss(_) => "Oss",
|
||||
Self::Azblob(_) => "Azblob",
|
||||
Self::Gcs(_) => "Gcs",
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
Self::Hdfs(_) => "Hdfs",
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
Self::Mysql(_) => "Mysql",
|
||||
}
|
||||
@@ -72,6 +80,8 @@ impl ObjectStoreConfig {
|
||||
Self::Oss(oss) => &oss.name,
|
||||
Self::Azblob(az) => &az.name,
|
||||
Self::Gcs(gcs) => &gcs.name,
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
Self::Hdfs(hdfs) => &hdfs.name,
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
Self::Mysql(mysql) => &mysql.name,
|
||||
};
|
||||
@@ -91,6 +101,8 @@ impl ObjectStoreConfig {
|
||||
Self::Oss(oss) => Some(&oss.cache),
|
||||
Self::Azblob(az) => Some(&az.cache),
|
||||
Self::Gcs(gcs) => Some(&gcs.cache),
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
Self::Hdfs(hdfs) => Some(&hdfs.cache),
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
Self::Mysql(mysql) => Some(&mysql.cache),
|
||||
}
|
||||
@@ -104,6 +116,8 @@ impl ObjectStoreConfig {
|
||||
Self::Oss(oss) => Some(&mut oss.cache),
|
||||
Self::Azblob(az) => Some(&mut az.cache),
|
||||
Self::Gcs(gcs) => Some(&mut gcs.cache),
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
Self::Hdfs(hdfs) => Some(&mut hdfs.cache),
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
Self::Mysql(mysql) => Some(&mut mysql.cache),
|
||||
}
|
||||
@@ -289,6 +303,42 @@ impl From<&GcsConnection> for Gcs {
|
||||
}
|
||||
}
|
||||
|
||||
/// Connection options for a Hadoop Distributed File System backend.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
|
||||
#[serde(default)]
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
pub struct HdfsConnection {
|
||||
/// Working directory for all object-store operations.
|
||||
pub root: String,
|
||||
/// HDFS NameNode URI, for example `hdfs://127.0.0.1:9000`.
|
||||
pub name_node: String,
|
||||
/// Additional options passed to the native HDFS client.
|
||||
pub options: HashMap<String, String>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
impl From<&HdfsConnection> for HdfsNative {
|
||||
fn from(connection: &HdfsConnection) -> Self {
|
||||
let root = util::normalize_dir(&connection.root);
|
||||
HdfsNative::default()
|
||||
.root(&root)
|
||||
.name_node(&connection.name_node)
|
||||
.options(connection.options.clone())
|
||||
}
|
||||
}
|
||||
|
||||
/// Hadoop Distributed File System object storage configuration.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
|
||||
#[serde(default)]
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
pub struct HdfsConfig {
|
||||
pub name: String,
|
||||
#[serde(flatten)]
|
||||
pub connection: HdfsConnection,
|
||||
#[serde(flatten)]
|
||||
pub cache: ObjectStorageCacheConfig,
|
||||
}
|
||||
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
|
||||
#[serde(default)]
|
||||
@@ -411,6 +461,20 @@ mod tests {
|
||||
assert_eq!("test", s3_config.config_name());
|
||||
assert_eq!("S3", s3_config.provider_name());
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
{
|
||||
let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig::default());
|
||||
assert_eq!("Hdfs", hdfs_config.config_name());
|
||||
assert_eq!("Hdfs", hdfs_config.provider_name());
|
||||
|
||||
let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig {
|
||||
name: "test".to_string(),
|
||||
..Default::default()
|
||||
});
|
||||
assert_eq!("test", hdfs_config.config_name());
|
||||
assert_eq!("Hdfs", hdfs_config.provider_name());
|
||||
}
|
||||
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
{
|
||||
let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig::default());
|
||||
@@ -438,6 +502,11 @@ mod tests {
|
||||
assert!(gcs_config.is_object_storage());
|
||||
let azblob_config = ObjectStoreConfig::Azblob(AzblobConfig::default());
|
||||
assert!(azblob_config.is_object_storage());
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
{
|
||||
let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig::default());
|
||||
assert!(hdfs_config.is_object_storage());
|
||||
}
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
{
|
||||
let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig::default());
|
||||
@@ -445,6 +514,46 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
#[test]
|
||||
fn test_hdfs_config_serde() {
|
||||
let config: ObjectStoreConfig = toml::from_str(
|
||||
r#"
|
||||
type = "Hdfs"
|
||||
name = "hdfs-store"
|
||||
root = "/greptimedb"
|
||||
name_node = "hdfs://127.0.0.1:9000"
|
||||
|
||||
[options]
|
||||
"dfs.client.block.write.replace-datanode-on-failure.enable" = "true"
|
||||
"#,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let ObjectStoreConfig::Hdfs(hdfs_config) = config else {
|
||||
unreachable!()
|
||||
};
|
||||
|
||||
assert_eq!("hdfs-store", hdfs_config.name);
|
||||
assert_eq!("/greptimedb", hdfs_config.connection.root);
|
||||
assert_eq!("hdfs://127.0.0.1:9000", hdfs_config.connection.name_node);
|
||||
assert_eq!(
|
||||
Some(&"true".to_string()),
|
||||
hdfs_config
|
||||
.connection
|
||||
.options
|
||||
.get("dfs.client.block.write.replace-datanode-on-failure.enable")
|
||||
);
|
||||
|
||||
let serialized = toml::to_string(&hdfs_config).unwrap();
|
||||
assert!(serialized.contains("name_node = \"hdfs://127.0.0.1:9000\""));
|
||||
assert!(
|
||||
serialized.contains(
|
||||
"\"dfs.client.block.write.replace-datanode-on-failure.enable\" = \"true\""
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
#[test]
|
||||
fn test_mysql_config_connection_string_serde() {
|
||||
|
||||
@@ -15,15 +15,21 @@
|
||||
use std::{fs, path};
|
||||
|
||||
use common_telemetry::info;
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use opendal::services::HdfsNative;
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
use opendal::services::Mysql;
|
||||
use opendal::services::{Fs, Gcs, Oss, S3};
|
||||
use snafu::prelude::*;
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use crate::config::HdfsConfig;
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
use crate::config::MysqlConfig;
|
||||
use crate::config::{AzblobConfig, FileConfig, GcsConfig, ObjectStoreConfig, OssConfig, S3Config};
|
||||
use crate::error::{self, Result};
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
use crate::layers::HdfsCompatibilityLayer;
|
||||
use crate::services::Azblob;
|
||||
use crate::util::{build_http_context, clean_temp_dir, join_dir, normalize_dir};
|
||||
use crate::{ATOMIC_WRITE_DIR, OLD_ATOMIC_WRITE_DIR, ObjectStore, util};
|
||||
@@ -39,11 +45,36 @@ pub async fn new_raw_object_store(
|
||||
ObjectStoreConfig::Oss(oss_config) => new_oss_object_store(oss_config).await,
|
||||
ObjectStoreConfig::Azblob(azblob_config) => new_azblob_object_store(azblob_config).await,
|
||||
ObjectStoreConfig::Gcs(gcs_config) => new_gcs_object_store(gcs_config).await,
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
ObjectStoreConfig::Hdfs(hdfs_config) => new_hdfs_object_store(hdfs_config).await,
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
ObjectStoreConfig::Mysql(mysql_config) => new_mysql_object_store(mysql_config).await,
|
||||
}
|
||||
}
|
||||
|
||||
/// Creates an object store backed by a native HDFS client.
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
pub async fn new_hdfs_object_store(hdfs_config: &HdfsConfig) -> Result<ObjectStore> {
|
||||
let root = util::normalize_dir(&hdfs_config.connection.root);
|
||||
info!(
|
||||
"The HDFS NameNode is: {}, root is: {}",
|
||||
hdfs_config.connection.name_node, root
|
||||
);
|
||||
|
||||
let builder = HdfsNative::from(&hdfs_config.connection);
|
||||
let compatibility_layer = HdfsCompatibilityLayer::new(
|
||||
&hdfs_config.connection.name_node,
|
||||
&hdfs_config.connection.root,
|
||||
&hdfs_config.connection.options,
|
||||
)
|
||||
.context(error::InitBackendSnafu)?;
|
||||
let operator = ObjectStore::new(builder)
|
||||
.context(error::InitBackendSnafu)?
|
||||
.layer(compatibility_layer);
|
||||
|
||||
Ok(operator)
|
||||
}
|
||||
|
||||
#[cfg(feature = "mysql-object-store")]
|
||||
pub async fn new_mysql_object_store(mysql_config: &MysqlConfig) -> Result<ObjectStore> {
|
||||
let root = util::normalize_dir(&mysql_config.root);
|
||||
@@ -144,3 +175,33 @@ pub async fn new_s3_object_store(s3_config: &S3Config) -> Result<ObjectStore> {
|
||||
|
||||
Ok(operator)
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "hdfs-object-store"))]
|
||||
mod tests {
|
||||
use opendal::services::HDFS_NATIVE_SCHEME;
|
||||
|
||||
use super::*;
|
||||
use crate::config::HdfsConnection;
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_new_hdfs_object_store() {
|
||||
let config = HdfsConfig {
|
||||
connection: HdfsConnection {
|
||||
root: "/greptimedb".to_string(),
|
||||
name_node: "hdfs://127.0.0.1:9000".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let store = new_hdfs_object_store(&config).await.unwrap();
|
||||
assert_eq!(HDFS_NATIVE_SCHEME, store.info().scheme());
|
||||
assert_eq!("/greptimedb/", store.info().root());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_new_hdfs_object_store_requires_name_node() {
|
||||
let result = new_hdfs_object_store(&HdfsConfig::default()).await;
|
||||
assert!(result.is_err());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,9 +12,13 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
mod hdfs;
|
||||
#[cfg(feature = "testing")]
|
||||
pub mod mock;
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
pub use hdfs::HdfsCompatibilityLayer;
|
||||
pub use opendal::layers::*;
|
||||
pub use prometheus::build_prometheus_metrics_layer;
|
||||
|
||||
|
||||
@@ -0,0 +1,551 @@
|
||||
// Copyright 2023 Greptime Team
|
||||
//
|
||||
// Licensed 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.
|
||||
|
||||
use std::fmt::{self, Debug};
|
||||
use std::sync::Arc;
|
||||
|
||||
use hdfs_native::{Client, ClientBuilder};
|
||||
use opendal::raw::oio::{Delete as _, Read as _, ReadStream as _, Write as _};
|
||||
use opendal::raw::{
|
||||
Layer, OpCompose, OpCopy, OpCreateDir, OpDelete, OpList, OpPresign, OpRead, OpRename, OpStat,
|
||||
OpWrite, RpCreateDir, RpPresign, RpRename, RpStat, Service, ServiceInfo, Servicer, oio,
|
||||
};
|
||||
use opendal::{Buffer, Capability, ErrorKind, Metadata, OperationContext, Result};
|
||||
use uuid::Uuid;
|
||||
|
||||
/// Adds atomic writes and streaming copies to the native HDFS backend.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct HdfsCompatibilityLayer {
|
||||
renamer: AtomicRenamer,
|
||||
}
|
||||
|
||||
impl HdfsCompatibilityLayer {
|
||||
/// Creates a compatibility layer for the HDFS connection.
|
||||
pub fn new(
|
||||
name_node: &str,
|
||||
root: &str,
|
||||
options: &std::collections::HashMap<String, String>,
|
||||
) -> Result<Self> {
|
||||
let mut config = std::collections::HashMap::new();
|
||||
let namenodes = name_node
|
||||
.split(',')
|
||||
.filter_map(|value| {
|
||||
let value = value
|
||||
.trim()
|
||||
.trim_start_matches("hdfs://")
|
||||
.trim_end_matches('/');
|
||||
(!value.is_empty()).then_some(value)
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
for (index, namenode) in namenodes.iter().enumerate() {
|
||||
config.insert(
|
||||
format!("dfs.namenode.rpc-address.nameservice.nn{index}"),
|
||||
(*namenode).to_string(),
|
||||
);
|
||||
}
|
||||
config.insert(
|
||||
"dfs.ha.namenodes.nameservice".to_string(),
|
||||
(0..namenodes.len())
|
||||
.map(|index| format!("nn{index}"))
|
||||
.collect::<Vec<_>>()
|
||||
.join(","),
|
||||
);
|
||||
config.extend(options.clone());
|
||||
let client = ClientBuilder::new()
|
||||
.with_url("hdfs://nameservice")
|
||||
.with_config(config)
|
||||
.build()
|
||||
.map_err(hdfs_error)?;
|
||||
Ok(Self {
|
||||
renamer: AtomicRenamer::Native {
|
||||
client,
|
||||
root: opendal::raw::normalize_root(root),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
/// Creates a compatibility layer backed by the inner service's rename.
|
||||
#[cfg(any(test, feature = "testing"))]
|
||||
pub fn new_for_test() -> Self {
|
||||
Self {
|
||||
renamer: AtomicRenamer::Raw,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
enum AtomicRenamer {
|
||||
Native {
|
||||
client: Client,
|
||||
root: String,
|
||||
},
|
||||
#[cfg(any(test, feature = "testing"))]
|
||||
Raw,
|
||||
}
|
||||
|
||||
impl AtomicRenamer {
|
||||
async fn rename(
|
||||
&self,
|
||||
_inner: &Servicer,
|
||||
_ctx: &OperationContext,
|
||||
from: &str,
|
||||
to: &str,
|
||||
) -> Result<()> {
|
||||
match self {
|
||||
Self::Native { client, root } => {
|
||||
// OpenDAL's HDFS rename removes an existing destination before
|
||||
// renaming. Use HDFS Rename2 with overwrite to keep replacement atomic.
|
||||
client
|
||||
.rename(
|
||||
&opendal::raw::build_rooted_abs_path(root, from),
|
||||
&opendal::raw::build_rooted_abs_path(root, to),
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.map_err(hdfs_error)
|
||||
}
|
||||
#[cfg(any(test, feature = "testing"))]
|
||||
Self::Raw => _inner
|
||||
.rename(_ctx, from, to, OpRename::new())
|
||||
.await
|
||||
.map(|_| ()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn hdfs_error(error: hdfs_native::HdfsError) -> opendal::Error {
|
||||
opendal::Error::new(ErrorKind::Unexpected, "native HDFS operation failed").set_source(error)
|
||||
}
|
||||
|
||||
impl Layer for HdfsCompatibilityLayer {
|
||||
fn apply_service(&self, inner: Servicer) -> Servicer {
|
||||
Arc::new(HdfsCompatibilityService {
|
||||
inner,
|
||||
renamer: self.renamer.clone(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
struct HdfsCompatibilityService {
|
||||
inner: Servicer,
|
||||
renamer: AtomicRenamer,
|
||||
}
|
||||
|
||||
impl Debug for HdfsCompatibilityService {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
f.debug_struct("HdfsCompatibilityService")
|
||||
.field("inner", &self.inner)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
|
||||
/// A writer that publishes non-append writes with an atomic rename.
|
||||
struct HdfsWriter(HdfsWriterInner);
|
||||
|
||||
enum HdfsWriterInner {
|
||||
Direct(oio::Writer),
|
||||
Atomic {
|
||||
inner: Servicer,
|
||||
context: OperationContext,
|
||||
renamer: AtomicRenamer,
|
||||
writer: Option<oio::Writer>,
|
||||
temporary_path: String,
|
||||
target_path: String,
|
||||
},
|
||||
}
|
||||
|
||||
impl oio::Write for HdfsWriter {
|
||||
async fn write(&mut self, buffer: Buffer) -> Result<()> {
|
||||
match &mut self.0 {
|
||||
HdfsWriterInner::Direct(writer) => writer.write(buffer).await,
|
||||
HdfsWriterInner::Atomic { writer, .. } => {
|
||||
writer
|
||||
.as_mut()
|
||||
.ok_or_else(writer_unavailable)?
|
||||
.write(buffer)
|
||||
.await
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn close(&mut self) -> Result<Metadata> {
|
||||
match &mut self.0 {
|
||||
HdfsWriterInner::Direct(writer) => writer.close().await,
|
||||
HdfsWriterInner::Atomic {
|
||||
inner,
|
||||
context,
|
||||
renamer,
|
||||
writer,
|
||||
temporary_path,
|
||||
target_path,
|
||||
} => {
|
||||
let mut writer = writer.take().ok_or_else(writer_unavailable)?;
|
||||
let metadata = match writer.close().await {
|
||||
Ok(metadata) => metadata,
|
||||
Err(error) => {
|
||||
drop(writer);
|
||||
let _ = delete_path(inner, context, temporary_path).await;
|
||||
return Err(error);
|
||||
}
|
||||
};
|
||||
drop(writer);
|
||||
|
||||
if let Err(error) = renamer
|
||||
.rename(inner, context, temporary_path, target_path)
|
||||
.await
|
||||
{
|
||||
let _ = delete_path(inner, context, temporary_path).await;
|
||||
return Err(error);
|
||||
}
|
||||
|
||||
Ok(metadata)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn abort(&mut self) -> Result<()> {
|
||||
match &mut self.0 {
|
||||
HdfsWriterInner::Direct(writer) => writer.abort().await,
|
||||
HdfsWriterInner::Atomic {
|
||||
inner,
|
||||
context,
|
||||
writer,
|
||||
temporary_path,
|
||||
..
|
||||
} => {
|
||||
let abort_result = if let Some(mut writer) = writer.take() {
|
||||
let result = writer.abort().await;
|
||||
drop(writer);
|
||||
result
|
||||
} else {
|
||||
Ok(())
|
||||
};
|
||||
let cleanup_result = delete_path(inner, context, temporary_path).await;
|
||||
|
||||
match abort_result {
|
||||
Err(error) if error.kind() != ErrorKind::Unsupported => Err(error),
|
||||
_ => cleanup_result,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn writer_unavailable() -> opendal::Error {
|
||||
opendal::Error::new(
|
||||
ErrorKind::Unexpected,
|
||||
"HDFS writer is unavailable after close or abort",
|
||||
)
|
||||
}
|
||||
|
||||
impl Service for HdfsCompatibilityService {
|
||||
type Reader = oio::Reader;
|
||||
type Writer = HdfsWriter;
|
||||
type Lister = oio::Lister;
|
||||
type Deleter = oio::Deleter;
|
||||
type Copier = oio::OneShotCopier;
|
||||
type Composer = oio::Composer;
|
||||
|
||||
fn info(&self) -> ServiceInfo {
|
||||
self.inner.info()
|
||||
}
|
||||
|
||||
fn capability(&self) -> Capability {
|
||||
let mut capability = self.inner.capability();
|
||||
capability.copy = true;
|
||||
capability
|
||||
}
|
||||
|
||||
async fn create_dir(
|
||||
&self,
|
||||
ctx: &OperationContext,
|
||||
path: &str,
|
||||
args: OpCreateDir,
|
||||
) -> Result<RpCreateDir> {
|
||||
self.inner.create_dir(ctx, path, args).await
|
||||
}
|
||||
|
||||
async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
|
||||
self.inner.stat(ctx, path, args).await
|
||||
}
|
||||
|
||||
fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
|
||||
self.inner.read(ctx, path, args)
|
||||
}
|
||||
|
||||
fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
|
||||
if args.append() {
|
||||
return self
|
||||
.inner
|
||||
.write(ctx, path, args)
|
||||
.map(|writer| HdfsWriter(HdfsWriterInner::Direct(writer)));
|
||||
}
|
||||
|
||||
let temporary_path = temporary_path(path);
|
||||
let writer = self.inner.write(ctx, &temporary_path, args)?;
|
||||
Ok(HdfsWriter(HdfsWriterInner::Atomic {
|
||||
inner: Arc::clone(&self.inner),
|
||||
context: ctx.clone(),
|
||||
renamer: self.renamer.clone(),
|
||||
writer: Some(writer),
|
||||
temporary_path,
|
||||
target_path: path.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
fn copy(
|
||||
&self,
|
||||
ctx: &OperationContext,
|
||||
from: &str,
|
||||
to: &str,
|
||||
args: OpCopy,
|
||||
) -> Result<Self::Copier> {
|
||||
if args.if_not_exists() || args.if_match().is_some() {
|
||||
return Err(opendal::Error::new(
|
||||
ErrorKind::Unsupported,
|
||||
"conditional copy is not supported by the HDFS fallback",
|
||||
));
|
||||
}
|
||||
|
||||
let inner = Arc::clone(&self.inner);
|
||||
let context = ctx.clone();
|
||||
let renamer = self.renamer.clone();
|
||||
let from = from.to_string();
|
||||
let to = to.to_string();
|
||||
Ok(oio::OneShotCopier::new(async move {
|
||||
copy_via_read_write(inner, &context, renamer, &from, &to).await
|
||||
}))
|
||||
}
|
||||
|
||||
fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
|
||||
self.inner.delete(ctx)
|
||||
}
|
||||
|
||||
fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
|
||||
self.inner.list(ctx, path, args)
|
||||
}
|
||||
|
||||
fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
|
||||
self.inner.compose(ctx, to, args)
|
||||
}
|
||||
|
||||
async fn rename(
|
||||
&self,
|
||||
ctx: &OperationContext,
|
||||
from: &str,
|
||||
to: &str,
|
||||
args: OpRename,
|
||||
) -> Result<RpRename> {
|
||||
self.inner.rename(ctx, from, to, args).await
|
||||
}
|
||||
|
||||
async fn presign(
|
||||
&self,
|
||||
ctx: &OperationContext,
|
||||
path: &str,
|
||||
args: OpPresign,
|
||||
) -> Result<RpPresign> {
|
||||
self.inner.presign(ctx, path, args).await
|
||||
}
|
||||
}
|
||||
|
||||
async fn copy_via_read_write(
|
||||
inner: Servicer,
|
||||
context: &OperationContext,
|
||||
renamer: AtomicRenamer,
|
||||
source_path: &str,
|
||||
target_path: &str,
|
||||
) -> Result<Metadata> {
|
||||
let reader = inner.read(context, source_path, OpRead::new())?;
|
||||
let (_, mut reader) = reader.open(opendal::BytesRange::from(..)).await?;
|
||||
let temporary_path = temporary_path(target_path);
|
||||
let mut writer = inner.write(context, &temporary_path, OpWrite::new())?;
|
||||
|
||||
loop {
|
||||
let buffer = match reader.read().await {
|
||||
Ok(buffer) => buffer,
|
||||
Err(error) => {
|
||||
abort_and_delete(&inner, context, writer, &temporary_path).await;
|
||||
return Err(error);
|
||||
}
|
||||
};
|
||||
if buffer.is_empty() {
|
||||
break;
|
||||
}
|
||||
if let Err(error) = writer.write(buffer).await {
|
||||
abort_and_delete(&inner, context, writer, &temporary_path).await;
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
|
||||
let metadata = match writer.close().await {
|
||||
Ok(metadata) => metadata,
|
||||
Err(error) => {
|
||||
drop(writer);
|
||||
let _ = delete_path(&inner, context, &temporary_path).await;
|
||||
return Err(error);
|
||||
}
|
||||
};
|
||||
drop(writer);
|
||||
|
||||
if let Err(error) = renamer
|
||||
.rename(&inner, context, &temporary_path, target_path)
|
||||
.await
|
||||
{
|
||||
let _ = delete_path(&inner, context, &temporary_path).await;
|
||||
return Err(error);
|
||||
}
|
||||
|
||||
Ok(metadata)
|
||||
}
|
||||
|
||||
async fn abort_and_delete(
|
||||
inner: &Servicer,
|
||||
context: &OperationContext,
|
||||
mut writer: oio::Writer,
|
||||
path: &str,
|
||||
) {
|
||||
let _ = writer.abort().await;
|
||||
drop(writer);
|
||||
let _ = delete_path(inner, context, path).await;
|
||||
}
|
||||
|
||||
async fn delete_path(inner: &Servicer, context: &OperationContext, path: &str) -> Result<()> {
|
||||
let mut deleter = inner.delete(context)?;
|
||||
deleter.delete(path, OpDelete::new()).await?;
|
||||
deleter.close().await
|
||||
}
|
||||
|
||||
// TODO(fengjiachun): Clean up temporary files left behind after a process crash.
|
||||
fn temporary_path(path: &str) -> String {
|
||||
let suffix = format!(".greptime-{}.tmp", Uuid::new_v4());
|
||||
match path.rsplit_once('/') {
|
||||
Some((parent, name)) => format!("{parent}/.{name}{suffix}"),
|
||||
None => format!(".{path}{suffix}"),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use opendal::services::Fs;
|
||||
use opendal::{Operator, Writer};
|
||||
use tempfile::TempDir;
|
||||
|
||||
use super::*;
|
||||
|
||||
fn test_store() -> (TempDir, Operator) {
|
||||
let directory = tempfile::tempdir().unwrap();
|
||||
let store = Operator::new(Fs::default().root(directory.path().to_str().unwrap()))
|
||||
.unwrap()
|
||||
.layer(HdfsCompatibilityLayer::new_for_test());
|
||||
(directory, store)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_service_operations() {
|
||||
let (_directory, store) = test_store();
|
||||
store.create_dir("data/").await.unwrap();
|
||||
store.write("data/source", "contents").await.unwrap();
|
||||
assert_eq!(8, store.stat("data/source").await.unwrap().content_length());
|
||||
assert_eq!(
|
||||
b"onte",
|
||||
store
|
||||
.read_with("data/source")
|
||||
.range(1..5)
|
||||
.await
|
||||
.unwrap()
|
||||
.to_bytes()
|
||||
.as_ref()
|
||||
);
|
||||
store.rename("data/source", "data/target").await.unwrap();
|
||||
let entries = store.list("data/").await.unwrap();
|
||||
assert!(entries.iter().any(|entry| entry.path() == "data/target"));
|
||||
assert!(!store.exists("data/source").await.unwrap());
|
||||
store.delete("data/target").await.unwrap();
|
||||
assert!(!store.exists("data/target").await.unwrap());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_atomic_write_keeps_old_data_after_abort() {
|
||||
let (_directory, store) = test_store();
|
||||
store.write("manifest.json", "old").await.unwrap();
|
||||
|
||||
let mut writer: Writer = store.writer("manifest.json").await.unwrap();
|
||||
writer.write("new").await.unwrap();
|
||||
writer.abort().await.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
b"old",
|
||||
store
|
||||
.read("manifest.json")
|
||||
.await
|
||||
.unwrap()
|
||||
.to_bytes()
|
||||
.as_ref()
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.list("")
|
||||
.await
|
||||
.unwrap()
|
||||
.iter()
|
||||
.all(|entry| !entry.path().contains(".greptime-"))
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_atomic_write_replaces_on_close() {
|
||||
let (_directory, store) = test_store();
|
||||
store.write("manifest.json", "old").await.unwrap();
|
||||
store.write("manifest.json", "new").await.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
b"new",
|
||||
store
|
||||
.read("manifest.json")
|
||||
.await
|
||||
.unwrap()
|
||||
.to_bytes()
|
||||
.as_ref()
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.list("")
|
||||
.await
|
||||
.unwrap()
|
||||
.iter()
|
||||
.all(|entry| !entry.path().contains(".greptime-"))
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_copy_fallback_streams_to_target() {
|
||||
let (_directory, store) = test_store();
|
||||
store.write("source.parquet", "contents").await.unwrap();
|
||||
store
|
||||
.copy("source.parquet", "target.parquet")
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
b"contents",
|
||||
store
|
||||
.read("target.parquet")
|
||||
.await
|
||||
.unwrap()
|
||||
.to_bytes()
|
||||
.as_ref()
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -30,6 +30,8 @@ pub mod secure_fs;
|
||||
pub mod test_util;
|
||||
pub mod util;
|
||||
|
||||
#[cfg(feature = "hdfs-object-store")]
|
||||
pub use config::HdfsConnection;
|
||||
pub use config::{AzblobConnection, GcsConnection, OssConnection, S3Connection};
|
||||
|
||||
/// The default object cache directory name.
|
||||
|
||||
Reference in New Issue
Block a user