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

如何批量处理异步读取流?批量入库异常问题排查

批量处理文件行写入数据库时的并发问题

我在处理文件读取结果批量写入数据库时,设置每次批量处理20条记录,但功能异常:预期persistToDB会执行5次,实际只执行1次,而且documentBatch.length会超出20的限制。我尝试通过暂停流、完成数据库操作后恢复流的方式控制,但问题依旧。persistToDB是外部库提供的异步方法,怀疑是并发导致的问题。

当前代码:

let documentBatch = [];
const processedMetrics = {
  succesfullyProcessed: 0,
  unsuccesfullyProcessed: 0,
};

rl.on('line', async (line) => {
  try {
    const document = JSON.parse(line);
    documentBatch.push(document);
    console.log(documentBatch.length);
    if (documentBatch.length === 20) {
      rl.pause();
      const batchMetrics = await persistToDB(documentBatch);
      documentBatch = [];
      processedMetrics.succesfullyProcessed +=
        batchMetrics.succesfullyProcessed;
      processedMetrics.unsuccesfullyProcessed +=
        batchMetrics.unsuccesfullyProcessed;
      rl.resume();
    }
  } catch (e) {
    logger.error(`Failed to save document ${line}`);
    throw e;
  }
});

问题根源

rl.on('line')的回调是异步函数,但Node.js的readline模块不会等待异步回调完成再触发下一个line事件。哪怕调用了rl.pause(),在回调的异步代码执行完成前,新的line事件已经被放入事件队列,导致documentBatch被继续填充,长度超过20,且批量处理逻辑只会触发一次。

修复方案

用标志位控制是否允许处理新行,避免批量写入时继续接收数据:

let documentBatch = [];
const processedMetrics = {
  succesfullyProcessed: 0,
  unsuccesfullyProcessed: 0,
};
// 标记是否正在执行批量写入
let isProcessing = false;

rl.on('line', async (line) => {
  // 正在处理批量时,先暂存当前行
  if (isProcessing) {
    try {
      documentBatch.push(JSON.parse(line));
    } catch (e) {
      logger.error(`Failed to parse document ${line}`);
      processedMetrics.unsuccesfullyProcessed += 1;
    }
    return;
  }

  try {
    const document = JSON.parse(line);
    documentBatch.push(document);

    if (documentBatch.length >= 20) {
      isProcessing = true;
      rl.pause();
      // 拷贝当前批次,避免后续修改影响数据库操作
      const batchToProcess = [...documentBatch];
      documentBatch = [];
      
      const batchMetrics = await persistToDB(batchToProcess);
      processedMetrics.succesfullyProcessed += batchMetrics.succesfullyProcessed;
      processedMetrics.unsuccesfullyProcessed += batchMetrics.unsuccesfullyProcessed;
      
      isProcessing = false;
      rl.resume();
      
      // 如果暂存了新数据,手动触发一次处理
      if (documentBatch.length > 0) {
        rl.emit('line', '');
      }
    }
  } catch (e) {
    logger.error(`Failed to process document ${line}`);
    processedMetrics.unsuccesfullyProcessed += 1;
  }
});

// 处理文件读取结束后剩余的不足20条数据
rl.on('close', async () => {
  if (documentBatch.length > 0 && !isProcessing) {
    const batchMetrics = await persistToDB(documentBatch);
    processedMetrics.succesfullyProcessed += batchMetrics.succesfullyProcessed;
    processedMetrics.unsuccesfullyProcessed += batchMetrics.unsuccesfullyProcessed;
    console.log('所有数据处理完成', processedMetrics);
  }
});

关键改动

  • 新增isProcessing标志位,确保同一时间仅一个批量写入操作执行
  • 在close事件中处理剩余的不足批量大小的数据
  • 优化错误处理,避免单个数据解析失败终止整个进程
  • 批量处理时拷贝数组,避免暂存数据干扰正在执行的数据库操作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 19:50:18