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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 14:53:13