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.