框架
版本
Debouncer API 参考
Throttler API 参考
速率限制器 API 参考
队列 API 参考
批处理器 API 参考
批处理器示例

异步队列指南

注意:队列指南 中的所有核心队列概念也适用于 AsyncQueuer。AsyncQueuer 通过并发(同时处理多个任务)和强大的错误处理等高级功能扩展了这些概念。如果你是队列新手,请从队列指南开始学习 FIFO/LIFO、优先级、过期、拒绝和队列管理。本指南重点介绍 AsyncQueuer 在异步和并发任务处理方面的独特之处和强大功能。

虽然 Queuer 提供了具有计时控制的同步队列,但 AsyncQueuer 专为处理并发异步操作而设计。它实现了传统上称为“任务池”或“工作池”的模式,允许同时处理多个操作,同时保持对并发和计时的控制。该实现主要复制自 Swimmer,这是 Tanner 最初的任务池实用程序,自 2017 年以来一直为 JavaScript 社区服务。

异步队列概念

异步队列通过添加并发处理功能扩展了基本队列概念。异步队列器可以同时处理多个项目,而不是一次处理一个项目,同时仍保持对执行的顺序和控制。这在处理 I/O 操作、网络请求或任何大部分时间都在等待而不是消耗 CPU 的任务时特别有用。

异步队列可视化

text
异步队列 (并发数: 2, 等待: 2 个滴答)
时间轴: [每个滴答 1 秒]
调用:        ⬇️  ⬇️  ⬇️  ⬇️     ⬇️  ⬇️     ⬇️
队列:       [ABC]   [C]    [CDE]    [E]    []
活动:      [A,B]   [B,C]  [C,D]    [D,E]  [E]
已完成:    -       A      B        C      D,E
             [=================================================================]
             ^ 与常规队列不同,可以并发处理多个项目

             [项目排队]   [一次处理 2 个]   [完成]
              繁忙时         处理间隔         所有项目
异步队列 (并发数: 2, 等待: 2 个滴答)
时间轴: [每个滴答 1 秒]
调用:        ⬇️  ⬇️  ⬇️  ⬇️     ⬇️  ⬇️     ⬇️
队列:       [ABC]   [C]    [CDE]    [E]    []
活动:      [A,B]   [B,C]  [C,D]    [D,E]  [E]
已完成:    -       A      B        C      D,E
             [=================================================================]
             ^ 与常规队列不同,可以并发处理多个项目

             [项目排队]   [一次处理 2 个]   [完成]
              繁忙时         处理间隔         所有项目

何时使用异步队列

当你需要执行以下操作时,异步队列特别有效:

  • 并发处理多个异步操作
  • 控制同时操作的数量
  • 通过适当的错���处理来处理基于 Promise 的任务
  • 在最大化吞吐量的同时保持顺序
  • 处理可以并行运行的后台任务

何时不使用异步队列

AsyncQueuer 非常通用,可以在许多情况下使用。如果你不需要并发处理,请改用队列。如果你不需要队列中的所有执行都通过,请改用节流

如果你想将操作组合在一起,请改用批处理

TanStack Pacer 中的异步队列

TanStack Pacer 通过简单的 asyncQueue 函数和更强大的 AsyncQueuer 类提供异步队列功能。所有队列类型和排序策略(FIFO、LIFO、优先级等)都与核心队列指南中的一样受支持。

asyncQueue 的基本用法

asyncQueue 函数提供了一种创建始终运行的异步队列的简单方法:

ts
import { asyncQueue } from '@tanstack/pacer'

// 创建一个最多可并发处理 2 个项目的队列
const processItems = asyncQueue(
  async (item: number) => {
    // 异步处理每个项目
    const result = await fetchData(item)
    return result
  },
  {
    concurrency: 2,
    onItemsChange: (queuer) => {
      console.log('活动任务:', queuer.peekActiveItems().length)
    }
  }
)

// 添加要处理的项目
processItems(1)
processItems(2)
import { asyncQueue } from '@tanstack/pacer'

