批量上传文件至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限制,导致请求阻塞。此外需检查权限配置:
修复步骤:
- 移除payload中的buffer数据:改为在Cloud Functions中从源头获取文件数据(比如生成buffer的服务接口、临时存储),仅传递文件标识信息。
- 检查权限:确保服务账号拥有
roles/cloudtasks.enqueuer角色,队列已正确创建。 - 验证网络:本地运行时,通过
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
相关产品推荐
相关产品推荐

