Create a new IterableQueueMapper, which uses IterableMapper underneath, and exposes a
queue interface for adding items that are not exposed via an iterator.
Function called for every enqueued item. Returns a Promise or value.
IterableQueueMapper options
Indicate that no more items will be enqueued.
Call after the last awaited enqueue. Finish consuming the async iterator to wait for mapped results; this method does not wait for processing to finish.
Add an item to the queue, waiting until it can be accepted. Resolves on acceptance, rather than completion of the mapper for this item. Await each enqueue while consuming results concurrently for producer backpressure.
Element to add
Used by the iterator returned from [Symbol.asyncIterator] Called every time an item is needed
Iterator result
Accepts queue items via
enqueueand calls themapperon them with specifiedconcurrency, storing themapperresult in a queue ofmaxUnreadsize, before being iterated / read by the caller. Theenqueuemethod will block if the queue is full, until an item is read.Also exported as
MappingQueue, with the same constructor and instance types. Each input uses the same callback supplied at construction; results are exposed through the async iterator as mapping completes.Remarks
Typical Use Cases
IterableQueueMapperSimple/WorkerQueue)maxUnreadis reachedError Handling
The mapper should ideally handle all errors internally to enable error handling closest to where they occur. However, if errors do escape the mapper:
When
stopOnMapperErroris true (default):AsyncIterator's next() callWhen
stopOnMapperErroris false:AggregateErrorafter all items completeUsage
await enqueue()methodawait enqueue()method will block until a slot is available, if queue is fulldone()after the last awaited enqueue, then finish consuming the iteratorenqueue()confirms acceptance of an input;done()closes input without waiting for work to finishSee
IterableMapper for underlying mapper implementation and examples of combined usage