// 创建一个最多可并发处理 2 个项目的队列
const processItems = asyncQueue(
  async (item: number) => {
    // 异步处理每个项目
    const result = await fetchData(item)
    return result
  },
  {
    concurrency: 2,
    onItemsChange: (queuer) => {
      console.log('活动任务:', queuer.peekActiveItems().length)
    }
  }
)

// 添加要处理的项目
processItems(1)
processItems(2)

要对队列进行更多控制,请直接使用 AsyncQueuer 类。

AsyncQueuer 类的高级用法

AsyncQueuer 类提供了对异步队列行为的完全控制,包括所有核心队列功能以及:

  • **并发:**一次处理多个项目(可通过 concurrency 配置)
  • **异步错误处理:**每个任务和全局错误回调,并控制错误传播
  • **活动和待处理任务跟踪:**监控哪些任务正在运行以及哪些任务已排队
  • 异步特定回调:onSuccessonErroronSettled 等。
ts
import { AsyncQueuer } from '@tanstack/pacer'

const queue = new AsyncQueuer(
  async (item: number) => {
    // 异步处理每个项目
    const result = await fetchData(item)
    return result
  },
  {
    concurrency: 2, // 一次处理 2 个项目
    wait: 1000,     // 在开始新项目之间等待 1 秒
    started: true   // 立即开始处理
  }
)

// 通过选项添加错误和成功处理程序
queue.setOptions({
  onError: (error, queuer) => {
    console.error('任务失败:', error)
    // 你可以在此处访问队列状态
    console.log('错误计数:', queuer.getErrorCount())
  },
  onSuccess: (result, queuer) => {
    console.log('任务已完成:', result)
    // 你可以在此处访问队列状态
    console.log('成功计数:', queuer.getSuccessCount())
  },
  onSettled: (queuer) => {
    // 每次执行(成功或失败)后调用
    console.log('总共已完成:', queuer.getSettledCount())
  }
})

// 添加要处理的项目
queue.addItem(1)
queue.addItem(2)
import { AsyncQueuer } from '@tanstack/pacer'

const queue = new AsyncQueuer(
  async (item: number) => {
    // 异步处理每个项目
    const result = await fetchData(item)
    return result
  },
  {
    concurrency: 2, // 一次处理 2 个项目
    wait: 1000,     // 在开始新项目之间等待 1 秒
    started: true   // 立即开始处理
  }
)

// 通过选项添加错误和成功处理程序
queue.setOptions({
  onError: (error, queuer) => {
    console.error('任务失败:', error)
    // 你可以在此处访问队列状态
    console.log('错误计数:', queuer.getErrorCount())
  },
  onSuccess: (result, queuer) => {
    console.log('任务已完成:', result)
    // 你可以在此处访问队列状态
    console.log('成功计数:', queuer.getSuccessCount())
  },
  onSettled: (queuer) => {
    // 每次执行(成功或失败)后调用
    console.log('总共已完成:', queuer.getSettledCount())
  }
})

// 添加要处理的项目
queue.addItem(1)
queue.addItem(2)

异步特定功能

所有队列类型和排序策略(FIFO、LIFO、优先级等)都受支持——有关详细信息,请参阅队列指南。AsyncQueuer 添加了:

  • **并发:**可以一次处理多个项目,由 concurrency 选项控制(可以是动态的)。
  • **异步错误处理:**使用 onErroronSuccessonSettled 进行强大的错误和结果跟踪。
  • **活动和待处理任务跟踪:**使用 peekActiveItems()peekPendingItems() 监控队列状态。
  • **异步过期和拒绝:**项目可以像核心队列指南中一样过期或被拒绝,但具有异步特定的回调。

示例:优先级异步队列

ts
const priorityQueue = new AsyncQueuer(
  async (item: { value: string; priority: number }) => {
    // 异步处理每个项目
    return await processTask(item.value)
  },
  {
    concurrency: 2,
    getPriority: (item) => item.priority // 数字越大优先级越高
  }
)

priorityQueue.addItem({ value: 'low', priority: 1 })
priorityQueue.addItem({ value: 'high', priority: 3 })
priorityQueue.addItem({ value: 'medium', priority: 2 })
// 处理顺序:high 和 medium 并发,然后是 low
const priorityQueue = new AsyncQueuer(
  async (item: { value: string; priority: number }) => {
    // 异步处理每个项目
    return await processTask(item.value)
  },
  {
    concurrency: 2,
    getPriority: (item) => item.priority // 数字越大优先级越高
  }
)

