You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确使用org.apache.hadoop.fs.MultipartUploader实现大文件流式上传?

Hadoop MultipartUploader 大文件流式上传实现指南

问题背景

需要实现无缓冲流式上传1GB+大文件到用作S3的Hadoop集群,org.apache.hadoop.fs.MultipartUploader接口匹配需求,但官方仅给出start/putPart/complete的步骤框架,缺少具体使用示例与参数说明。

已完成的初始化步骤

已成功创建FileSystemMultipartUploader实例并完成上传初始化:

FileSystemMultipartUploader mu = 
    (FileSystemMultipartUploader) fs.createMultipartUploader(new Path("base/path/for/uploading/to/hdfs")); 
Path targetPath = new Path("base/path/plus/result/filename");

// 初始化上传并获取UploadHandle
UploadHandle uh = mu.startUpload(targetPath).get();

核心疑问

调用putPart方法时,对以下参数的取值逻辑不明确:

  • partNumber:分片的序号规则
  • inputStream:分片的输入流如何生成
  • lengthInBytes:分片的字节长度如何确定

原预期接口会自动处理文件分片,但实际需要上层自行实现分片逻辑。

参数解释与完整实现示例

参数说明

  • partNumber:分片的唯一序号,从1开始递增,最大支持10000个分片;不同分片可并行上传,序号无需按顺序传递。
  • inputStream:对应单个分片的输入流,需自行从源文件中截取对应分段生成。
  • lengthInBytes:当前分片的准确字节数,必须与输入流中实际读取的字节数一致。

完整代码示例

以下是实现大文件流式分片上传的完整逻辑,包含分片读取、并行上传、完成上传的步骤:

import org.apache.hadoop.fs.*;
import java.io.*;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;

public class HdfsMultipartUploadExample {
    private static final int PART_SIZE = 1024 * 1024 * 100; // 100MB分片大小,可根据集群配置调整

    public static void uploadLargeFile(FileSystem fs, Path sourceLocalPath, Path targetHdfsPath) throws Exception {
        // 创建MultipartUploader实例
        FileSystemMultipartUploader uploader = 
            (FileSystemMultipartUploader) fs.createMultipartUploader(targetHdfsPath.getParent());
        
        // 初始化上传
        UploadHandle uploadHandle = uploader.startUpload(targetHdfsPath).get();
        List<CompletableFuture<PartHandle>> partFutures = new ArrayList<>();

        try (FileInputStream fis = new FileInputStream(sourceLocalPath.toUri().getPath())) {
            byte[] buffer = new byte[PART_SIZE];
            int partNumber = 1;
            int bytesRead;

            // 分片读取源文件并提交上传
            while ((bytesRead = fis.read(buffer)) != -1) {
                // 为当前分片创建输入流(接口实现会负责关闭流)
                InputStream partStream = new ByteArrayInputStream(buffer, 0, bytesRead);
                
                // 异步提交分片上传
                CompletableFuture<PartHandle> partFuture = uploader.putPart(
                    uploadHandle,
                    partNumber++,
                    targetHdfsPath,
                    partStream,
                    bytesRead
                );
                partFutures.add(partFuture);
            }

            // 等待所有分片上传完成并收集PartHandle
            List<PartHandle> partHandles = new ArrayList<>();
            for (CompletableFuture<PartHandle> future : partFutures) {
                partHandles.add(future.get());
            }

            // 完成多段上传
            uploader.completeUpload(uploadHandle, targetHdfsPath, partHandles).get();
            System.out.println("大文件上传完成");
        } catch (Exception e) {
            // 上传失败时取消上传,清理临时文件
            uploader.abortUpload(uploadHandle, targetHdfsPath).get();
            throw new RuntimeException("上传失败,已取消多段上传任务", e);
        }
    }
}

关键注意事项

  • 分片大小建议设置为64MB到5GB之间,Hadoop默认支持的最小分片为5MB,最大为5GB。
  • 分片上传支持并行执行,可根据集群资源调整并发数。
  • 必须确保所有分片上传完成后再调用completeUpload,若上传失败需调用abortUpload清理临时文件。

内容的提问来源于stack exchange,提问作者Michael

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 05:22:03