Merge pull request #79 from avandeweghe/patch-1
Refactor parseStreaming to use pendingBatches array
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