BlobClient.setHttpHeaders()执行后upload代码未执行的问题排查
批量上传Byte[]到Azure存储时upload语句未执行的问题
我尝试用以下代码批量将byte[]格式的文件上传到Azure存储账户,设置HttpHeaders的操作看似正常,但后续的upload语句从未执行。在IntelliJ中给setHttpHeaders()和upload()行都加了断点,程序会在setHttpHeaders()处停下,但步进后直接进入下一个任务的添加流程。
原代码
public void uploadFilesInBatch(Map<Document, byte[]> documents) throws InterruptedException, IOException { // int numThreads = Runtime.getRuntime().availableProcessors(); int numThreads = 6; ExecutorService executorService = Executors.newFixedThreadPool(numThreads); List<Runnable> tasks = new ArrayList<>(); TikaConfig config = TikaConfig.getDefaultConfig(); for (Map.Entry<Document, byte[]> entry : documents.entrySet()) { MediaType mediaType = config.getMimeRepository().detect(new ByteArrayInputStream(entry.getValue()), new Metadata()); tasks.add(() -> { String fileName = entry.getKey().getFileName(); BlobClient blobClient = blobContainerClient.getBlobClient(fileName); // synchronized (blobClient) { blobClient.setHttpHeaders(new BlobHttpHeaders().setContentType(mediaType.getType())); // blobClient.setHttpHeaders(new BlobHttpHeaders().setContentType(mediaType.toString())); blobClient.upload(new ByteArrayInputStream(entry.getValue())); // } }); } for (Runnable task : tasks) { executorService.submit(task); } executorService.shutdown(); executorService.awaitTermination(Long.MAX_VALUE, TimeUnit.SECONDS); }
环境信息
- Java版本:11
- 相关依赖版本:
<azure-core.version>1.36.0</azure-core.version> <spring-cloud-azure-starter-keyvault-secrets.version>4.3.0</spring-cloud-azure-starter-keyvault-secrets.version> <!-- Version 5.0.0 requires Java major version 61 instead of 55 --> <azure-storage-blob.version>12.21.0</azure-storage-blob.version> <azure-storage-blob-batch.version>12.18.0</azure-storage-blob-batch.version>
- Azurite:运行在Docker容器,镜像版本
mcr.microsoft.com/azure-storage/azurite:3.23.0 - 调用上下文:
uploadFilesInBatch()方法在CompletableFuture.supplyAsync()的实现中被调用
问题原因及解决方案
核心问题1:setHttpHeaders调用时机错误
当目标Blob尚未创建时,调用BlobClient.setHttpHeaders()会抛出BlobNotFoundException(运行时异常),任务因未捕获异常直接终止,后续的upload代码根本没有执行机会。
核心问题2:Lambda循环变量捕获的闭包问题
在for循环中直接捕获entry变量,Lambda实际引用的是循环变量的同一个实例,当任务异步执行时,可能出现多个任务引用同一entry的情况,导致逻辑混乱。
修正后的代码
public void uploadFilesInBatch(Map<Document, byte[]> documents) throws InterruptedException, IOException { int numThreads = 6; ExecutorService executorService = Executors.newFixedThreadPool(numThreads); TikaConfig config = TikaConfig.getDefaultConfig(); for (Map.Entry<Document, byte[]> entry : documents.entrySet()) { // 提取局部变量,避免Lambda闭包引用循环变量 Document doc = entry.getKey(); byte[] fileBytes = entry.getValue(); MediaType mediaType = config.getMimeRepository().detect(new ByteArrayInputStream(fileBytes), new Metadata()); executorService.submit(() -> { String fileName = doc.getFileName(); BlobClient blobClient = blobContainerClient.getBlobClient(fileName); try { // 合并上传与Header设置,一次请求完成,避免先设置Header的错误 BlobUploadOptions uploadOptions = new BlobUploadOptions(new ByteArrayInputStream(fileBytes)) .setHeaders(new BlobHttpHeaders().setContentType(mediaType.toString())); blobClient.uploadWithResponse(uploadOptions, null, Context.NONE); } catch (Exception e) { // 添加异常捕获,避免任务静默失败 System.err.printf("上传文件%s失败: %s%n", fileName, e.getMessage()); e.printStackTrace(); } }); } executorService.shutdown(); executorService.awaitTermination(Long.MAX_VALUE, TimeUnit.SECONDS); }
关键优化点
- 合并上传与Header设置:使用
BlobUploadOptions在上传时直接设置ContentType,通过uploadWithResponse一次HTTP请求完成Blob创建与Header配置,既避免了先设置Header的错误,也提升了效率。 - 修复闭包引用问题:在循环内将
entry的key、value和mediaType提取为局部变量,确保每个Lambda捕获的是当前迭代的独立变量。 - 添加异常处理:在任务内部捕获所有异常并输出日志,避免任务因异常静默终止,方便排查问题。
内容的提问来源于stack exchange,提问作者du-it
相关产品推荐
相关产品推荐

