Firebase函数Bulk Writer批量写入遇ECONNRESET崩溃求助
Firebase函数处理CSV写入Firestore时出现ECONNRESET及崩溃问题排查
问题概述
近期将.csv文件上传至Firebase Storage指定文件夹时,触发的云函数在读取文件内容并写入Firestore时出现异常。函数此前运行正常,近期可能更新过Firebase SDK和firebase-tools,已尝试更新SDK、添加try-catch块排查,但问题未解决,怀疑与Bulk Writer代码或Node.js v16有关。
错误信息
Error: Invalid response body while trying to fetch...: read ECONNRESET- 函数崩溃提示:
Function execution took 90053 ms, finished with status: 'crash'- 伴随错误:
Client network socket disconnected before secure TLS connection was established
补充信息
- 单文件最多10万行数据,预估30万行也不会超时
- 所有文档ID唯一,无重复写入场景
- 文件按顺序上传,等待前一个函数执行完成后再处理下一个
函数代码
// Create new documents from file on firebase storage export const createDocumentsFromCSV = functions.region(functionLocation).runWith({ timeoutSeconds: 540, memory: '8GB' }).storage.bucket().object().onFinalize(async (object: any) => { // Check that the file is a CSV file and located in the specified folder if (!object.name.endsWith('.csv') || !object.name.startsWith('carga-documentos/')) { return null; } const bulkWriter = firestoredb.bulkWriter(); let writeCount = 0; let batchCount = 0; const file = bucket.file(object.name); functions.logger.log(`Starting, fileName: ${object.name}`); const headers = ['id', 'client', 'creationTime', 'modificationTime']; // Read the CSV file from Firebase Storage const stream = file.createReadStream(); return new Promise<void>((resolve, reject) => { stream.pipe(csv({ headers, skipLines: 1, separator: ';' })) .on('data', (row: any) => { try { row.client = JSON.parse(row.client); row.creationTime = Timestamp.fromDate(new Date(row.creationTime)); row.modificationTime = Timestamp.fromDate(new Date(row.modificationTime)); writeCount++; if (writeCount % 500 === 0) { batchCount++; functions.logger.log(`Batch ${batchCount} committed with ${writeCount} writes`); } const docRef = firestoredb.collection('clients').doc(row.client.id).collection('xxxxx').doc(row.xxxxx).collection('xxxxx').doc(row.id); bulkWriter.set(docRef, row); } catch (error) { functions.logger.log(`Row causing error: ${JSON.stringify(row)}`); functions.logger.log(`Error: ${error}`); } }) .on('end', async () => { functions.logger.log(`Estimated number of batches: ${Math.ceil(writeCount / 500)}`); functions.logger.log(`Number of documents: ${writeCount}`); await bulkWriter.close(); functions.logger.log(`Finished, fileName: ${object.name}`); resolve(); }) .on('error', (error: any) => { reject(error); }); }); });
排查思路与解决方案
1. 解决流处理的背压问题
当前代码中CSV流的读取速度远快于Bulk Writer的写入速度,会导致内存堆积,进而触发网络连接中断。需在data事件中控制流的暂停与恢复:
.on('data', async (row: any) => { stream.pause(); // 暂停流避免背压 try { // 原有数据处理逻辑 row.client = JSON.parse(row.client); row.creationTime = Timestamp.fromDate(new Date(row.creationTime)); row.modificationTime = Timestamp.fromDate(new Date(row.modificationTime)); writeCount++; if (writeCount % 500 === 0) { batchCount++; functions.logger.log(`Batch ${batchCount} committed with ${writeCount} writes`); } const docRef = firestoredb.collection('clients').doc(row.client.id).collection('xxxxx').doc(row.xxxxx).collection('xxxxx').doc(row.id); await bulkWriter.set(docRef, row); // 等待写入完成 } catch (error) { functions.logger.log(`Row causing error: ${JSON.stringify(row)}`); functions.logger.log(`Error: ${error}`); } finally { stream.resume(); // 恢复流继续读取 } })
2. 调整Firestore客户端连接池配置
Node.js v16的TLS或连接池配置可能与新版Firebase SDK存在适配问题,显式配置连接参数优化稳定性:
const firestoredb = initializeFirestore(app, { ssl: true, poolSize: 10, keepAlive: true, maxIdleConnections: 5, idleConnectionTimeout: 30000 });
3. 自定义Bulk Writer重试与超时策略
默认重试策略无法应对网络波动,增加重试次数与超时时间:
const bulkWriter = firestoredb.bulkWriter({ maxRetries: 5, backoffFactor: 1.5, maxBackoffMs: 30000, timeout: 60000 });
4. 优化函数资源与终止逻辑
虽然已配置高内存与超时,但长时间处理可能导致连接被中断,优化函数终止前的资源释放:
.on('end', async () => { functions.logger.log(`Estimated number of batches: ${Math.ceil(writeCount / 500)}`); functions.logger.log(`Number of documents: ${writeCount}`); await bulkWriter.close(); functions.logger.log(`Finished, fileName: ${object.name}`); // 增加短暂延迟确保连接完全关闭 await new Promise(resolve => setTimeout(resolve, 1000)); resolve(); })
5. 完善流错误处理
确保流出错时正确关闭Bulk Writer,避免资源泄漏:
.on('error', async (error: any) => { functions.logger.error(`Stream error: ${error.message}`, error); await bulkWriter.close().catch(err => functions.logger.error(`Bulk writer close error: ${err}`)); reject(error); })
内容的提问来源于stack exchange,提问作者Nicolas Hurtado
相关产品推荐
相关产品推荐

