【发布时间】: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