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

Cloud Run接收Cloud Storage文件上传Pub/Sub推送时重复触发问题

问题:Cloud Run服务因Pub/Sub推送重复触发导致重复处理

我有一个基于Express端点的Cloud Run服务,当文件上传至Cloud Storage存储桶时,该服务会通过Pub/Sub推送通知触发。但目前每次文件上传都会导致Cloud Run服务被多次触发,进而造成重复处理。我已将确认期限设置为10分钟,但并未解决问题。

现有代码示例

主Express端点代码

app.post("/", async (req, res) => {
    
    try {
        console.log("Post endpoint is hit")
        const file = decodeBase64Json(req.body.message.data);
        const csvData = await downloadFile(storage, file.bucket, file.name);
        await removeHydrofoilsConfigurations();
        console.log("Processing of csv started")
        await processData(csvData);
        res.status(200).send("Cloud run finished csv processing successfully");
    } catch (ex) {
        console.log(`Error: ${ex}`);
        res.status(400).send("Cloud run failed during csv processing");
    }
    
});

downloadFile函数代码

async function downloadFile(storage, bucketName, fileName) {
    const options = { destination: `/tmp/${fileName}` };
    await storage.bucket(bucketName).file(fileName).download(options);
    const filePath = `/tmp/${fileName}`;
    const csvData = [];
    return new Promise((resolve, reject) => {
        fs.createReadStream(filePath)
            .pipe(csv())
            .on("data", (data) => csvData.push(data))
            .on("end", () => {
                try {
                    if (fs.existsSync(filePath)) {
                        fs.unlinkSync(filePath);
                        console.log("File deleted successfully");
                    } else {
                        console.log("File does not exist");
                    }
                } catch (err) {
                    console.error(err);
                }
                resolve(csvData);
            })
            .on("error", (error) => reject(error));
    });
}

解决方案

1. 强制实现幂等性(核心方案)

Pub/Sub的推送机制无法完全避免重复消息,因此必须让你的处理逻辑支持重复执行不产生副作用:

  • 利用Pub/Sub消息自带的messageId作为唯一标识,处理前先检查该ID是否已经被处理过(可存储到Cloud Firestore、Redis或Cloud Storage元数据中)
  • 也可以使用上传文件的唯一标识(如文件名+上传时间戳)来判断是否已处理

示例修改(加入幂等性检查):

// 假设已有checkIfProcessed和markAsProcessed函数实现存储校验逻辑
app.post("/", async (req, res) => {
    const messageId = req.body.message.messageId;
    // 先校验是否已处理
    const isProcessed = await checkIfProcessed(messageId);
    if (isProcessed) {
        console.log(`消息 ${messageId} 已处理,跳过`);
        return res.status(200).send("已处理");
    }

    try {
        console.log("Post endpoint is hit")
        const file = decodeBase64Json(req.body.message.data);
        const csvData = await downloadFile(storage, file.bucket, file.name);
        await removeHydrofoilsConfigurations();
        console.log("Processing of csv started")
        await processData(csvData);
        // 处理完成后标记为已处理
        await markAsProcessed(messageId);
        res.status(200).send("Cloud run finished csv processing successfully");
    } catch (ex) {
        console.log(`Error: ${ex}`);
        // 注意:不要返回400,否则Pub/Sub会认为处理失败并重试
        // 若为不可恢复错误,标记为已处理避免重复重试
        await markAsFailed(messageId);
        res.status(200).send("Cloud run failed during csv processing");
    }
});

2. 优化响应时机,异步处理任务

不要等到所有处理完成才返回200给Pub/Sub,避免因处理耗时过长导致的重试:

  • 收到请求后立即返回200确认消息,再通过Cloud Tasks等服务异步执行后续处理逻辑
  • 这样既可以快速确认Pub/Sub消息,又能保证任务被可靠执行

3. 检查配置合理性

  • 确保Cloud Run服务的超时时间≥Pub/Sub的确认期限(你设置了10分钟,需将Cloud Run超时调整为10分钟)
  • 调整Pub/Sub订阅的重试策略:设置合理的最大重试次数和退避时间,避免短时间内重复推送
  • 检查Cloud Run的并发设置,避免因实例扩容导致的重复触发

4. 排查代码性能瓶颈

  • 检查downloadFile函数的耗时:若文件过大,下载和解析会占用大量时间,可能导致响应延迟触发重试
  • 确保所有异步操作都正确使用await,避免未处理的Promise导致响应异常

内容的提问来源于stack exchange,提问作者Muhammad Bilal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 12:37:47