Sax批量插入PostgreSQL时记录丢失问题求助
问题排查:sax+zlib读取gz文件批量更新PostgreSQL时部分记录丢失
场景与问题
使用sax解析XML、zlib解压gz文件批量更新PostgreSQL时,部分记录未被处理。代码无报错,但日志显示记录数不匹配,存在失败对象,且每次运行数值波动。
代码片段
主处理逻辑
// 异常行为观测变量 let recordsNum = 1; let recordsNum2 = 1; let numOfFailed = 0; let arrOfRecoreds = []; const keyTransforms = {}; // 假设此处有键映射逻辑 async function boo(...args){ const saxStream = sax.createStream(true); // 严格模式 let currentElement = null; let currentObject = null; saxStream.on("opentag", (node) => { if (node.name === "Item") { currentObject = {}; } else { currentElement = keyTransforms[node.name]; } }); saxStream.on("closetag", async (nodeName) => { if (nodeName === "Item") { recordsNum++; if (arrOfRecoreds.length <= 100 && currentObject) { arrOfRecoreds.push(currentObject); recordsNum2++; } if (!currentObject) { numOfFailed++; } if (arrOfRecoreds.length === 100) { await insertBatch(arrOfRecoreds); arrOfRecoreds = []; } currentObject = null; } else { currentElement = null; } }); const fileName = args[1]; const readStream = fs .createReadStream("./temp/" + fileName) .pipe(zlib.createGunzip()) .pipe(saxStream); await new Promise((resolve, reject) => { readStream.on("end", async () => { if (arrOfRecoreds.length > 0) { await insertBatch(arrOfRecoreds); resolve(); } }); readStream.on("error", reject); }); return new Promise((res) => { pool.end(); console.log("Finished processing file."); res({ statusCode: 201, message: "update db" }); }); } // 调用示例 boo("str", "filename.rar").then(() => { console.log("num of recordes ", recordsNum); console.log("num of recordes2 ", recordsNum2); console.log("num of failed ", numOfFailed); });
批量插入函数
async function insertBatch(records) { const client = await pool.connect(); try { await client.query("BEGIN"); await doSomthing(records); // 批量更新逻辑 await client.query("COMMIT"); } catch(e) { await client.query("ROLLBACK"); } finally { client.release(); } }
日志输出
num of records 100 num of recordes2 45 num of failed 6
核心疑问
recordsNum与recordsNum2数值为何不相等?- 已确认XML有效,为何存在失败对象?
recordsNum2与失败对象数值为何每次运行都变化?
问题根源与修复方案
1. 异步回调与流处理的冲突
sax的closetag事件回调是异步函数,但Node.js流会同步触发事件,不会等待异步操作完成。当你在closetag中await insertBatch时,流会继续解析后续XML内容,导致:
currentObject被下一个Item的opentag提前覆盖为新对象,原本的currentObject可能未被推入数组- 批量插入期间,新的
Item会因数组长度判断逻辑错误被跳过,或currentObject被置为null,触发numOfFailed计数
2. 关键逻辑缺失与错误
- 未处理XML文本内容:当前代码仅处理了标签的开闭,没有把元素的文本值存入
currentObject,如果keyTransforms映射有问题或文本事件未处理,会导致currentObject无有效数据 - 数组长度判断错误:
arrOfRecoreds.length <= 100会在数组已满(100条)时仍尝试推入新记录,与后续的批量插入逻辑冲突 end事件Promise未正确resolve:当arrOfRecoreds为空时,resolve()未被调用,导致Promise永远pending,后续逻辑异常
具体修复步骤
(1)暂停流以等待异步操作完成
在批量插入时暂停流,避免异步期间流继续推送数据:
if (arrOfRecoreds.length === 100) { readStream.pause(); // 暂停解析 await insertBatch(arrOfRecoreds); arrOfRecoreds = []; readStream.resume(); // 恢复解析 }
(2)补充XML文本内容处理
添加text事件回调,将元素文本存入currentObject:
saxStream.on("text", (text) => { if (currentElement && currentObject) { currentObject[currentElement] = text.trim(); } });
(3)修正数组长度判断逻辑
将arrOfRecoreds.length <= 100改为arrOfRecoreds.length < 100,避免数组已满时继续推入:
if (arrOfRecoreds.length < 100 && currentObject) { arrOfRecoreds.push(currentObject); recordsNum2++; }
(4)修复end事件的Promise逻辑
无论数组是否为空,都要调用resolve():
readStream.on("end", async () => { if (arrOfRecoreds.length > 0) { await insertBatch(arrOfRecoreds); } resolve(); // 确保Promise始终完成 });
(5)避免全局变量竞争
将recordsNum、recordsNum2、numOfFailed移至boo函数内部,避免多次调用时的数值混乱:
async function boo(...args){ let recordsNum = 1; let recordsNum2 = 1; let numOfFailed = 0; let arrOfRecoreds = []; // ... 其他逻辑 // 在函数内输出日志 console.log("num of recordes ", recordsNum); console.log("num of recordes2 ", recordsNum2); console.log("num of failed ", numOfFailed); }
内容的提问来源于stack exchange,提问作者Bennyh961
相关产品推荐
相关产品推荐

