Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions core/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion core/services/hdfs/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ opendal-core = { path = "../../core", version = "0.59.3", default-features = fal

bytes = { workspace = true }
futures = { workspace = true }
hdrs = { version = "0.3.2", features = ["async_file"] }
hdrs = { version = "0.3.3", features = ["async_file"] }
log = { workspace = true }
serde = { workspace = true, features = ["derive"] }
tokio = { workspace = true, features = ["rt"] }
22 changes: 15 additions & 7 deletions core/services/hdfs/src/backend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,8 @@ impl Builder for HdfsBuilder {

list: true,

copy: true,

rename: true,
rename_with_if_not_exists: true,

Expand All @@ -195,7 +197,7 @@ impl Service for HdfsBackend {
type Writer = HdfsLazyWriter;
type Lister = Option<HdfsLister>;
type Deleter = oio::OneShotDeleter<HdfsDeleter>;
type Copier = ();
type Copier = oio::OneShotCopier;
type Composer = ();

fn info(&self) -> ServiceInfo {
Expand Down Expand Up @@ -258,14 +260,20 @@ impl Service for HdfsBackend {
fn copy(
&self,
_ctx: &OperationContext,
_from: &str,
_to: &str,
from: &str,
to: &str,
_args: OpCopy,
) -> Result<Self::Copier> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
let core = self.core.clone();
let from = from.to_string();
let to = to.to_string();
// Recreate the future so retry layers can rerun the copy after a temporary error.
Ok(oio::OneShotCopier::new_with(move || {
Comment thread
erickguan marked this conversation as resolved.
let core = core.clone();
let from = from.clone();
let to = to.clone();
async move { core.hdfs_copy(&from, &to).await }
}))
}

async fn rename(
Expand Down
71 changes: 55 additions & 16 deletions core/services/hdfs/src/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -124,8 +124,7 @@ impl HdfsCore {
});

if !target_exists {
let parent = get_parent(&target_path);
self.client.create_dir(parent).map_err(new_std_io_error)?;
self.ensure_parent_dir(&target_path)?;
}
if !should_append {
initial_size = 0;
Expand Down Expand Up @@ -175,20 +174,7 @@ impl HdfsCore {
return Err(new_std_io_error(err));
}

let parent = std::path::PathBuf::from(&to_path)
.parent()
.ok_or_else(|| {
Error::new(
ErrorKind::Unexpected,
"path should have parent but not, it must be malformed",
)
.with_context("to", &to_path)
})?
.to_path_buf();

self.client
.create_dir(&parent.to_string_lossy())
.map_err(new_std_io_error)?;
self.ensure_parent_dir(&to_path)?;
}
Ok(metadata) => {
if metadata.is_file() {
Expand All @@ -215,6 +201,59 @@ impl HdfsCore {

Ok(())
}

pub async fn hdfs_copy(&self, from: &str, to: &str) -> Result<Metadata> {
let from_path = build_rooted_abs_path(&self.root, from);
// OpenDAL copy is file-to-file only. Reject directory sources before
// the HDFS API can recursively copy their contents.
let from_meta = self.client.metadata(&from_path).map_err(new_std_io_error)?;
if !from_meta.is_file() {
return Err(
Error::new(ErrorKind::IsADirectory, "from path should be a file")
.with_context("from", &from_path),
);
}

let to_path = build_rooted_abs_path(&self.root, to);
match self.client.metadata(&to_path) {
Ok(meta) => {
// OpenDAL treats `to` as the exact destination file path. Reject
// an existing directory instead of copying the source into it.
if meta.is_dir() {
return Err(
Error::new(ErrorKind::IsADirectory, "to path should be a file")
.with_context("to", &to_path),
);
}
// The HDFS copy API does not replace an existing destination,
// so remove it first to preserve OpenDAL's overwrite semantics.
self.client
.remove_file(&to_path)
.map_err(new_std_io_error)?;
}
Err(err) if err.kind() == io::ErrorKind::NotFound => {
// hdrs copy_file requires the destination parent to already exist.
self.ensure_parent_dir(&to_path)?;
}
Err(err) => return Err(new_std_io_error(err)),
}

let client = self.client.clone();
let copy_from = from_path.clone();
let copy_to = to_path.clone();
tokio::task::spawn_blocking(move || client.copy_file(&copy_from, &copy_to))
.await
.map_err(|e| Error::new(ErrorKind::Unexpected, "tokio task join failed").set_source(e))?
.map_err(new_std_io_error)?;

Ok(MetadataBuilder::file(from_meta.len()).build())
}

fn ensure_parent_dir(&self, path: &str) -> Result<()> {
self.client
.create_dir(get_parent(path))
.map_err(new_std_io_error)
}
}

#[cfg(test)]
Expand Down
10 changes: 9 additions & 1 deletion core/services/hdfs/src/docs.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,13 +10,21 @@ Depending on its configuration and the backing system, this service can expose:
- [x] write
- [x] delete
- [x] list
- [ ] copy
- [x] copy
- [x] rename
- [ ] ~~presign~~

Inspect the effective capability set with [`opendal_core::Operator::info`] and
[`opendal_core::OperatorInfo::capability`] after building an operator.

## Copy behavior

The HDFS service copies files to exact destination file paths and rejects
directory sources and destinations. Copy overwrites an existing destination by
removing it before copying the source. This replacement is not atomic: if the
copy fails after removing the destination, the destination can be missing or
incomplete.

## Differences with webhdfs

The [WebHDFS service](https://docs.rs/opendal-service-webhdfs) uses HDFS's
Expand Down
Loading