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

批量上传文件至GCS偶发ECONNRESET及Cloud Tasks调用冻结问题

批量上传GCS偶发ECONNRESET异常及Cloud Tasks调用无响应的解决方案

问题背景

使用@google-cloud/storage v6.9.5批量上传文件至Google Cloud Storage时,偶发无规律的ECONNRESET异常。场景为:30个进程同时运行,每个进程调用约57次uploadFileToGCS函数,上传文件包含35张20KB的jpg、1张500KB的tiff、7个150KB的json。已尝试修改流配置、设置超时等方案无效,确认总数据量未超50GB,云监控无带宽数据异常。

错误信息

error - unhandledRejection: FetchError: request to https://storage.googleapis.com/upload/storage/v1/b/satellite-photos/o?uploadType=multipart&name=64d162a80103b55f87a6cdb4_64db59204f7cd1f34a65e5c7_2023_6_31_raw.tiff failed, reason: read ECONNRESET
    at ClientRequest.<anonymous> (D:\WORK\nirby-project\node_modules\next\dist\compiled\node-fetch\index.js:1:65756)
    at ClientRequest.emit (node:events:525:35)
    at TLSSocket.socketErrorListener (node:_http_client:496:9)
    at TLSSocket.emit (node:events:513:28)
    at emitErrorNT (node:internal/streams/destroy:151:8)
    at emitErrorCloseNT (node:internal/streams/destroy:116:3)
    at process.processTicksAndRejections (node:internal/process/task_queues:82:21) {
  type: 'system',
  errno: 'ECONNRESET',
  code: 'ECONNRESET'
}

相关代码

saveLayers.ts

import { saveSinglePhoto } from "~/helpers/sentinel/saveSinglePhoto";

interface INdviData {
   min: number;
   max: number;
   averageNDVI: number;
   buffer: Buffer;
}

const saveLayers = async (
   fileNameBase: string,
   trueColorData: Buffer,
   ndviData: INdviData,
   contrastNdviData: Buffer,
   sclData: Buffer,
   clmData: Buffer
) => {
   console.log(fileNameBase);

   const data = await Promise.all([
       saveSinglePhoto(trueColorData, fileNameBase, "TRUECOLOR"),
       saveSinglePhoto(ndviData.buffer, fileNameBase, "NDVI"),
       saveSinglePhoto(contrastNdviData, fileNameBase, "CONTRAST_NDVI"),
       saveSinglePhoto(sclData, fileNameBase, "SCL"),
       saveSinglePhoto(clmData, fileNameBase, "CLM"),
   ]);

   return data;
};

export default saveLayers;

saveSinglePhoto.ts

import { SatelliteData } from "~/interfaces/sentinel/SatelliteData";
import uploadFileToGCS from "../gcs/uploadFileToGCS";

const baseUrl = `https://storage.googleapis.com/${process.env.GCLOUD_STORAGE_BUCKET}/`;

export const saveSinglePhoto = async (
    buffer: Buffer,
    fileNameBase: string,
    layerId: string
): Promise<SatelliteData | null> => {
    const fileName = fileNameBase + layerId + ".jpg";

    await uploadFileToGCS(fileName, buffer, "image/jpeg");

    const satelliteData: SatelliteData = {
        layerId: layerId,
        fileUrl: baseUrl + fileName,
    };

    return satelliteData;
};

uploadFileToGCS.ts

import { Storage } from "@google-cloud/storage";

const gcsKey = JSON.parse(Buffer.from(process.env.GCLOUD_CRED_FILE, "base64").toString());

const storage = new Storage({
    credentials: {
        client_email: gcsKey.client_email,
        private_key: gcsKey.private_key,
    },
    projectId: process.env.GCLOUD_PROJECT_ID,
});

const uploadFileToGCS = (filename: string, data: any, contentType: string) => {
    return new Promise((resolve, reject) => {
        const file = storage.bucket(process.env.GCLOUD_STORAGE_BUCKET).file(filename);

        const stream = file.createWriteStream({
            metadata: {
                contentType,
            },
            resumable: false,
            validation: false,
            timeout: 86400,
        });

        stream.on("error", (err) => {
            reject(err);
        });

        stream.on("finish", () => {
            resolve("ok");
        });

        stream.end(data);
    });
};

export default uploadFileToGCS;

补充:Cloud Tasks尝试及问题

尝试用Cloud Tasks异步处理上传,但调用createTask时服务器卡住无响应,程序冻结,相关代码:

import { v2beta3 } from "@google-cloud/tasks";
import { google } from "@google-cloud/tasks/build/protos/protos";

const client = new v2beta3.CloudTasksClient();

const project = process.env.GCLOUD_PROJECT_ID;
const location = process.env.GCLOUD_QUEUE_LOCATION;
const queue = process.env.GCLOUD_QUEUE_NAME;

const parent = client.queuePath(project, location, queue);

const gcsKey = JSON.parse(Buffer.from(process.env.GCLOUD_CRED_FILE, "base64").toString());
const email = gcsKey.client_email;

const createUploadTask = async (filename: string, bucket: string, contentType: string, data: any) => {
    const payload = {
        filename,
        bucket,
        contentType,
        data,
    };

    const url = `${process.env.GCLOUD_FUNCTIONS_URL}/uploadToCS`;

    const task = {
        httpRequest: {
            httpMethod: google.cloud.tasks.v2beta3.HttpMethod.POST,
            url,
            oidcToken: {
                serviceAccountEmail: email,
                audience: url,
            },
            headers: {
                "Content-Type": "application/json",
            },
            body: Buffer.from(JSON.stringify(payload)).toString("base64"),
        },
    };

    const [response] = await client.createTask({ parent, task });
    const name = response.name ?? "name not exists";

    return name;
};