priorityQueue.addItem({ value: 'low', priority: 1 })
priorityQueue.addItem({ value: 'high', priority: 3 })
priorityQueue.addItem({ value: 'medium', priority: 2 })
// 处理顺序:high 和 medium 并发,然后是 low

示例:错误处理

ts
const queue = new AsyncQueuer(
  async (item: number) => {
    // 异步处理每个项目
    if (item < 0) throw new Error('负数项目')
    return await processTask(item)
  },
  {
    onError: (error, queuer) => {
      console.error('任务失败:', error)
      // 你可以在此处访问队列状态
      console.log('错误计数:', queuer.getErrorCount())
    },
    throwOnError: true, // 即使有 onError 处理程序也会抛出错误
    onSuccess: (result, queuer) => {
      console.log('任务成功:', result)
      // 你可以在此处访问队列状态
      console.log('成功计数:', queuer.getSuccessCount())
    },
    onSettled: (queuer) => {
      // 每次执行(成功或失败)后调用
      console.log('总共已完成:', queuer.getSettledCount())
    }
  }
)

queue.addItem(-1) // 将触发错误处理
queue.addItem(2)
const queue = new AsyncQueuer(
  async (item: number) => {
    // 异步处理每个项目
    if (item < 0) throw new Error('负数项目')
    return await processTask(item)
  },
  {
    onError: (error, queuer) => {
      console.error('任务失败:', error)
      // 你可以在此处访问队列状态
      console.log('错误计数:', queuer.getErrorCount())
    },
    throwOnError: true, // 即使有 onError 处理程序也会抛出错误
    onSuccess: (result, queuer) => {
      console.log('任务成功:', result)
      // 你可以在此处访问队列状态
      console.log('成功计数:', queuer.getSuccessCount())
    },
    onSettled: (queuer) => {
      // 每次执行(成功或失败)后调用
      console.log('总共已完成:', queuer.getSettledCount())
    }
  }
)

queue.addItem(-1) // 将触发错误处理
queue.addItem(2)

示例:动态并发

ts
const queue = new AsyncQueuer(
  async (item: number) => {
    // 异步处理每个项目
    return await processTask(item)
  },
  {
    // 基于系统负载的动态并发
    concurrency: (queuer) => {
      return Math.max(1, 4 - queuer.peekActiveItems().length)
    },
    // 基于队列大小的动态等待时间
    wait: (queuer) => {
      return queuer.getSize() > 10 ? 2000 : 1000
    }
  }
)
const queue = new AsyncQueuer(
  async (item: number) => {
    // 异步处理每个项目
    return await processTask(item)
  },
  {
    // 基于系统负载的动态并发
    concurrency: (queuer) => {
      return Math.max(1, 4 - queuer.peekActiveItems().length)
    },
    // 基于队列大小的动态等待时间
    wait: (queuer) => {
      return queuer.getSize() > 10 ? 2000 : 1000
    }
  }
)

队列管理和监控

AsyncQueuer 提供了核心队列指南中的所有队列管理和监控方法,以及异步特定的方法:

  • peekActiveItems() — 当前正在处理的项目
  • peekPendingItems() — 等待处理的项目
  • getSuccessCount()getErrorCount()getSettledCount() — 执行统计信息
  • start()stop()clear()reset() 等。

有关队列管理概念的更多信息,请参阅队列指南

任务过期和拒绝

AsyncQueuer 支持过期和拒绝,就像核心队列器一样:

  • 使用 expirationDurationgetIsExpiredonExpire 处理任务过期
  • 使用 maxSizeonReject 处理队列溢出

有关详细信息和示例,请参阅队列指南

框架适配器

每个框架适配器都围绕异步队列器类构建了方便的钩子和函数。诸如 useAsyncQueueruseAsyncQueuedState 之类的钩子是小型包装器,可以减少某些常见用例中你自己代码所需的样板代码。