【问题标题】:How to improve the processing of large csv file and insert data to remote database (MYSQL) in node.js?如何改进 node.js 中大型 csv 文件的处理并将数据插入远程数据库(MYSQL)?
【发布时间】:2022-01-09 03:37:25
【问题描述】:

这是一个工作示例,我正在寻找进一步改进它的方法,并在代码级别提供性能提升,如果不寻找其他可扩展性选项。

我应用的东西:

  • Nodejs 流 -> 改进了具有背压的大型 csv 文件的读取和解析操作。
  • 使用 await 循环进行批处理 -> 用于处理背压并删除了 nodejs 内存堆错误
  • 在 Promise.all 中传递任务以进行异步操作
  • 数据库池 -> 提高了性能,但需要注意最大池大小
  • 子进程拆分文件然后应用上述操作 -> 尚未实现任何性能提升。
  • 更好地处理每个异步操作的错误,以了解哪个进程正常工作以及哪个进程抛出错误。

我需要处理 csv 文件中的数据,该文件可能有 100 万行。

处理完每一行后,必须检查数据库中的值以进行一些验证,然后将处理后的数据插入到不同的表中。

我从简单的节点流和异步库开始,但是速度太慢了,在 3-5 分钟内 2k 行,有时它会引发 nodejs 内存错误

经过一些迭代和研究,在 20-25 秒内将速度提高到 5k 行,并避免了内存已满错误。

瓶颈是即使我改变batch size时间也没有改变,为数据库增加了pooling。

我可以增加池大小以提高速度,但是需要知道我们如何确定最大池大小。如果池连接适用于每个连接或整体连接,鉴于我无法更改 mysql 的默认连接,因为没有管理员访问权限

有哪些方法可以提高这个速度?

这是代码#

insertBig.ts

import { Request, Response } from 'express';
import * as fs from 'fs';
import * as path from 'path';
import { getProductById, insertProduct } from '../repo';
import { performance } from 'perf_hooks';

import Papa from 'papaparse';

let processedNum = 0;

const insertBig = async (req: Request, res: Response) => {
  try {
    const filePath = path.normalize(`${__dirname}./../assets/test.csv`);
    importCSV(filePath).catch((error) => {
      return res.status(500).send('Some error here in isert biG 123');
    });
    return res.status(200).json({ data: 'Data send for processing 123' });
  } catch (error) {
    return res.status(500).send('Some error here in isert biG 123');
  }
};

async function importCSV(filePath: fs.PathLike) {
  let parsedNum = 0;
  const dataStream = fs.createReadStream(filePath);
  const parseStream = Papa.parse(Papa.NODE_STREAM_INPUT, {
    header: true
  });
  dataStream.pipe(parseStream);
  let buffer = [];
  let totalTime = 0;
  const startTime = performance.now();
  for await (const row of parseStream) {
    //console.log('PA#', parsedNum, ': parsed');
    buffer.push(row);
    parsedNum++;
    if (parsedNum % 400 == 0) {
      await dataForProcessing(buffer);
      buffer = [];
    }
  }
  totalTime = totalTime + (performance.now() - startTime);
  console.log(`Parsed ${parsedNum} rows and took ${totalTime} seconds`);
}

const wrapTask = async (promise: any) => {
  try {
    return await promise;
  } catch (e) {
    return e;
  }
};

const handle = async (promise: Promise<any>) => {
  try {
    const data = await promise;
    return [data, undefined];
  } catch (error) {
    return await Promise.resolve([undefined, error]);
  }
};

const dataForProcessing = async (arrayItems: any) => {
  const tasks = arrayItems.map(task);
  const startTime = performance.now();
  console.log(`Tasks starting...`);
  console.log('DW#', processedNum, ': dirty work START');
  try {
    await Promise.all(tasks.map(wrapTask));

    console.log(
      `Task finished in ${performance.now() - startTime} miliseconds with,`
    );
    processedNum++;
  } catch (e) {
    console.log('should not happen but we never know', e);
  }
};

const task = async (item: any) => {
  let table = 'Product';

  if (item.contactNumber == '9999999999') {
    table = 'random table'; // to create read error
  }
  if (item.contactNumber == '11111111111') {
    item.randomRow = 'random'; // to create insert error
  }

  // To add some read process
  const [data, readError] = await handle(getProductById(2, table));
  if (readError) {
    return 'Some error in read of table';
  }
  //console.log(JSON.parse(JSON.stringify(data))[0]['customerName']);
  data;
  // To add some write process
  const [insertId, insertErr] = await handle(insertProduct(item));
  if (insertErr) {
    return `Some error to log and continue process for ${item}`;
  }
  return `Done for ${insertId}`;
};

export { insertBig };

repo.ts

import pool from './dbConfig';

