65 lines
1.7 KiB
TypeScript
65 lines
1.7 KiB
TypeScript
import { Transform, type TransformCallback } from 'node:stream';
|
|
|
|
/**
|
|
* Backpressure on readable side is limits how many batches are in flight.
|
|
*
|
|
* @param size - maximum items per emitted batch
|
|
*/
|
|
export const batch = <T>(size: number): Transform => {
|
|
if (!Number.isInteger(size) || size < 1) {
|
|
throw new TypeError('batch(size) requires a positive integer size');
|
|
}
|
|
|
|
let pending: T[] = [];
|
|
|
|
return new Transform({
|
|
objectMode: true,
|
|
transform(item: T, _encoding: BufferEncoding, callback: TransformCallback) {
|
|
pending.push(item);
|
|
|
|
if (pending.length < size) {
|
|
callback();
|
|
return;
|
|
}
|
|
|
|
const full = pending;
|
|
pending = [];
|
|
callback(null, full);
|
|
},
|
|
flush(callback: TransformCallback) {
|
|
if (pending.length === 0) {
|
|
callback();
|
|
return;
|
|
}
|
|
|
|
const remainder = pending;
|
|
pending = [];
|
|
callback(null, remainder);
|
|
},
|
|
});
|
|
};
|
|
|
|
export const jsonArray = (): Transform => {
|
|
let wroteFirst = false;
|
|
|
|
return new Transform({
|
|
writableObjectMode: true,
|
|
transform(item: unknown, _encoding: BufferEncoding, callback: TransformCallback) {
|
|
let serialized: string;
|
|
try {
|
|
serialized = JSON.stringify(item);
|
|
} catch (err) {
|
|
callback(err as Error);
|
|
return;
|
|
}
|
|
|
|
const prefix = wroteFirst ? ',' : '[';
|
|
wroteFirst = true;
|
|
callback(null, prefix + serialized);
|
|
},
|
|
flush(callback: TransformCallback) {
|
|
callback(null, wroteFirst ? ']' : '[]');
|
|
},
|
|
});
|
|
};
|