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

Files touched

--- 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;
+}

Review this PR

Node practice