const getProducts = (table: any) => {
  return new Promise((resolve, reject) => {
    pool.query(`SELECT * FROM ${table}`, (error, results) => {
      if (error) {
        reject(error);
      }
      resolve(results);
    });
  });
};
const getProductById = (id: any, table: any) => {
  return new Promise((resolve, reject) => {
    pool.query(`SELECT * FROM ${table} WHERE id = ${id}`, (error, results) => {
      if (error) {
        reject(error);
      }
      resolve(results);
    });
  });
};

const insertProduct = (data: any) => {
  return new Promise((resolve, reject) => {
    pool.query(
      `INSERT INTO Product SET ?`,
      [
        {
          ...data
        }
      ],
      (error, results) => {
        if (error) {
          reject(error);
        }
        resolve(results);
      }
    );
  });
};

export { insertProduct, getProducts, getProductById };

dbConfig.ts

import mysql from 'mysql';
import * as dotenv from 'dotenv';
dotenv.config();

const dbConn = {
  connectionLimit: 80,
  host: process.env.DB_HOST,
  user: process.env.DB_USER,
  password: process.env.DB_PASSWORD,
  database: process.env.DB_NAME
};

const pool = mysql.createPool(dbConn);

pool.getConnection((err, connection) => {
  if (err) {
    if (err.code === 'PROTOCOL_CONNECTION_LOST') {
      console.error('Database connection was closed.');
    }
    if (err.code === 'ER_CON_COUNT_ERROR') {
      console.error('Database has to many connections');
    }
    if (err.code === 'ECONNREFUSED') {
      console.error('Database connection was refused');
    }
  }

  if (connection) {
    connection.release();
  }
  console.log('DB pool is Connected');
  return;
});

// pool.query = promisify(pool.query);

export default pool;

种子数据

import { OkPacket } from 'mysql';
import pool from './dbConfig';

const seed = () => {
  const queryString = `CREATE TABLE IF NOT EXISTS Product (
  id int(11) NOT NULL,
  customerName varchar(100) DEFAULT NULL,
  contactNumber varchar(100) DEFAULT NULL,
  modelName varchar(255) NOT NULL,
  retailerName varchar(100) NOT NULL,
  dateOfPurchase varchar(100) NOT NULL,
  voucherCode varchar(100) NOT NULL,
  voucherValue int(10) DEFAULT NULL,
  surveyUrl varchar(255) NOT NULL,
  surveyId varchar(255) NOT NULL,
  createdAt timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP,
  updatedAt timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;`;

  pool.query(queryString, (err, result) => {
    if (err) {
      console.log(err);
    }

    const insertId = (<OkPacket>result).insertId;
    console.log(insertId);
  });
};

export default seed;

可以做索引来增加从mysql的读取

增加子进程是否会产生任何影响,因为 nodejs 本身正在使用所有可用的池连接异步调用 dataProcessing。

根据一条评论: 我将文件分成 8 个部分,它实际上减少了处理,并且在性能提升方面似乎无关紧要。

child.ts

import * as fs from 'fs';
import * as path from 'path';
import { getProductById, insertProduct } from '../repo';
import { performance } from 'perf_hooks';

import Papa from 'papaparse';

let processedNum = 0;

process.on('message', async function (message: any) {
  console.log('[child] received message from server:', message);
  JSON.stringify(process.argv);
  const filePath = path.normalize(
    `${__dirname}./../output/output.csv.${message}`
  );

  let time = await importCSV(filePath, message);
  if (process.send) {
    process.send({
      child: process.pid,
      result: message + 1,
      time: time
    });
  }

  process.disconnect();
});

async function importCSV(filePath: fs.PathLike, message: any) {
  let parsedNum = 0;
  const dataStream = fs.createReadStream(filePath);
  const parseStream = Papa.parse(Papa.NODE_STREAM_INPUT, {
    header: true
  });
  dataStream.pipe(parseStream);
  let buffer = [];
  let totalTime = 0;
  const startTime = performance.now();
  for await (const row of parseStream) {
    // console.log('Child Server # :', message, 'PA#', parsedNum, ': parsed');
    buffer.push(row);
    parsedNum++;
    if (parsedNum % 400 == 0) {
      await dataForProcessing(buffer, message);
      buffer = [];
    }
  }

  totalTime = totalTime + (performance.now() - startTime);
  //   console.log(
  //     `Child Server ${message} : Parsed ${parsedNum} rows and took ${totalTime} seconds`
  //   );
  return totalTime;
}

const wrapTask = async (promise: any) => {
  try {
    return await promise;
  } catch (e) {
    return e;
  }
};

const handle = async (promise: Promise<any>) => {
  try {
    const data = await promise;
    return [data, undefined];
  } catch (error) {
    return await Promise.resolve([undefined, error]);
  }
};

const dataForProcessing = async (arrayItems: any, message: any) => {
  const tasks = arrayItems.map(task);
  const startTime = performance.now();
  console.log(`Tasks starting... from server ${message}`);
  console.log('CS#: ', message, 'DW#:', processedNum, ': dirty work START');
  try {
    await Promise.all(tasks.map(wrapTask));

    console.log(
      `Task finished in ${performance.now() - startTime} miliseconds with,`
    );
    processedNum++;
  } catch (e) {
    console.log('should not happen but we never know', e);
  }
};

