Refactor parseStreaming to use pendingBatches array
Was running into an issue where epg wouldn't process passed the first batch. Refactored the batching to use an array and added backpressure to keep memory usage low.
This commit is contained in:
@@ -266,9 +266,11 @@ async function* fetchAndParseStreaming(url, batchSize = 1000) {
|
||||
* @yields {{ channels: Array|null, programmes: Array, isLast: boolean }}
|
||||
*/
|
||||
async function* parseStreaming(input, batchSize = 1000) {
|
||||
const channels = [];
|
||||
let channels = [];
|
||||
let programmeBatch = [];
|
||||
let channelsYielded = false;
|
||||
const maxPenmdingBatches = 4;
|
||||
let paused = false;
|
||||
|
||||
// We need to convert SAX events to an async iterator
|
||||
// This requires collecting events and yielding when batch is full
|
||||
@@ -279,7 +281,7 @@ async function* parseStreaming(input, batchSize = 1000) {
|
||||
let currentObject = null;
|
||||
let textBuffer = '';
|
||||
let resolveNext = null;
|
||||
let pendingBatch = null;
|
||||
const pendingBatches = [];
|
||||
let ended = false;
|
||||
let error = null;
|
||||
|
||||
@@ -345,13 +347,27 @@ async function* parseStreaming(input, batchSize = 1000) {
|
||||
isLast: false
|
||||
};
|
||||
channelsYielded = true;
|
||||
|
||||
// allow channels array to be GC'd after included in first batch
|
||||
try { if (channels && channels.length) channels = null; } catch (e) {}
|
||||
|
||||
programmeBatch = [];
|
||||
|
||||
if (resolveNext) {
|
||||
resolveNext(batch);
|
||||
resolveNext = null;
|
||||
} else {
|
||||
pendingBatch = batch;
|
||||
pendingBatches.push(batch);
|
||||
|
||||
// Apply backpressure when queue grows too large
|
||||
try {
|
||||
if (!paused && pendingBatches.length >= maxPenmdingBatches) {
|
||||
if (input && typeof input.pause === 'function') {
|
||||
input.pause();
|
||||
paused = true;
|
||||
}
|
||||
}
|
||||
} catch (e) {}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -398,7 +414,7 @@ async function* parseStreaming(input, batchSize = 1000) {
|
||||
resolveNext(batch);
|
||||
resolveNext = null;
|
||||
} else {
|
||||
pendingBatch = batch;
|
||||
pendingBatches.push(batch);
|
||||
}
|
||||
});
|
||||
|
||||
@@ -413,10 +429,21 @@ async function* parseStreaming(input, batchSize = 1000) {
|
||||
input.pipe(saxStream);
|
||||
|
||||
// Yield batches as they become available
|
||||
while (!ended || pendingBatch) {
|
||||
if (pendingBatch) {
|
||||
const batch = pendingBatch;
|
||||
pendingBatch = null;
|
||||
|
||||
while (!ended || pendingBatches.length > 0) {
|
||||
if (pendingBatches.length > 0) {
|
||||
const batch = pendingBatches.shift();
|
||||
|
||||
// If paused and queue drained below threshold, resume
|
||||
try {
|
||||
if (paused && pendingBatches.length < maxPenmdingBatches) {
|
||||
if (input && typeof input.resume === 'function') {
|
||||
input.resume();
|
||||
paused = false;
|
||||
}
|
||||
}
|
||||
} catch (e) {}
|
||||
|
||||
yield batch;
|
||||
if (batch.isLast) break;
|
||||
} else if (!ended) {
|
||||
@@ -445,4 +472,3 @@ module.exports = {
|
||||
getProgrammesForChannel,
|
||||
getCurrentAndUpcoming
|
||||
};
|
||||
|
||||
|
||||
Reference in New Issue
Block a user