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

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);
}

关键优化点

  1. 合并上传与Header设置:使用BlobUploadOptions在上传时直接设置ContentType,通过uploadWithResponse一次HTTP请求完成Blob创建与Header配置,既避免了先设置Header的错误,也提升了效率。
  2. 修复闭包引用问题:在循环内将entry的key、value和mediaType提取为局部变量,确保每个Lambda捕获的是当前迭代的独立变量。
  3. 添加异常处理:在任务内部捕获所有异常并输出日志,避免任务因异常静默终止,方便排查问题。

内容的提问来源于stack exchange,提问作者du-it

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 16:24:54