Refactor the task merge-queue: batch, dedupe by key, priority-order, flush
Node · Node · advanced · modification
Refactors MergeQueue from v1's interval-only coalescing to a v2 that also orders each batch by priority, takes a pluggable comparator and tracks flush metrics. Dedup still collapses jobs that share a key — merging payloads and raising priority — and the batch is now sorted (higher priority first, stable on arrival order) before onFlush. The drain-aware flushing and the field-by-field payload merge are both covered. Verified against a seeded stream: keys collapse correctly, priorities order as expected and the batch flushes on time.
MergeQueue sits in the hot ingest path of a task pipeline; enqueue is called synchronously on the event loop as jobs stream in. Under load a single batch fills to tens of thousands of distinct keys before it flushes. The v1 it replaces stored the pending batch in a Map keyed by job key; v2 switched to an array to make priority sorting straightforward.
Requirements
- Coalesce a high-volume stream of jobs into periodic batches: producers call enqueue(jobs) with one or many jobs, and each distinct job `key` survives at most once per flush.
- When a later job repeats a key already in the batch, merge it into the existing entry: raise the entry's priority to the strongest seen and combine payloads field-by-field.
- Flush on a fixed interval, or immediately when a caller invokes flushNow(); the flushed batch is ordered by priority (higher first), ties broken by arrival order so the ordering is stable.
- The queue must be drain-aware: work enqueued while an onFlush callback is running accumulates into the next batch, and the timer is only rearmed after the in-flight batch is handed off.
- The batch can hold tens of thousands of distinct keys, and enqueue runs on the event loop in the ingest path, so accepting and deduplicating N jobs must be O(N) total with expected O(1) lookup per job — the flush's O(K log K) priority sort on its K-key batch is separate and required.
Files touched
- src/mergeQueue.js
--- src/mergeQueue.js
/**
- * MergeQueue (v1)
+ * MergeQueue (v2)
* ----------------
- * Coalesces a stream of jobs into periodic, de-duplicated batches. Producers
- * call `enqueue(jobs)` with one or many jobs; jobs that share a `key` collapse
- * into a single pending entry so a downstream consumer processes each key at
- * most once per flush. The batch is handed to `onFlush` on a fixed interval.
+ * Coalesces a high-volume stream of jobs into periodic, de-duplicated batches.
+ *
+ * Producers call `enqueue(jobs)` with one or many jobs. Jobs that share a
+ * `key` are merged into a single pending entry — payloads combined, priority
+ * raised to the strongest seen — so a downstream consumer processes each key at
+ * most once per flush. On a fixed interval, or via an explicit `flushNow()`,
+ * the queue orders the batch by priority and hands it to `onFlush`.
+ *
+ * v2 adds priority ordering, a pluggable comparator and flush metrics on top of
+ * v1's interval coalescing. The queue is drain-aware: while a flush callback is
+ * running, new work keeps accumulating for the next batch, and the flush timer
+ * is only rearmed once the in-flight batch has been handed off.
*/
const DEFAULT_FLUSH_INTERVAL_MS = 50;
+function normalizePriority(value) {
+ const n = Number(value);
+ return Number.isFinite(n) ? n : 0;
+}