You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

适配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_PARALLELISM limit 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 push method works with both single tasks and arrays, just like async.queue

内容的提问来源于stack exchange,提问作者Moshe Shmukler

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 09:32:54