适配Async/Await的async.queue等价实现及代码现代化咨询
Replace async.queue with Async/Await-Compatible Queue Implementation
Got it, let's break this down. You've modernized your insertCustomer function to use Async/Await, which means the old async.queue (built for callback-based tasks) won't play nicely anymore. Here's a clean, equivalent Promise-based queue implementation that matches the behavior of your original code:
Step 1: Implement a Custom Async Queue Class
This class handles concurrent task execution just like async.queue, but supports Async/Await task functions:
class AsyncQueue { constructor(taskHandler, concurrency) { this.taskHandler = taskHandler; // Your Async/Await-enabled insertCustomer function this.concurrency = concurrency; // Same DATABASE_PARALLELISM value as before this.queue = []; this.activeTasks = 0; this.onDrain = null; // Callback for when all tasks are done } // Push single or multiple tasks to the queue push(tasks) { const tasksToAdd = Array.isArray(tasks) ? tasks : [tasks]; this.queue.push(...tasksToAdd); this._processQueue(); } async _processQueue() { // Keep processing tasks while we have queued items and capacity while (this.queue.length > 0 && this.activeTasks < this.concurrency) { const task = this.queue.shift(); this.activeTasks++; try { // Execute your Async/Await task await this.taskHandler(task); } catch (error) { // Add error handling here (e.g., log failures) logger.error(`Failed to process customer: ${error.message}`, error); } finally { this.activeTasks--; // Continue processing remaining tasks this._processQueue(); // Trigger drain callback when queue is empty and no active tasks left if (this.queue.length === 0 && this.activeTasks === 0 && this.onDrain) { this.onDrain(); } } } } }
Step 2: Replace Your Original Queue Usage
Swap out the old async.queue setup with this new class—its API will feel familiar:
// Create your queue with the same concurrency limit const customerQueue = new AsyncQueue(insertCustomer, DATABASE_PARALLELISM); // Set up the drain callback just like before customerQueue.onDrain = function() { logger.info('all customers loaded'); airportCodeMappingQueue.push(airportCodeMappings); }; // Push tasks exactly as you did before (single item or array) // customerQueue.push(singleCustomer); // customerQueue.push(batchOfCustomers);
Key Details That Match Original Behavior
- Concurrency Control: Maintains the same
DATABASE_PARALLELISMlimit as your original queue - Drain Callback: Triggers only when all queued tasks have completed (including any in-flight ones)
- Error Resilience: A failed task won't block the rest of the queue (you can adjust error handling to your needs)
- API Compatibility: The
pushmethod works with both single tasks and arrays, just likeasync.queue
内容的提问来源于stack exchange,提问作者Moshe Shmukler
相关产品推荐
相关产品推荐

