Skip to content

GitHub 源码lib/src/utils/fast_async_queue/

概览

FastAsyncQueueFIFO 顺序串行执行异步任务(AsyncJob),适合「必须一条条做完」的场景(上传队列、串行 API、离线同步)。支持手动 start()自动启动 工厂、任务 labelJobInfo 跟踪、失败后在任务内调用 retry() 的重试,以及 QueueEvent 监听。

创建方式行为概要典型场景
FastAsyncQueue()入队后不执行,需 await start()批量攒任务后统一跑
FastAsyncQueue.autoStart()每次 addJob 后若空闲则自动 start()来一条处理一条、仍保证串行

FastAsyncQueue

队列内部为链表;同一时刻 最多一个 任务处于 running。任务在 catch 中调用 retry() 可将当前任务标记为 pendingRetry,在本轮 start() 循环内再次执行;retryTime允许的重试次数(默认 1),-1 表示不限次数。

基础使用示例

手动启动:先 addJobstart()

dart
final queue = FastAsyncQueue();

queue.addJob(
  () => Future.delayed(const Duration(seconds: 1), () => uploadChunk(1)),
  label: 'chunk-1',
  description: '第一片',
);
queue.addJob(
  () => Future.delayed(const Duration(seconds: 1), () => uploadChunk(2)),
  label: 'chunk-2',
);

await queue.start(); // 按入队顺序串行执行
dart
final queue = FastAsyncQueue.autoStart();

queue.addJob(() => syncProfile()); // 若队列空闲,立即开始处理
queue.addJob(() => syncSettings()); // 上一任务未完成则排队
dart
final queue = FastAsyncQueue.autoStart();

queue.addJob(
  () async {
    try {
      await requestWithTransientError();
    } catch (_) {
      queue.retry(); // 在 retryTime 限额内再次执行当前任务
    }
  },
  label: 'sync-order',
  retryTime: 3,
);

队列监听

通过 addQueueListener 订阅 QueueEvent(含 typecurrentQueueSizejobLabel、发生时间)。

dart
final queue = FastAsyncQueue();
queue.addQueueListener((event) {
  debugPrint('$event'); // QueueEvent [currentQueueSize: ..., type: ..., ...]
});

常见 QueueEventTypequeueStart / beforeJob / afterJob / queueEndnewJobAddedretryJob / retryLimitReachedqueueClosedqueueStoppedviolateAddWhenClosed

完整 API 参考


构造函数

dart
final manual = FastAsyncQueue();
final autoRun = FastAsyncQueue.autoStart();
说明
FastAsyncQueue()普通队列;任务入队后等待 start()
FastAsyncQueue.autoStart()工厂;addJob 后若未在运行则自动调用 start()

FastAsyncQueue.addJob

向队尾追加异步任务。队列已 close()静默拒绝 并发出 violateAddWhenClosedlabel 重复时抛出 DuplicatedLabelException。未传 label 时使用 ISO8601 时间字符串。

dart
queue.addJob(
  () async => doWork(),
  label: 'job-a',
  description: '可选说明',
  retryTime: 1,
);
参数类型必填说明
jobAsyncJob无参、返回 Future 的异步函数。
labelString?任务唯一标识;用于 getJobInfo 与事件中的 jobLabel
descriptionString?可读描述,写入 JobInfo
retryTimeint允许的重试次数,默认 1-1 为无限重试。

FastAsyncQueue.addJobThrow

addJob 相同,但队列已关闭时抛出 ClosedQueueException,便于调用方显式处理。

dart
queue.addJobThrow(() async => doWork(), label: 'critical');
参数类型必填说明
jobAsyncJob要执行的任务。
labelString?addJob
descriptionString?addJob
retryTimeintaddJob

FastAsyncQueue.start

若队列非空且当前未在运行,则依次执行队首任务直到队列为空、被 stop() 打断或循环因 stop 清空。队列为空、已在运行或已关闭时 立即返回

dart
await queue.start();
返回值类型说明
Future<void>本轮串行处理结束(或被 stop 终止)后完成。

FastAsyncQueue.retry

当前正在执行 的任务体内调用(通常在 catch 中)。将队首任务设为 pendingRetry,以便在本轮 start() 中再次 run();超过 retryTime 时设为 failed 并发出 retryLimitReached

dart
try {
  await apiCall();
} catch (_) {
  queue.retry();
}

无参数。


FastAsyncQueue.stop

强制停止:清空待执行链表与 _map,重置 size,发出 queueStopped。可选 callBack 在停止逻辑开始时同步调用。

dart
queue.stop(() => onQueueAborted());
参数类型必填说明
callBackFunction?停止时立即执行的回调。

FastAsyncQueue.clear

内部调用 stop(callBack) 并清空任务映射,等价于停止并丢弃全部任务信息。

dart
queue.clear();
参数类型必填说明
callBackFunction?传给 stop 的回调。

FastAsyncQueue.close

禁止再 addJob(已关闭后再加会走 violateAddWhenClosedClosedQueueException);已在队中或正在执行的任务会继续跑完。发出 queueClosed

dart
queue.close();

无参数。


FastAsyncQueue.addQueueListener

注册单个监听器(后注册会覆盖前者),接收全部 QueueEvent

dart
queue.addQueueListener((QueueEvent event) { /* ... */ });
参数类型必填说明
listenerQueueListenervoid Function(QueueEvent event)

FastAsyncQueue.getJobInfo / list

dart
final info = queue.getJobInfo('chunk-1');
final all = queue.list();
方法返回值说明
getJobInfo(label)JobInfo无此 label 时抛出 InvalidJobLabelException
list()List<JobInfo>当前映射中所有任务的快照。

JobInfo 字段:labeldescriptionstateJobState)、retryCountmaxRetry


属性

类型说明
sizeint(getter)待处理 + 当前执行中的任务数(链表长度)。
isClosedbool(getter)是否已调用 close()

类型与异常

JobState

含义
pending已入队,等待执行
running正在执行
pendingRetry已请求重试,等待再次执行
done成功结束(即将出队)
failed失败且不再重试(即将出队)

相关异常

类型触发条件
DuplicatedLabelExceptionaddJob 使用了已存在的 label
ClosedQueueException对已关闭队列调用 addJobThrow
InvalidJobLabelExceptiongetJobInfo 找不到 label

基于 MIT 许可证发布