如何正确使用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
相关产品推荐
相关产品推荐

