Node.js服务器间JSON流传输大数据处理报错解决方案咨询
大集合流式传输时Premature Close错误的解决方案
问题背景
需要跨Node.js服务器实现MongoDB数据流式迁移:从A服务器读取MongoDB集合数据,以JSON数组形式通过REST接口流式发送到B服务器,B服务器使用stream-json解析后批量写入本地MongoDB。小数据量场景下运行正常,但处理1.2GB大集合时,接收端持续报Premature close错误。
发送端代码
export const streamData = async (res: Response) => { try { res.type('json'); const amountOfItems = await MyModel.count(); if (JSON.stringify(amountOfItems) !== '0'){ const cursor = MyModel.find().cursor(); let first = true; cursor.on('error', (err) => { logger.error(err); }); cursor.on('data', (doc) => { if (first) { // open json array res.write('['); first = false; } else { // add the delimiter before every object that isn't the first res.write(','); } // add json object res.write(`${JSON.stringify(doc)}`); }); cursor.on('end', () => { // close json array res.write(']'); res.end(); logger.info('REST-API-Call to fetchAllItems: Streamed all items to the receiver.'); }); } else { res.write('[]'); res.end(); logger.info('REST-API-Call to fetchAllItems: Streamed an empty response to the receiver.'); } } catch (err) { logger.error(err); return []; } };
接收端代码
import { MyModel } from '../models/my-model'; import axios from 'axios'; import { logger } from '../services/logger'; import StreamArray from 'stream-json'; import { streamArray } from 'stream-json/streamers/StreamArray'; import { pipeline } from 'stream'; const persistItems = async (items:Item[], ip: string) => { try { await MyModel.bulkWrite(items.map(item => { return { updateOne: { filter: { 'itemId': item.itemId }, update: item, upsert: true, }, }; })); logger.info(`${ip}: Successfully upserted items to mongoDB`); } catch (err) { logger.error(`${ip}: Upserting items to mongoDB failed due to the following error: ${err}`); } }; const getAndPersistDataStream = async (ip: string) => { try { const axiosStream = await axios(`http://${ip}:${process.env.PORT}/api/items`, { responseType: 'stream' }); const jsonStream = StreamArray.parser({ jsonStreaming : true }); let items : Item[] = []; const stream = pipeline(axiosStream.data, jsonStream, streamArray(), (error) => { if ( error ){ logger.error(`Error: ${error}`); } else { logger.info('Pipeline successful'); } }, ); stream.on('data', (i: any) => { items.push(<Item> i.value); // wait until the array contains 500 objects, than bulkWrite them to database if (items.length === 500) { persistItems(items, ip); items = []; } }); stream.on('end', () => { // bulkwrite the last items to the mongodb persistItems(items, ip); }); stream.on('error', (err: any) => { logger.error(err); }); await new Promise(fulfill => stream.on('finish', fulfill)); } catch (err) { if (err) { logger.error(err); } } }
报错信息
ERROR: Premature close err: { "type": "NodeError", "message": "Premature close", "stack": Error [ERR_STREAM_PREMATURE_CLOSE]: Premature close at IncomingMessage.onclose (internal/streams/end-of-stream.js:75:15) at IncomingMessage.emit (events.js:314:20) at Socket.socketCloseListener (_http_client.js:384:11) at Socket.emit (events.js:326:22) at TCP.<anonymous> (net.js:676:12) "code": "ERR_STREAM_PREMATURE_CLOSE" }
问题原因及解决方案
核心问题分析
- 发送端未处理背压:直接调用
res.write()时未检查返回值,当响应缓冲区满时未暂停MongoDB游标,导致数据积压,最终连接被强制关闭。 - 接收端未控制流速:异步批量写入未等待完成,也未暂停流,数据库写入速度跟不上读取速度时,内存积压大量待处理数据,触发流异常关闭。
- JSON序列化效率低:逐个调用
JSON.stringify(doc)产生额外性能开销,大集合下发送速度不稳定。
优化后的发送端代码
export const streamData = async (res: Response) => { try { res.type('json'); const count = await MyModel.countDocuments(); // MongoDB 4.0+推荐用countDocuments替代count if (count === 0) { res.end('[]'); logger.info('REST-API-Call to fetchAllItems: Streamed empty response.'); return; } const cursor = MyModel.find().cursor(); let isFirst = true; let isPaused = false; // 写入数组开头 res.write('['); cursor.on('data', (doc) => { if (isPaused) return; if (!isFirst) { // 写入分隔符并处理背压 if (!res.write(',')) { isPaused = true; cursor.pause(); } } isFirst = false; // 序列化文档并处理背压 if (!res.write(JSON.stringify(doc))) { isPaused = true; cursor.pause(); } }); // 响应缓冲区排空后恢复游标 res.on('drain', () => { isPaused = false; cursor.resume(); }); cursor.on('end', () => { res.end(']'); logger.info('REST-API-Call to fetchAllItems: Streamed all items successfully.'); }); cursor.on('error', (err) => { logger.error('Cursor error:', err); res.status(500).end(JSON.stringify({ error: 'Stream failed' })); }); res.on('close', () => { // 客户端断开时销毁游标,避免资源泄漏 cursor.close(); logger.info('Client closed connection prematurely.'); }); } catch (err) { logger.error('Stream initialization error:', err); res.status(500).end(JSON.stringify({ error: 'Stream initialization failed' })); } };
优化后的接收端代码
import { MyModel } from '../models/my-model'; import axios from 'axios'; import { logger } from '../services/logger'; import { streamArray } from 'stream-json/streamers/StreamArray'; import { pipeline } from 'stream/promises'; // 使用Promise版pipeline,简化异步处理 import { Transform } from 'stream'; const persistItems = async (items: Item[], ip: string) => { try { await MyModel.bulkWrite(items.map(item => ({ updateOne: { filter: { itemId: item.itemId }, update: { $set: item }, // 用$set避免覆盖整个文档,可按需调整 upsert: true, }, }))); logger.info(`${ip}: Upserted ${items.length} items successfully.`); } catch (err) { logger.error(`${ip}: Upsert failed:`, err); throw err; // 抛出错误让pipeline统一处理 } }; // 自定义Transform流,控制批量写入速度 class BatchTransform extends Transform { private batchSize: number; private currentBatch: Item[]; private ip: string; constructor(batchSize: number, ip: string) { super({ objectMode: true }); this.batchSize = batchSize; this.currentBatch = []; this.ip = ip; } async _transform(chunk: any, encoding: BufferEncoding, callback: (error?: Error | null) => void) { this.currentBatch.push(chunk.value as Item); if (this.currentBatch.length >= this.batchSize) { try { await persistItems(this.currentBatch, this.ip); this.currentBatch = []; callback(); } catch (err) { callback(err as Error); } } else { callback(); } } async _flush(callback: (error?: Error | null) => void) { // 处理剩余的批量数据 if (this.currentBatch.length > 0) { try { await persistItems(this.currentBatch, this.ip); callback(); } catch (err) { callback(err as Error); } } else { callback(); } } } const getAndPersistDataStream = async (ip: string) => { try { const axiosStream = await axios(`http://${ip}:${process.env.PORT}/api/items`, { responseType: 'stream', timeout: 0, // 禁用超时,避免大文件传输时断开 }); // 使用Promise版pipeline,自动处理错误和流销毁 await pipeline( axiosStream.data, streamArray(), // streamArray已包含parser,无需单独创建 new BatchTransform(500, ip) ); logger.info(`${ip}: Data stream processed completely.`); } catch (err) { logger.error(`${ip}: Stream processing failed:`, err); } };
额外优化建议
- 发送端简化JSON流处理:使用
JSONStream库自动处理背压和JSON数组格式,替代手动拼接:import JSONStream from 'JSONStream'; // 替换原cursor处理逻辑 res.type('json'); cursor.stream().pipe(JSONStream.stringify('[', ',', ']')).pipe(res); - 接收端优化批量写入:根据数据库性能调整
batchSize,或在bulkWrite中添加ordered: false实现并行写入(无需顺序一致性时),提升写入速度。 - 添加心跳机制:在HTTP请求中加入自定义心跳,避免长时间无数据传输时连接被防火墙/负载均衡器断开。
内容的提问来源于stack exchange,提问作者gold
相关产品推荐
相关产品推荐

