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

JavaScript中如何用流更优雅实现Parquet转CSV以降低内存占用

流式Parquet转CSV优化方案

核心思路

利用Node.js内置的流转换能力,将Parquet库返回的AsyncIterator游标直接转为可读流,通过管道串联CSV转换流和文件写入流,全程无全量数据内存暂存,内存占用仅和单条数据大小相关,完全解决大文件处理时的内存溢出问题。

优化后完整代码

import pts from 'parquets'
let { ParquetSchema, ParquetWriter, ParquetReader } = pts

import { createWriteStream } from 'fs'
import { Readable } from 'stream'
import { pipeline } from 'stream/promises'
import stringify from 'csv-stringify'

// 为`PI`表声明schema
let schema = new ParquetSchema({
    Source: { type: 'UTF8' },
    TagID: { type: 'UTF8' },
    Timestamp: { type: 'TIMESTAMP_MILLIS' },
    Value: { type: 'DOUBLE' },
});

const WriterParquet = async () => {
    // 创建写入'pi.parquet'的ParquetWriter实例
    let writer = await ParquetWriter.openFile(schema, 'pi.parquet')
    // 向文件追加若干行数据
    await writer.appendRow({Source: 'PI/NO-SVG-PISRV01', TagID: 'OGP8TI198Z.PV', Timestamp: new Date(), Value: 410 })
    await writer.appendRow({Source: 'PI/NO-SVG-PISRV01', TagID: 'OGP8TI198Z.PV', Timestamp: new Date(), Value: 420 }) 
    await writer.close()
}

const WriterCSV = async () => {
    // 创建读取'pi.parquet'的ParquetReader实例
    let reader = await ParquetReader.openFile('pi.parquet')
    // 创建游标
    let cursor = reader.getCursor()

    try {
        // 1. 将AsyncIterator游标转为Node可读流,开启对象模式处理JSON格式的行数据
        const parquetReadStream = Readable.from(cursor, { objectMode: true })
        // 2. 创建CSV转换流,开启表头输出
        const csvTransformStream = stringify({ header: true })
        // 3. 创建CSV文件写入流
        const csvWriteStream = createWriteStream('./pi.csv')

        // 串联全链路流式处理,自动处理背压、错误销毁
        await pipeline(
            parquetReadStream,
            csvTransformStream,
            csvWriteStream
        )
    } finally {
        // 无论处理成功失败都关闭Parquet读取器,避免文件句柄泄漏
        await reader.close()
    }
}

const Main = async () => {
    console.log('writing parquet...')
    await WriterParquet()
    console.log('reading parquet and writing csv...')
    await WriterCSV()
    console.log('convert finished')
}

Main().catch(console.error)

关键优化说明

  • 去掉了全量数据暂存的records数组,逐行读取逐行转换写入,内存占用稳定在KB级别
  • 采用Node官方推荐的pipeline方法串联流,比原生pipe更安全,任意环节出错时自动销毁所有流避免资源泄漏
  • 自动处理读写速度不匹配的背压问题,不需要手动控制读取节奏
  • 逻辑更简洁,避免了回调嵌套,全流程支持async/await异步等待

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 16:09:03