Запись быстрых данных в файлы через writable stream Node.js
Задача: Организовать запись в .csv файлы данных поступающих через websocket. При этом создание нового файла для записи при достижении размера файла( количество строк) определенного значения, допустим 1000 строк. Необходимо записать все данные поступившие по websocket. Точная скорость, периодичность, объемы данных неизвестны. предположительно волнами, предположительно пиковые скорости держатся до 1 минут, а потом запись прекращается на неизвестное время.
Пытаюсь организовать запись с помощью двух writable Stream потока:
let testWriteableStream_2 = fs.createWriteStream(`logs/test_profit_2.csv`, { flags: 'a' });;
let testFlag = { number: 1 };
let testCount = { number: 0 };
let testCountAll = { number: 0 };
function TestWritable(testWriteableStream_1, testWriteableStream_2, testFlag, testCount, testCountAll) {
console.log('testCountAll.number=', testCountAll.number);
console.log('function TestWritable----------------------------------------------------------------------------------------------------');
console.log(`testCount=${testCount.number}----------------------------------------------------------------------------------------------------`);
if (testCount.number > 10) {
testCount.number = 0;
if (testFlag.number === 1) {
testFlag.number = 2;
console.log('testFlag=2----------------------------------------------------------------------------------------------------');
testWriteableStream_1.end();
// testWriteableStream_1.close();
// if (writeableStream._writableState.closed) {
let time = new Date().getTime();
console.log('time:', time);
testWriteableStream_2 = fs.createWriteStream(`logs/test2_profit${time}.csv`, { flags: 'a' });
return
}
testFlag.number = 1;
console.log('testFlag=1----------------------------------------------------------------------------------------------------');
testWriteableStream_2.end();
// testWriteableStream_2.close();
// if (writeableStream._writableState.closed) {
let time = new Date().getTime();
console.log('time:', time);
testWriteableStream_1 = fs.createWriteStream(`logs/test1_profit${time}.csv`, { flags: 'a' });
};
console.log(`testFlag.number=${testFlag.number},----------------------------------------------------------------------------------------------------`);
if (testFlag.number === 1) {
console.log('writeableStream_1');
testWriteableStream_1.write(`writeableStream_1${testCountAll.number}\r\n`);
}
if (testFlag.number === 2) {
console.log('writeableStream_2');
testWriteableStream_2.write(`writeableStream_2${testCountAll.number}\r\n`);
}
testCount.number++;
testCountAll.number++;
}
ws.onmessage = function (message) {
TestWritable(testWriteableStream_1, testWriteableStream_2, testFlag, testCount, testCountAll);
}
Я хотел организовать поочередную запись, сначалов один поток, потом в другой. Но скорость поступающих данных быстрее скорости записи в файл. И когда я завершаю testWriteableStream_.end(); то в первоначальные созданные стримы данные успевают записаться, а уже в перезаписываемые в переменных testWriteableStream_1, testWriteableStream_2 записывается либо одна строка либо ни одной.
Пробовал не завершать .end сразу стримы, но все равно данные пропускались или вовсе не записывались.
Какими способами можно реализовать данную задачу?
Можно увеличить буфер стримов, но это кардинально не решит задачу.