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

在Firebase云函数中下载处理大型JSON文件的异常排查

Firebase云函数处理大型JSON文件流式无响应问题解决

问题描述

将约2GB的大型JSON文件存储在Firebase Storage桶中,使用Node.js编写Firebase云函数,通过JSONStream包流式处理文件内容。部署后上传文件,函数日志仅输出“processing file”和“file downloaded”,之后无任何流处理相关日志,也无报错,函数运行约90秒后停止,怀疑容器在流执行前被清理。

原代码:

exports.importJSON = onObjectFinalized({ memory: "8GiB", cpu: 2, timeoutSeconds: 3600 }, (event) => {
    logger.info("processing file", { structuredData: true });

    const fileBucket = event.data.bucket; 
    const filePath = event.data.name; 

    var tempFilePath = tmp.tmpNameSync();

    const bucket = getStorage().bucket(fileBucket);

    return bucket.file(filePath).download({ destination: tempFilePath })
    .then(() => {
        logger.log('file downloaded');
        let stream = fs.createReadStream(tempFilePath);
        stream
            .pipe(JSONStream.parse('*'))
            .pipe(es.mapSync(function (data) {
                logger.log(data);
                return database.doc(`/ticketmaster_events/${uuidv4()}`).create({
                    name: data.eventName,
                    date: new Date(Date.parse(data.date)),
                    venue: data.venue.venueName,
                    // venueAddress: data.venueAddress,
                    locale: data.locale
                });
            }))
    })
    .catch((error) => {
        logger.error(error);
    }) 
});

核心原因

  1. Promise链未等待流完成:云函数的返回Promise在download完成后就结束了,流操作是异步的,但没有被纳入Promise链,云函数认为任务已完成,提前回收容器,导致流处理未执行。
  2. 同步map无法处理异步操作:es.mapSync是同步方法,无法正确处理Firestore create返回的Promise,会导致异步任务被忽略,且错误无法被捕获。
  3. 流错误未监听:流的错误事件未被监听,若流初始化或解析出错,错误会被吞掉,无日志输出。

解决方案

1. 将流操作包装为Promise

通过Promise包裹流的end和error事件,确保云函数等待流处理完成后再结束。

2. 使用异步map处理Firestore操作

替换es.mapSync为es.map(异步版本),并正确返回和等待Firestore的Promise,避免异步任务丢失。

3. 监听流错误事件

为读取流、JSONStream添加错误监听,确保错误能被捕获并输出日志。

4. 清理临时文件

处理完成或出错后删除临时文件,避免占用云函数临时存储空间。

修改后的代码

const fs = require('fs');
const tmp = require('tmp');
const JSONStream = require('JSONStream');
const es = require('event-stream');
const { v4: uuidv4 } = require('uuid');
const { onObjectFinalized } = require('firebase-functions/v2/storage');
const { getStorage } = require('firebase-admin/storage');
const { getFirestore } = require('firebase-admin/firestore');
const logger = require('firebase-functions/logger');

const database = getFirestore();

exports.importJSON = onObjectFinalized({ memory: "8GiB", cpu: 2, timeoutSeconds: 3600 }, (event) => {
    logger.info("processing file", { structuredData: true });

    const fileBucket = event.data.bucket; 
    const filePath = event.data.name; 
    const tempFilePath = tmp.tmpNameSync();
    const bucket = getStorage().bucket(fileBucket);

    return bucket.file(filePath).download({ destination: tempFilePath })
    .then(() => {
        logger.log('file downloaded');
        return new Promise((resolve, reject) => {
            const stream = fs.createReadStream(tempFilePath);
            
            // 监听读取流错误
            stream.on('error', (err) => {
                logger.error('Read stream error:', err);
                reject(err);
            });

            stream
                .pipe(JSONStream.parse('*'))
                // 监听JSON解析错误
                .on('error', (err) => {
                    logger.error('JSONStream parse error:', err);
                    reject(err);
                })
                .pipe(es.map((data, callback) => {
                    logger.log('Processing data:', data);
                    // 处理Firestore异步操作,完成后调用callback
                    database.doc(`/ticketmaster_events/${uuidv4()}`).create({
                        name: data.eventName,
                        date: new Date(Date.parse(data.date)),
                        venue: data.venue.venueName,
                        locale: data.locale
                    })
                    .then(() => callback(null))
                    .catch((err) => {
                        logger.error('Firestore create error:', err);
                        callback(err);
                    });
                }))
                .on('end', () => {
                    logger.log('Stream processing completed');
                    // 清理临时文件
                    fs.unlink(tempFilePath, (err) => {
                        if (err) logger.error('Failed to delete temp file:', err);
                        resolve();
                    });
                })
                .on('error', (err) => {
                    logger.error('Event stream error:', err);
                    // 出错时也清理临时文件
                    fs.unlink(tempFilePath, (unlinkErr) => {
                        if (unlinkErr) logger.error('Failed to delete temp file on error:', unlinkErr);
                        reject(err);
                    });
                });
        });
    })
    .catch((error) => {
        logger.error('Overall error:', error);
        // 出错时清理临时文件
        fs.unlink(tempFilePath, (unlinkErr) => {
            if (unlinkErr) logger.error('Failed to delete temp file on catch:', unlinkErr);
        });
        throw error; // 抛出错误让云函数标记为失败
    });
});

额外注意事项

  • 确保event-stream和JSONStream依赖已正确安装(执行npm install event-stream JSONStream)。
  • 若JSON文件结构不是顶级数组,需调整JSONStream.parse('*')的路径,比如对象下的数组用'data.*'。
  • 可考虑批量写入Firestore(如每100条批量提交),减少网络请求次数,提升处理效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 08:32:50