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

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"
    }

问题原因及解决方案

核心问题分析

  1. 发送端未处理背压:直接调用res.write()时未检查返回值,当响应缓冲区满时未暂停MongoDB游标,导致数据积压,最终连接被强制关闭。
  2. 接收端未控制流速:异步批量写入未等待完成,也未暂停流,数据库写入速度跟不上读取速度时,内存积压大量待处理数据,触发流异常关闭。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 10:45:45