p-map
Map over promises concurrently
Useful when you need to run promise-returning & async functions multiple times with different inputs concurrently.
This is different from Promise.all() in that you can control the concurrency and also decide whether or not to stop iterating when there's an error.
Install
npm install p-map
Usage
import pMap from 'p-map';
import got from 'got';
const sites = [
getWebsiteFromUsername('sindresorhus'), //=> Promise
'https://avajs.dev',
'https://github.com'
];
const mapper = async site => {
const {requestUrl} = await got.head(site);
return requestUrl;
};
const result = await pMap(sites, mapper, {concurrency: 2});
console.log(result);
//=> ['https://sindresorhus.com/', 'https://avajs.dev/', 'https://github.com/']
API
Except for pMapIterable, each function is a version of a Promise method that maps each element and limits the concurrency:
| Promise | p-map |
| --- | --- |
| Promise.all() | pMap |
| Promise.allSettled() | pMapSettled |
| Promise.allKeyed() | pMapKeyed |
| Promise.allSettledKeyed() | pMapSettledKeyed |
Promise.allKeyed() and Promise.allSettledKeyed() are a Stage 3 proposal, so they may not be in your JavaScript engine yet. pMapKeyed and pMapSettledKeyed work without them.
pMap(input, mapper, options?)
Returns a Promise that is fulfilled when all promises in input and the ones returned from mapper are fulfilled, or rejects if any of the promises reject. The fulfilled value is an Array of the fulfilled values returned from mapper in input order.
pMapSettled(input, mapper, options?)
A version of pMap for Promise.allSettled(): an error from mapper, or an element of input that rejects, does not reject the returned Promise.
Returns a Promise that is fulfilled when all promises in input and the ones returned from mapper are settled. The fulfilled value is an Array of objects in input order, the same as from Promise.allSettled(): {status: 'fulfilled', value} for each value returned from mapper, and {status: 'rejected', reason} for each error from mapper or from input.
import {pMapSettled} from 'p-map';
const results = await pMapSettled(urls, async url => {
const response = await fetch(url);
return response.json();
}, {concurrency: 4});
for (const result of results) {
if (result.status === 'fulfilled') {
console.log(result.value);
} else {
console.error(result.reason);
}
}
An element of input that rejects gets a {status: 'rejected', reason} result, and mapper is not called for it. A pMapSkip result is left out, and pMapStop stops like in pMap.
The returned Promise still rejects when iterating input throws, like Promise.allSettled(), and when the signal aborts.
It takes the same options as pMap, except stopOnError.
[!NOTE]
Promises that are already ininputare already running, soconcurrencycannot limit them. With a finiteconcurrency, a promise that rejects beforepMapSettledreaches it is reported as an unhandled rejection. For promises you already have, usePromise.allSettled(). WithpMapSettled, give the values ininputand start the async work inmapper.
pMapKeyed(input, mapper, options?)
A version of pMap for the values of an object or a Map, like Promise.allKeyed() is for Promise.all().
Returns a Promise that is fulfilled with an object with the same keys, in the same order, and the values returned from mapper, or rejects if any of the values or the mappers reject. The object has a null prototype, like the one from Promise.allKeyed(). For a Map, it is a new Map.
import {pMapKeyed} from 'p-map';
const {user, posts} = await pMapKeyed({
user: 'https://example.com/api/user',
posts: 'https://example.com/api/posts',
}, async url => {
const response = await fetch(url);
return response.json();
}, {concurrency: 2});
input is an object or a Map. For an object, it uses the own enumerable keys, strings and symbols, like Promise.allKeyed(). For a Map, it uses the entries, and a key can be any value. A Map from another realm, like from node:vm or an iframe, is used as an object, so its entries are not mapped. All the values are read right away, and each value is await'd before mapper is called for it, so a value may be a Promise. With TypeScript, give an input of type any, like the result of response.json(), a specific type, like as Record, as the result type is not useful otherwise.
mapper(value, key) is called with each value and its key. Expected to return a Promise or value.
A value that rejects gets a handler right away, so it is never an unhandled rejection, even when concurrency has not reached it yet.
A pMapSkip result leaves out the key, and pMapStop keeps only the keys before it.
It takes the same options as pMap.
pMapSettledKeyed(input, mapper, options?)
A version of pMapKeyed for Promise.allSettledKeyed(): an error from mapper, or a value of input that rejects, does not reject the returned Promise.
Returns a Promise that is fulfilled with an object with the same keys, in the same order, when all the values and mappers are settled. Each value is {status: 'fulfilled', value} with the value returned from mapper, or {status: 'rejected', reason} with the error from mapper or from input, the same as from Promise.allSettledKeyed(). The object has a null prototype. For a Map, it is a new Map.
import {pMapSettledKeyed} from 'p-map';
const {user, posts} = await pMapSettledKeyed({
user: 'https://example.com/api/user',
posts: 'https://example.com/api/posts',
}, async url => {
const response = await fetch(url);
return response.json();
}, {concurrency: 2});
if (user.status === 'fulfilled') {
console.log(user.value);
}
if (posts.status === 'rejected') {
console.error(posts.reason);
}
It takes the same input and mapper as pMapKeyed, and the same options as pMap, except stopOnError.
A value of input that rejects gets a {status: 'rejected', reason} result, and mapper is not called for it. Each value gets a handler right away, so it is never an unhandled rejection. A pMapSkip result leaves out the key, and pMapStop keeps only the keys before it.
The returned Promise still rejects when reading input throws, and when the signal aborts.
pMapIterable(input, mapper, options?)
Returns an async iterable that streams each return value from mapper in order, or as soon as each is ready when preserveOrder is false.
import {pMapIterable} from 'p-map';
// Multiple posts are fetched concurrently, with limited concurrency and backpressure
for await (const post of pMapIterable(postIds, getPostMetadata, {concurrency: 8})) {
console.log(post);
};
[!NOTE]
Like an async generator, areturn()orthrow()on the iterator that comes while anext()call is still pending waits for thatnext()to settle. Abreakin afor awaitloop is not affected, as the loop only stops between results. Ifinputcan block innext()for good, like a queue that waits for an item, stop it at the source, for example with anAbortSignalthat it listens to.
input
Type: AsyncIterable
Synchronous or asynchronous iterable that is iterated over concurrently, calling the mapper function for each element. Each iterated item is await'd before the mapper is invoked so the iterable may return a Promise that resolves to an item.
Asynchronous iterables (different from synchronous iterables that return Promise that resolves to an item) can be used when the next item may not be ready without waiting for an asynchronous process to complete and/or the end of the iterable may be reached after the asynchronous process completes. For example, reading from a remote queue when the queue has reached empty, or reading lines from a stream.
For pMapKeyed and pMapSettledKeyed, see pMapKeyed.
mapper(element, index)
Type: Function
Expected to return a Promise or value.
For pMapKeyed and pMapSettledKeyed, see pMapKeyed.
options
Type: object
##### concurrency
Type: number (Integer)\
Default: Infinity\
Minimum: 1
Number of concurrently pending promises returned by mapper.
[!NOTE]
pMapandpMapSettledcan drop items that pulls frominputalready had in flight: when they stop early, because a mapper returnspMapStop, a mapper throws (only forpMap), or thesignalaborts, and, withconcurrencyabove 1, when one pull reportsdonewhile another still gets an item. That is invisible for a re-readable input, like an array or a string, but a source that cannot be re-read, such as a queue, loses those items. Withconcurrency: 1, only an abort can drop an item.
##### backpressure
Only for pMapIterable
Type: number (Integer)\
Default: options.concurrency\
Minimum: options.concurrency
Maximum number of elements taken from input that the consumer of the async iterable has not collected yet. That counts the elements being read, waiting for mapper, being mapped, and the results waiting to be collected. Calls to mapper will be limited so that there is never too much backpressure.
Useful whenever you are consuming the iterable slower than what the mapper function can produce concurrently. For example, to avoid making an overwhelming number of HTTP requests if you are saving each of the results to a database.
A backpressure above concurrency also lets pMapIterable read ahead: up to backpressure - concurrency elements, but no more than concurrency, are read from input while the mappers run, so a slow input does not leave a mapper waiting for its next element. Useful when reading input is slow, for example when it is another pMapIterable. With concurrency: 1, backpressure: 2 overlaps reading each element with mapping the one before it. When the iteration ends early, the elements read ahead are not mapped, so a source that cannot be re-read, such as a queue, loses them.
A result of pMapSkip is collected by the consumer stepping over it, so it counts towards the limit until then. A mapper that skips most elements therefore still reads no further ahead than backpressure allows.
[!NOTE]
backpressuredefaults toconcurrency, which defaults toInfinity, so by default nothing bounds how farinputis read ahead of the consumer. A source that never ends is then read without limit until memory runs out. Give such a source a finiteconcurrency. That also boundsbackpressure, since it defaults toconcurrency.
##### preserveOrder
Only for pMapIterable
Type: boolean\
Default: true
Whether to yield the results in input order.
When false, each result is yielded as soon as it is ready. A slow element then no longer holds back the results after it, and those results no longer fill up backpressure while they wait, so mapper keeps running at full concurrency. Use it when the order of the results does not matter.
An error from mapper is then also thrown as soon as it happens, instead of after the results before it. The index passed to mapper is still the position in input.
import {pMapIterable} from 'p-map';
// Each page is saved as soon as it is downloaded, without waiting for slower pages before it
for await (const page of pMapIterable(urls, downloadPage, {concurrency: 8, preserveOrder: false})) {
await savePage(page);
}
##### stopOnError
Only for pMap and pMapKeyed
Type: boolean\
Default: true
When true, the first mapper rejection will be rejected back to the consumer.
When false, instead of stopping when a promise rejects, it will wait for all the promises to settle and then reject with an AggregateError containing all the errors from the rejected promises.
The errors are listed in the order of the input, not in the order the promises happened to settle in, so the same input always produces the same AggregateError.
Caveat: When true, any already-started async mappers will continue to run until they resolve or reject. In the case of infinite concurrency with sync iterables, all mappers are invoked on startup and will continue after the first rejection. Use the signal option for abort control.
##### signal
Not for pMapIterable
Type: AbortSignal
You can abort the promises using AbortController.
import pMap from 'p-map';
import delay from 'delay';
const abortController = new AbortController();
setTimeout(() => {
abortController.abort();
}, 500);
const mapper = async value => value;
await pMap([delay(1000), delay(1000)], mapper, {signal: abortController.signal});
// Throws AbortError (DOMException) after 500 ms.
[!NOTE]
An abort that comes from a timer or an I/O callback, like asetTimeout, cannot take effect while the mapping is inside a stretch of work that never waits on a timer or I/O. Draining a large input with mappers that just return a value is one such stretch, and there the run finishes first. This is the same forPromise.alland for a plain loop, and it goes away as soon as the mapper or the source waits on a timer or I/O.
pMapSkip
Return this value from a mapper function to skip including the value in the returned array, or, for pMapKeyed and pMapSettledKeyed, to leave out the key.
import pMap, {pMapSkip} from 'p-map';
import got from 'got';
const sites = [
getWebsiteFromUsername('sindresorhus'), //=> Promise
'https://avajs.dev',
'https://example.invalid',
'https://github.com'
];
const mapper = async site => {
try {
const {requestUrl} = await got.head(site);
return requestUrl;
} catch {
return pMapSkip;
}
};
const result = await pMap(sites, mapper, {concurrency: 2});
console.log(result);
//=> ['https://sindresorhus.com/', 'https://avajs.dev/', 'https://github.com/']
pMapStop
Return this value from a mapper function to stop iterating, like break in a loop.
pMap resolves with the results of the elements before it, in order, without waiting for the elements after it. Like in a loop, an error from an earlier element still rejects. An error from a later element rejects only if it happens before the stop. With stopOnError: false, only the errors from earlier elements are collected.
pMapKeyed stops the same way, and keeps only the keys before it. pMapSettled and pMapSettledKeyed also stop the same way, but an error from an element before the stop gives a {status: 'rejected', reason} result instead of a rejection.
pMapIterable ends the stream where the stop would have been yielded: after the elements before it with preserveOrder, or after the results that were ready before it without.
pMap, pMapSettled, and pMapIterable then read no more elements from input, and close input. pMapKeyed and pMapSettledKeyed have read all the values already. Mappers that are still running are not canceled, and their results are dropped. The element that returns pMapStop adds no value to the result.
import pMap, {pMapStop} from 'p-map';
function * pageNumbers() {
for (let pageNumber = 1; ; pageNumber++) {
yield pageNumber;
}
}
const mapper = async pageNumber => {
const page = await getPage(pageNumber);
return page.items.length === 0 ? pMapStop : page;
};
const pages = await pMap(pageNumbers(), mapper, {concurrency: 4});
console.log(pages);
//=> Every page before the first empty page, in order
Recipes
Rate limiting
This package controls how many mapper promises run concurrently. To limit how often the mapper starts, compose it with a rate limiter like p-throttle:
import pThrottle from 'p-throttle';
import pMap from 'p-map';
const throttle = pThrottle({
limit: 1,
interval: 1000,
strict: true,
});
const result = await pMap(input, throttle(mapper), {concurrency: 2});
For more advanced scheduling, use p-queue.
Dynamic concurrency
To change the concurrency while mapping, compose the mapper with p-limit, whose concurrency can be changed at any time:
import pLimit from 'p-limit';
import pMap from 'p-map';
const limit = pLimit(2);
// For example, raise it when the server is less busy
setTimeout(() => {
limit.concurrency = 8;
}, 10_000);
try {
// The concurrency option is the highest value limit.concurrency can have
await pMap(input, (element, index) => limit(mapper, element, index), {concurrency: 8});
} finally {
// Do not start the queued elements after a stop or an error
limit.clearQueue();
}
Related
- p-all - Run promise-returning & async functions concurrently with optional limited concurrency
- p-filter - Filter promises concurrently
- p-times - Run promise-returning & async functions a specific number of times concurrently
- p-props - Like
Promise.all()but forMapandObject - p-map-series - Map over promises serially
- More…