export default createUploadTask;

可行解决方案

1. 为GCS添加上传重试机制

ECONNRESET属于网络瞬态错误,添加重试逻辑可有效解决偶发失败。使用p-retry库实现重试:

修改uploadFileToGCS.ts:

import { Storage } from "@google-cloud/storage";
import pRetry from "p-retry"; // 先执行安装:npm install p-retry

const gcsKey = JSON.parse(Buffer.from(process.env.GCLOUD_CRED_FILE, "base64").toString());

const storage = new Storage({
    credentials: {
        client_email: gcsKey.client_email,
        private_key: gcsKey.private_key,
    },
    projectId: process.env.GCLOUD_PROJECT_ID,
});

const uploadFileToGCS = async (filename: string, data: any, contentType: string) => {
    return pRetry(async () => {
        return new Promise((resolve, reject) => {
            const file = storage.bucket(process.env.GCLOUD_STORAGE_BUCKET).file(filename);
            const stream = file.createWriteStream({
                metadata: { contentType },
                resumable: false,
                validation: false,
                timeout: 30000, // 缩短超时配合重试
            });

            stream.on("error", (err) => {
                // 标记网络类错误为可重试
                if (err.code === "ECONNRESET" || err.type === "system") {
                    throw err;
                }
                reject(err);
            });

            stream.on("finish", () => resolve("ok"));
            stream.end(data);
        });
    }, {
        retries: 3, // 最多重试3次
        minTimeout: 1000, // 初始间隔1秒
        factor: 2, // 指数退避
    });
};

export default uploadFileToGCS;

2. 控制上传并发数

当前每个进程通过Promise.all同时发起5个上传请求,30个进程合计150并发,可能超出GCS的并发限制。使用p-limit控制单进程并发数:

修改saveLayers.ts:

import { saveSinglePhoto } from "~/helpers/sentinel/saveSinglePhoto";
import pLimit from "p-limit"; // 先执行安装:npm install p-limit

interface INdviData {
   min: number;
   max: number;
   averageNDVI: number;
   buffer: Buffer;
}

const limit = pLimit(2); // 单进程最多2个并发上传

const saveLayers = async (
   fileNameBase: string,
   trueColorData: Buffer,
   ndviData: INdviData,
   contrastNdviData: Buffer,
   sclData: Buffer,
   clmData: Buffer
) => {
   console.log(fileNameBase);

   const tasks = [
       limit(() => saveSinglePhoto(trueColorData, fileNameBase, "TRUECOLOR")),
       limit(() => saveSinglePhoto(ndviData.buffer, fileNameBase, "NDVI")),
       limit(() => saveSinglePhoto(contrastNdviData, fileNameBase, "CONTRAST_NDVI")),
       limit(() => saveSinglePhoto(sclData, fileNameBase, "SCL")),
       limit(() => saveSinglePhoto(clmData, fileNameBase, "CLM")),
   ];

   const data = await Promise.all(tasks);
   return data;
};

export default saveLayers;

3. 修复Cloud Tasks调用问题

Cloud Tasks卡住的核心原因是payload过大:将文件buffer放入任务payload会超过Cloud Tasks的100KB payload限制,导致请求阻塞。此外需检查权限配置:

修复步骤:

  1. 移除payload中的buffer数据:改为在Cloud Functions中从源头获取文件数据(比如生成buffer的服务接口、临时存储),仅传递文件标识信息。
  2. 检查权限:确保服务账号拥有roles/cloudtasks.enqueuer角色,队列已正确创建。
  3. 验证网络:本地运行时,通过gcloud tasks queues describe命令测试是否能访问Cloud Tasks API。

修改后的createUploadTask示例:

import { v2beta3 } from "@google-cloud/tasks";
import { google } from "@google-cloud/tasks/build/protos/protos";

const client = new v2beta3.CloudTasksClient();

const project = process.env.GCLOUD_PROJECT_ID;
const location = process.env.GCLOUD_QUEUE_LOCATION;
const queue = process.env.GCLOUD_QUEUE_NAME;

const parent = client.queuePath(project, location, queue);

const gcsKey = JSON.parse(Buffer.from(process.env.GCLOUD_CRED_FILE, "base64").toString());
const email = gcsKey.client_email;

// 仅传递文件元数据,不传递buffer
const createUploadTask = async (filename: string, bucket: string, contentType: string, fileId: string) => {
    const payload = {
        filename,
        bucket,
        contentType,
        fileId, // 用于Cloud Functions获取源文件的标识
    };

    const url = `${process.env.GCLOUD_FUNCTIONS_URL}/uploadToCS`;

    const task = {
        httpRequest: {
            httpMethod: google.cloud.tasks.v2beta3.HttpMethod.POST,
            url,
            oidcToken: {
                serviceAccountEmail: email,
                audience: url,
            },
            headers: {
                "Content-Type": "application/json",
            },
            body: Buffer.from(JSON.stringify(payload)).toString("base64"),
        },
    };

    const [response] = await client.createTask({ parent, task });
    const name = response.name ?? "name not exists";

    return name;
};

export default createUploadTask;

4. 升级@google-cloud/storage版本

v6.9.5存在部分已知的网络和重试逻辑缺陷,升级至最新稳定版(如v8.x)可修复底层问题,同时新版本默认提供更完善的重试机制。

执行升级命令:

npm install @google-cloud/storage@latest

内容的提问来源于stack exchange,提问作者Michał Dubrowski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 04:43:10