在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); }) });
核心原因
- Promise链未等待流完成:云函数的返回Promise在
download完成后就结束了,流操作是异步的,但没有被纳入Promise链,云函数认为任务已完成,提前回收容器,导致流处理未执行。 - 同步map无法处理异步操作:
es.mapSync是同步方法,无法正确处理Firestorecreate返回的Promise,会导致异步任务被忽略,且错误无法被捕获。 - 流错误未监听:流的错误事件未被监听,若流初始化或解析出错,错误会被吞掉,无日志输出。
解决方案
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
相关产品推荐
相关产品推荐

