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

使用多线程向GCS上传大量CSV文件时持续报错并中断

问题

在同一GCP环境的Java App Engine应用中,需要向GCS上传36000+个CSV文件。最初采用parallelStream.forEach实现并行上传,核心代码如下:

void createFiles()
{
    try {
        Map<Integer, FileDetail> collection = GetFiles();
        collection.values().entrySet().parallelStream().forEach((entry) -> {
            try{
                createFileInGCS();
            } catch (Exception e) {
                LOG.error("Error occurred while creating FilesInGCS : " + e.getMessage());
                throw new RuntimeException(e);
            }
        });
    }
    catch (Exception e) {
        LOG.error("Error occurred while creating files : "+ e.getMessage());
        throw new CustomException("Unable to process files");
    }
}


public void createFileInGCS()
{
    try {
        String fileName ="file_" + fileId + ".csv";
        String gcsBucket = config.getGcsBucket();
        var blobId = BlobId.of(gcsBucket, fileName);
        var blobInfo = BlobInfo.newBuilder(blobId).setContentType("text/plain").build();
        var storage = StorageOptions.newBuilder().setProjectId(config.getProject())
                .build().getService();
        storage.create(blobInfo, csvContent.getBytes());
    }
    catch(Exception e){
        LOG.error("Storage failed with exception :"+e.getMessage());
    }
}

当文件数量增加后,应用出现故障:

  • 请求初期频繁报错:ERROR: onFailure exception: com.google.cloud.logging.LoggingException: io.grpc.StatusRuntimeException: INTERNAL: Panic! This is a bug!
  • 后续出现Broken pipe、Remote host terminated the handshake错误并导致应用终止

尝试将文件拆分为10批(每批约3000个),前9批4分钟内完成,但最后一批运行超1小时仍失败。曾尝试GCS的StorageBatch,但该工具没有创建文件的方法,询问是否有更优的批量上传方案。

优化方案

1. 复用Storage客户端实例

当前代码每次上传都创建新的Storage实例,会导致大量连接资源浪费,是引发连接类错误的核心原因。需改为全局单例复用:

// 全局单例Storage实例
private final Storage storage;

// 在类初始化时创建一次
public YourService(Config config) {
    this.storage = StorageOptions.newBuilder()
            .setProjectId(config.getProject())
            .build().getService();
}

public void createFileInGCS() {
    try {
        String fileName ="file_" + fileId + ".csv";
        String gcsBucket = config.getGcsBucket();
        var blobId = BlobId.of(gcsBucket, fileName);
        var blobInfo = BlobInfo.newBuilder(blobId).setContentType("text/plain").build();
        // 复用全局实例
        storage.create(blobInfo, csvContent.getBytes());
    } catch(Exception e){
        LOG.error("Storage failed with exception :"+e.getMessage());
    }
}

2. 手动控制并发度

parallelStream的并发度由JVM自动管理,对于IO密集型的GCS上传,过高并发会触发限流或耗尽资源。建议用自定义线程池控制并发数(比如50-100,根据实例规格调整):

private final ExecutorService uploadExecutor = Executors.newFixedThreadPool(50);

void createFiles() {
    try {
        Map<Integer, FileDetail> collection = GetFiles();
        List<CompletableFuture<Void>> futures = collection.values().entrySet().stream()
                .map(entry -> CompletableFuture.runAsync(() -> {
                    try {
                        createFileInGCS();
                    } catch (Exception e) {
                        LOG.error("Error creating file: " + e.getMessage());
                        throw new CompletionException(e);
                    }
                }, uploadExecutor))
                .collect(Collectors.toList());
        
        // 等待所有上传任务完成
        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
    } catch (Exception e) {
        LOG.error("Batch upload failed: " + e.getMessage());
        throw new CustomException("Unable to process files");
    } finally {
        uploadExecutor.shutdown();
    }
}

3. 增加重试机制处理临时错误

Broken pipe、Remote host terminated the handshake属于临时网络错误,通过指数退避重试可解决大部分此类问题。可以启用GCP客户端内置的重试策略:

public YourService(Config config) {
    RetrySettings retrySettings = RetrySettings.newBuilder()
            .setMaxAttempts(5)
            .setInitialRetryDelay(Duration.ofMillis(100))
            .setRetryDelayMultiplier(2.0)
            .setMaxRetryDelay(Duration.ofSeconds(10))
            .build();

    this.storage = StorageOptions.newBuilder()
            .setProjectId(config.getProject())
            .setRetrySettings(retrySettings)
            .build().getService();
}

4. 批量上传替代方案

如果单文件并发上传仍有瓶颈,可考虑以下方案:

  • 打包后上传再解压:将多个CSV打包成ZIP上传到GCS,再通过Cloud Functions或Cloud Run自动解压到目标路径(适合小文件场景)。
  • 使用GCS Transfer Service:如果文件来源是本地或其他云存储,可通过Transfer Service批量迁移,无需在App Engine内处理。

5. 调整App Engine实例配置

  • 升级实例规格:使用F4_1G或更高配置的实例,确保有足够的CPU和内存处理并发任务。
  • 检查配额:确认App Engine的出站连接数、CPU等配额未达上限,必要时申请调整。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 06:44:53