feat(services/hdfs): add copy support - #8332
Conversation
Enable HDFS copy via OneShotCopier by streaming bytes through hdrs AsyncFile APIs, creating parent directories and overwriting existing destination files when needed.
Bump hdrs to 0.3.3 and replace the AsyncFile streaming copy with native Client::copy_file.
Move hdrs copy_file onto spawn_blocking, drop the pre-delete now that hdfsCopy overwrites, and use OneShotCopier::new_with so temporary IO errors can be retried.
Restore remove_file for existing copy destinations and document that hdfsCopy has been verified to overwrite natively.
eaaa7d9 to
0327d2f
Compare
|
Hi, @Xuanwo @erickguan . Have self-reviewed by different coding agents. Could you please review another time when have free time? Thanks very much!!! |
erickguan
left a comment
There was a problem hiding this comment.
Please try to understand your code or your agentic code pipeline. Above all, what you want to achieve for the project besides code.
Supporting copy with hdfs, solid use case. But there are some issues I don't know what you want to achieve:
- Do you want copy atomicity? Does HDFS support it?
- Linking another crate's code for explanation is okay. Though I suggest documenting critical assumptions and decisions.
|
|
||
| pub async fn hdfs_copy(&self, from: &str, to: &str) -> Result<Metadata> { | ||
| let from_path = build_rooted_abs_path(&self.root, from); | ||
| // FileUtil.copy recurses when the source is a directory. |
There was a problem hiding this comment.
OpenDAL defines this operation as file-to-file copy. Hadoop FileUtil.copy also
accepts directory sources and recursively copies their contents, which would
produce behavior outside OpenDAL's copy contract.
I will rewrite the comment to describe the OpenDAL constraint directly:
// OpenDAL copy is file-to-file only. Reject directory sources before
// Hadoop can recursively copy their contents.
| let to_path = build_rooted_abs_path(&self.root, to); | ||
| match self.client.metadata(&to_path) { | ||
| Ok(meta) => { | ||
| // FileUtil.checkDest rewrites a directory destination to dst/<srcName> |
There was a problem hiding this comment.
Why do we mention FileUtil.checkDest here?
There was a problem hiding this comment.
We do not need to mention the internal FileUtil.checkDest implementation here.
The intended constraint is that OpenDAL treats to as the exact destination
file path. If to is an existing directory, Hadoop may copy the source under
that directory instead, while OpenDAL should return IsADirectory.
I will rewrite the comment as:
// OpenDAL treats to as the exact destination file path. Reject an
// existing directory instead of copying the source into it.
| // hdfsCopy has been verified to overwrite natively via | ||
| // FileUtil.copy(..., overwrite=true). | ||
| self.client | ||
| .remove_file(&to_path) | ||
| .map_err(new_std_io_error)?; | ||
| } |
There was a problem hiding this comment.
This comment really contradicts the code.
There was a problem hiding this comment.
You are right. The comment is incorrect and contradicts the implementation.
The hdrs copy_file binding calls hdfsCopy, and this path does not pass an
overwrite=true argument. The existing destination is removed explicitly to
implement OpenDAL's default overwrite semantics.
I will remove the incorrect claim and document the actual behavior:
// hdfsCopy does not replace an existing destination, so remove the
// destination first to implement OpenDAL's overwrite semantics.
//
// This replacement is not atomic. If the copy fails after removal, the
// destination may be missing or incomplete.
Thanks for the guidance. The motivating use case is Lance's opt-in native copy path for immutable data Therefore, this copy operation does not require atomicity. HDFS hdfsCopy / The existing-destination removal is needed to implement OpenDAL's default copy I will also rewrite the comments around directory handling in terms of OpenDAL's |
Thanks for your very valuable suggestions! |
|
Happy to help. |
Which issue does this PR close?
Closes #.
Rationale for this change
HDFS currently reports
copyas unsupported. Enabling nativeOperator::copyfor the HDFS service makes same-cluster file duplication consistent with other filesystem-like backends such asfs.hdrs0.3.3 exposesClient::copy_file, which wraps libhdfshdfsCopy/ HadoopFileUtil.copywithoverwrite=true. This PR uses that API instead of streaming bytes throughhdrs::AsyncFile.What changes are included in this PR?
hdrsfrom 0.3.2 to 0.3.3, which also pullshdfs-sys0.3.0 → 0.3.1 (JNI attach/detach tracking,cargo:rustc-link-arg=-Wl,-rpathon non-Windows, and dropping bindgenstatic_flag(true))copycapabilityhdfs_copy()that:dst/<srcName>)copy_filerequires the destination parent to exist)hdfsCopyhas also been verified to overwrite nativelycopy_fileonspawn_blockingso the data-path JNI copy does not block the async workerService::copytooio::OneShotCopier::new_withso temporary IO errors remain retryablecopyas supported in service docsAre there any user-facing changes?
Yes. HDFS operators that previously got
Unsupportedfromcopycan now copy files. Capability discovery will reportcopy: true.Breaking changes
AI Usage Statement
hdrs0.3.3copy_file, and applied review follow-ups. Behavior tests against a live HDFS cluster were not run on this machine.