Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Doing work concurrently

concurrentMap and concurrentUnorderedMap are map with several calls running at once. The first gives results in input order; the second gives them as they finish.

import { enumerate, sleep } from "@j50n/proc";

const jobs = [
  { name: "slow", ms: 300 },
  { name: "fast", ms: 100 },
  { name: "medium", ms: 200 },
];

async function work(job: { name: string; ms: number }) {
  await sleep(job.ms);
  return job.name;
}

// In input order: "fast" and "medium" wait behind "slow".
console.log(
  await enumerate(jobs).concurrentMap(work, { concurrency: 3 }).collect(),
);

// In the order they finish.
console.log(
  await enumerate(jobs).concurrentUnorderedMap(work, { concurrency: 3 })
    .collect(),
);
[ "slow", "fast", "medium" ]
[ "fast", "medium", "slow" ]

Which one to use

concurrentMap keeps the order, and pays for it: a slow item holds back the results after it, and while it does, no new call starts. With one slow item in every few, fewer calls run than you asked for. Use it when the output must line up with the input.

concurrentUnorderedMap starts a new call as soon as any finishes, so every slot stays busy. Use it when order doesn’t matter, or carry the input along in the result (return { file, size }) so you can tell which result is which.

Plain map makes one call at a time, even when the callback is async.

How many at once

The concurrency option sets how many calls may be running. It defaults to navigator.hardwareConcurrency, the number of CPUs. A fraction rounds up, and a value below 1 throws an Error when reading starts, not at the call.

Calls overlap only while they wait: on a timer, a file, the network, or a child process. The JavaScript inside the callbacks still runs one piece at a time, so CPU-heavy work in a callback gets no faster. Put that work in a child process, where it does run in parallel. The CPU count suits commands that keep one CPU busy each; work that mostly waits on the network can go higher.

Running several commands at once

import { enumerate, run } from "@j50n/proc";

const files = ["a.txt", "b.txt", "c.txt"];
for (const file of files) {
  await Deno.writeTextFile(file, `${file}\n`.repeat(10_000));
}

// Up to two gzip processes at a time.
const results = await enumerate(files)
  .concurrentMap(async (file) => {
    await run("gzip", "-k", file).collect(); // wait for it, reading its output
    const before = (await Deno.stat(file)).size;
    const after = (await Deno.stat(`${file}.gz`)).size;
    return `${file}: ${before} -> ${after} bytes`;
  }, { concurrency: 2 })
  .collect();

console.log(results);
[
  "a.txt: 60000 -> 137 bytes",
  "b.txt: 60000 -> 137 bytes",
  "c.txt: 60000 -> 137 bytes"
]

Each call must consume its command’s output, here with collect(), even when there is none to speak of. That is how the call waits for the command to finish, and how a failed command throws: as an ExitCodeError from the callback, which then surfaces as below. A command that writes to stdout and is never read blocks, and its call never finishes (see Key ideas).

When a call fails

import { enumerate, sleep } from "@j50n/proc";

try {
  await enumerate(["a", "b", "c"])
    .concurrentMap(async (name) => {
      await sleep(name === "c" ? 100 : 10);
      if (name === "b") throw new Error(`${name} failed`);
      console.log(`${name} done`);
      return name;
    }, { concurrency: 3 })
    .forEach((name) => console.log(`got ${name}`));
} catch (error) {
  if (error instanceof Error) console.log(`caught: ${error.message}`);
}

// The error didn't stop "c", which was already running.
await sleep(200);
a done
got a
caught: b failed
c done

An error thrown by the callback comes out of the consumer, so one try around the await catches it. concurrentMap throws it when the failed item’s turn comes, after the results before it; concurrentUnorderedMap throws it in the order it happened.

Nothing cancels the calls already running: “c” finished after the error was caught. Their results are dropped, but whatever they do (write a file, start a command) still happens. A consumer that stops early, with take or break, leaves running calls going in the same way.

To let every job finish and see each outcome, catch inside the callback and return the outcome as the result:

import { enumerate, ExitCodeError, run } from "@j50n/proc";

// Catch inside the callback to let every job finish and report each one.
const checks = [["true"], ["sh", "-c", "exit 3"], ["echo", "fine"]];

const results = await enumerate(checks)
  .concurrentMap(async ([cmd, ...args]) => {
    try {
      await run(cmd, ...args).collect();
      return `${cmd}: ok`;
    } catch (error) {
      if (error instanceof ExitCodeError) return `${cmd}: exit ${error.code}`;
      throw error;
    }
  })
  .collect();

console.log(results);
[ "true: ok", "sh: exit 3", "echo: ok" ]

For a larger worked example, see Running jobs in parallel.