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

Node.js中如何等待可写流结束并封装为异步函数

将Node.js CSV写入代码封装为异步函数的解决方案

问题描述

需要将以下Node.js代码封装为异步函数,确保文件内容完全写入后再resolve,但为stringifier和writableStream添加的close/finish事件监听器均未被触发:

const fs = require("fs");
const { stringify } = require("csv-stringify");
const db = require("./db");
const filename = "saved_from_db.csv";
const writableStream = fs.createWriteStream(filename);

const columns = [
  "year_month",
  "month_of_release",
  "passenger_type",
  "direction",
  "sex",
  "age",
  "estimate",
];

const stringifier = stringify({ header: true, columns: columns });
db.each(`select * from migration`, (error, row) => {
  if (error) {
    return console.log(error.message);
  }
  stringifier.write(row);
});
stringifier.pipe(writableStream);
console.log("Finished writing data");

添加的事件监听器代码:

stringifier.on('close', ()=> console.log( `STR CLOSED`))
stringifier.on('finish', ()=> console.log( `STR FINISH`))
writableStream.on('close', ()=> console.log( `WS CLOSED`))
writableStream.on('finish', ()=> console.log( `WS FINISH`))

核心原因

事件未触发的关键原因是:db.each是异步遍历数据库记录的方法,遍历过程中stringifier一直在接收数据,但遍历结束后没有手动结束stringifier流,导致流始终处于等待写入的状态,不会触发finish或close事件。

解决方案

将代码封装为Promise,利用db.each的第三个回调函数(遍历完成时触发)来结束stringifier,同时监听流的finish事件来resolve Promise:

const fs = require("fs");
const { stringify } = require("csv-stringify");
const db = require("./db");

async function exportDbToCsv() {
  return new Promise((resolve, reject) => {
    const filename = "saved_from_db.csv";
    const writableStream = fs.createWriteStream(filename);

    const columns = [
      "year_month",
      "month_of_release",
      "passenger_type",
      "direction",
      "sex",
      "age",
      "estimate",
    ];

    const stringifier = stringify({ header: true, columns: columns });

    // 监听流的finish事件,确认所有数据写入完成
    writableStream.on('finish', () => {
      console.log("Finished writing data");
      resolve(filename);
    });

    // 统一处理各环节错误
    stringifier.on('error', (err) => reject(err));
    writableStream.on('error', (err) => reject(err));

    let hasError = false;
    db.each(`select * from migration`, (error, row) => {
      if (error) {
        hasError = true;
        console.log(error.message);
        reject(error);
      }
      if (!hasError) {
        stringifier.write(row);
      }
    }, (err, count) => {
      // 数据库遍历完成后,结束stringifier流
      if (!hasError && !err) {
        stringifier.end();
      }
    });

    stringifier.pipe(writableStream);
  });
}

// 使用示例
// exportDbToCsv().then(filename => console.log(`导出完成:${filename}`)).catch(err => console.error(err));

关键改动说明

  • 用Promise包裹整个逻辑,成功时返回文件名,失败时抛出错误
  • 借助db.each的第三个回调,在数据库遍历结束后调用stringifier.end(),告知流不再有新数据写入
  • 监听writableStream的finish事件,该事件会在所有数据写入底层文件系统后触发,此时resolve Promise
  • 补充了数据库遍历、CSV序列化、文件写入各环节的错误捕获逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 16:34:55