const task = async (item: any) => {
  let table = 'Product';

  if (item.contactNumber == '8800210524') {
    table = 'random table'; // ro create read error
  }
  if (item.contactNumber == '9134743017') {
    item.randomRow = 'random'; // to create insert error
  }

  // To add some read process
  const [data, readError] = await handle(getProductById(2, table));
  if (readError) {
    return 'Some error in read of table';
  }
  //console.log(JSON.parse(JSON.stringify(data))[0]['customerName']);
  data;
  // To add some write process
  const [insertId, insertErr] = await handle(insertProduct(item));
  if (insertErr) {
    return `Some error to log and continue process for ${item}`;
  }
  return `Done for ${insertId}`;
};

parent.ts

import { Request, Response } from 'express';
var child_process = require('child_process');

const insertBigChildProcess = async (req: Request, res: Response) => {
  try {
    var numchild = require('os').cpus().length;
    var done = 0;
    let totalProcessTime: any[] = [];
    for (var i = 1; i <= numchild; i++) {
      const child = child_process.fork(__dirname + '/child.ts');
      child.send(i);
      child.on('message', function (message: any) {
        console.log('[parent] received message from child:', message);
        totalProcessTime.push(message.time);
        const sum = totalProcessTime.reduce(
          (partial_sum, a) => partial_sum + a,
          0
        );
        console.log(sum); // 6
        console.log(totalProcessTime);
        done++;
        if (done === numchild) {
          console.log('[parent] received all results');
        }
      });
    }

  
    return res.status(200).json({ data: 'Data send for processing 123' });
  } catch (error) {
    return res.status(500).send('Some error here in isert biG 123');
  }
};

【问题讨论】:

    标签: mysql node.js async-await


    【解决方案1】:

    处理每一行后,必须检查数据库中的值 一些验证,然后将处理后的数据插入到不同的 表格。

    是否可以从缓存中进行此验证?它可以提高性能,因为当我们谈论数百万行时,每行验证都需要很长时间。

    您可以将csv 文件分成更小的csv 文件,然后您可以并行处理这些文件。可以使用this包来实现。

    一旦你这样做了,使用子进程并行运行这些文件。

    // parent.js
    var child_process = require('child_process');
    
    var numchild  = require('os').cpus().length;
    var done      = 0;
    
    for (var i = 0; i < numchild; i++){
      var child = child_process.fork('./child');
      child.send((i + 1) * 1000);
      child.on('message', function(message) {
        console.log('[parent] received message from child:', message);
        done++;
        if (done === numchild) {
          console.log('[parent] received all results');
          ...
        }
      });
    }
    
    // child.js
    process.on('message', function(message) {
      console.log('[child] received message from server:', message);
      setTimeout(function() {
        process.send({
          child   : process.pid,
          result  : message + 1
        });
        process.disconnect();
      }, (0.5 + Math.random()) * 5000);
    });
    

    #从this 线程复制。你可以试一试,看看现在需要多少时间。

    【讨论】:

    • 另外,您将使用哪种托管服务? AWS?
    • 你在说哪个“池”? “连接池”可能无关紧要。 “buffer_pool”应该根据可用的 RAM 量来设置;它有助于大多数处理。并行处理可能是没有用的——它们很可能会相互干扰,从而无法提高性能。
    • @ApoorvaChikara 实现了这一点,但我认为连接池导致连接过多,结果也很突然,尽管我在子进程中使用相同的功能。
    【解决方案2】:

    我从“我可以用 SQL 完成这一切”来处理这样的任务。这可能是性能最高的。我会在“验证”和“处理”两种方式之间进行选择:

    LOAD DATA期间

        LOAD DATA ...
            ( ... @a, ..., @b, ...)
            SET cola = ... @a ...,
                colb = ... @b ...
    

    解释一下:

    • 在读取行时,会将一些列放入@变量中。
    • 然后在表达式/函数中使用这些变量来计算实际列的所需值。
    • 请注意,这是一种“忽略”列(不在SET 中使用它)或合并列的方法。

    LOADing之后

    运行UPDATE 语句以进行整体后处理。这可能比一次修复一行要快很多

    【讨论】:

    • 所以这是应用程序级别的限制,由于查询瓶颈而无法进一步改进,需要考虑对多个服务调用和数据库优化和查询进行优化?
    • @MrAJ - 有时最好的优化涉及重新设计整个应用程序。我的建议是在中间某处妥协。需要加载需要编辑的百万行文件听起来像是架构设计缺陷。
    猜你喜欢
    • 2021-09-20
    • 1970-01-01
    • 2014-05-03
    • 2016-01-26
    • 1970-01-01
    • 2021-03-15
    • 1970-01-01
    • 2011-10-19
    • 2020-06-22
    相关资源
    最近更新 更多