如何批量处理异步读取流?批量入库异常问题排查
批量处理文件行写入数据库时的并发问题
我在处理文件读取结果批量写入数据库时,设置每次批量处理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
相关产品推荐
相关产品推荐

