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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 19:38:34