如何在Node.js/NestJS中将PostgreSQL导出的CSV转为Parquet文件?
Node.js/NestJS 下 CSV 转 Parquet 实现方案
在Node.js/NestJS(≥16)环境下,无需切换到Python栈,可通过以下两种主流方案实现CSV到Parquet的转换:
方案一:使用 parquetjs + csv-parser
这是轻量级的组合,适合中小规模文件转换。
安装依赖
npm install parquetjs csv-parser fs-extra # 或 yarn add parquetjs csv-parser fs-extra
实现代码
const fs = require('fs-extra'); const csv = require('csv-parser'); const parquet = require('parquetjs'); async function convertCsvToParquet(csvPath, parquetPath, schema) { const records = []; // 流式读取并解析CSV await new Promise((resolve, reject) => { fs.createReadStream(csvPath) .pipe(csv()) .on('data', (data) => records.push(data)) .on('end', resolve) .on('error', reject); }); // 创建Parquet写入器 const writer = await parquet.ParquetWriter.openFile(schema, parquetPath); // 处理数据类型并写入 for (const record of records) { const processedRecord = Object.fromEntries( Object.entries(record).map(([key, value]) => { const field = schema.fields.find(f => f.name === key); if (!field) return [key, value]; // 根据schema转换数据类型 switch (field.type) { case 'int64': return [key, BigInt(value)]; case 'double': return [key, parseFloat(value)]; case 'timestamp': return [key, new Date(value)]; default: return [key, value]; } }) ); await writer.appendRow(processedRecord); } await writer.close(); } // 示例调用(需根据你的表结构调整schema) const targetSchema = new parquet.ParquetSchema({ id: { type: 'int64' }, username: { type: 'utf8' }, age: { type: 'int32' }, register_time: { type: 'timestamp' }, account_balance: { type: 'double' } }); convertCsvToParquet('./source.csv', './result.parquet', targetSchema) .then(() => console.log('转换完成')) .catch(err => console.error('转换失败:', err));
方案二:使用 apache-arrow + csv-parser
基于Apache Arrow(Parquet的底层支持标准),性能更优,适合大规模数据处理。
安装依赖
npm install apache-arrow csv-parser fs # 或 yarn add apache-arrow csv-parser fs
实现代码
const fs = require('fs'); const csv = require('csv-parser'); const { Table, Schema, RecordBatch, writeParquet } = require('apache-arrow'); async function convertCsvToParquet(csvPath, parquetPath, arrowSchema) { const records = []; // 读取CSV数据 await new Promise((resolve, reject) => { fs.createReadStream(csvPath) .pipe(csv()) .on('data', (data) => records.push(data)) .on('end', resolve) .on('error', reject); }); // 转换为Arrow格式并写入Parquet const batch = RecordBatch.from(records, arrowSchema); const table = new Table(batch); await writeParquet(table, parquetPath); } // 示例调用(调整schema匹配你的数据结构) const targetSchema = new Schema([ { name: 'id', type: 'int64' }, { name: 'username', type: 'utf8' }, { name: 'age', type: 'int32' }, { name: 'register_time', type: 'timestamp(ms)' }, { name: 'account_balance', type: 'float64' } ]); convertCsvToParquet('./source.csv', './result.parquet', targetSchema) .then(() => console.log('转换完成')) .catch(err => console.error('转换失败:', err));
NestJS 集成示例
将转换逻辑封装为服务,便于在NestJS项目中调用:
import { Injectable } from '@nestjs/common'; import { Table, Schema, RecordBatch, writeParquet } from 'apache-arrow'; import * as csv from 'csv-parser'; import * as fs from 'fs'; @Injectable() export class CsvToParquetService { async convert(csvPath: string, parquetPath: string, schema: Schema): Promise<void> { const records: Record<string, any>[] = []; await new Promise((resolve, reject) => { fs.createReadStream(csvPath) .pipe(csv()) .on('data', (data) => records.push(data)) .on('end', resolve) .on('error', reject); }); const batch = RecordBatch.from(records, schema); const table = new Table(batch); await writeParquet(table, parquetPath); } }
关键注意事项
- 必须定义schema:CSV是无类型文本,Parquet是强类型列存,需根据数据库表结构准确设置字段类型,避免数据失真。
- 大文件处理:若CSV文件过大,不要一次性读取全部数据到内存,可改用流式写入逻辑(
parquetjs和apache-arrow均支持流式API)。 - 数据类型转换:CSV中的数字、日期以字符串存储,需手动转换为对应类型(如代码中示例的类型处理逻辑)。
内容的提问来源于stack exchange,提问作者Falyoun
相关产品推荐
相关产品推荐

