bounded-queue helps solves the producer–consumer problem.
graph LR;
P(Producer);
B[bounded-queue];
C(Consumer);
P-- batched item -->B;
B-- batched item -->C;
style B fill:#99E,stroke:#333
The bounded-queue allows the producer and consumer to operate in
Imagine you have to read records from a database and write those to another database. A simple way to that move the records is to first read from database A and sequentially write each record to database B.
letbatchNr=0;letitems2produce=10;/** * Mockup database A, producer */constdbA={readRecord: async()=>{++batchNr;console.log("Producing batch #",batchNr);awaitnewPromise(resolve=>setTimeout(resolve,100));console.log("Produced batch #",batchNr);returnbatchNr;},moreRecordsAvailable: ()=>batchNr<items2produce}/** * Mockup database B, consumer */constdbB={asyncwriteRecord(batchNr){console.log("Consuming batch #",batchNr);awaitnewPromise(resolve=>setTimeout(resolve,100));console.log("Consumed batch #",batchNr);}}/** * Sequential conversion */asyncfunctionconvertDatabaseRecords(){while(dbA.moreRecordsAvailable()){constrecord=awaitdbA.readRecord();// ConsumerawaitdbB.writeRecord(record);// expensive async write (consume) operation}}(async()=>{console.time("no-queue");awaitconvertDatabaseRecords();console.timeEnd("no-queue");})();In the previous example, we either read from database A, or write to database B.
It would be faster if read from database A, while we write to database B, at the same time.
As dbA.readRecord() and dbB.readRecord()areasync` functions, there is no need to introduce threading to accomplish that.
The bounded-queue helps you with that. The following example uses bounded-queue, with a maximum of 3 queued records:
import{queue}from'bounded-queue';letbatchNr=0;letitems2produce=10;/** * Mockup database A, producer */constdbA={readRecord: async()=>{++batchNr;console.log("Producing batch #",batchNr);awaitnewPromise(resolve=>setTimeout(resolve,100));console.log("Produced batch #",batchNr);returnbatchNr;},moreRecordsAvailable: ()=>batchNr<items2produce}/** * Mockup database B, consumer */constdbB={asyncwriteRecord(batchNr){console.log("Consuming batch #",batchNr);awaitnewPromise(resolve=>setTimeout(resolve,100));console.log("Consumed batch #",batchNr);}}/** * Conversion using bounded-queue */asyncfunctionconvertDatabaseRecords(){awaitqueue(3,()=>{// ProducerreturndbA.moreRecordsAvailable() ? dbA.readRecord() : null;// expenive async read (produce) operation},record=>{// ConsumerreturndbB.writeRecord(record);// expensive async write (consume) operation});}(async()=>{console.time("bounded-queue");awaitconvertDatabaseRecords();console.timeEnd("bounded-queue");})();Using the bounded-queue, the conversion will complete in roughly half the time.
npm install bounded-queue- Parameters:
maxQueueSize(number): Maximum number of items that can be in the queue.producer(Producer): A function that produces items to be added to the queue.consumer(Consumer): A function that consumes items from the queue.
- Description: Initializes a new
BoundedQueueinstance with the specifiedmaxQueueSize,producer, andconsumer.
- Description: Returns the number of items currently in the queue.
- Returns: The number of items in the queue.
- Description: Initiates the asynchronous processing of items in the queue. It starts filling the queue with items produced by the
producerfunction and concurrently consumes items using theconsumerfunction. - Returns: A Promise that resolves when all items have been produced and consumed.
import{queue}from'your-module';// Create and run a BoundedQueue with a maximum queue size of 10queue(10,producer,consumer).then(()=>{console.log('Queue processing completed');}).catch((error)=>{console.error('Error during queue processing:',error);});queue<ItemType>(maxQueueSize: number, producer: Producer<ItemType>, consumer: Consumer<ItemType>): Promise<void>
- Parameters:
maxQueueSize(number): Maximum number of items that can be in the queue.producer(Producer): A function that produces items to be added to the queue.consumer(Consumer): A function that consumes items from the queue.
- Description: A convenience function for creating and running a
BoundedQueueinstance. It takes the same parameters as theBoundedQueueconstructor and returns a Promise that resolves when all items have been produced and consumed. - Returns: A Promise that resolves when all items have been produced and consumed. Resolving
nullindicated the end of the production.
import{queue}from'your-module';// Create and run a BoundedQueue with a maximum queue size of 10queue(10,producer,consumer).then(()=>{console.log('Queue processing completed');}).catch((error)=>{console.error('Error during queue processing:',error);});