| title | 任务定义与接口 |
|---|---|
| description | queue.define、execute、onSuccess 与 BatchTask 的 start、get、cancel:如何定义任务、推进分页、提交和查询 Run。 |
Queuebit 的用户侧入口分两层:
- 用
queue.define<Q, S>()定义任务名称、版本、事件、执行函数和回调。 - 使用返回的
BatchTask<Q, S>提交、查询和取消某一任务定义下的 Run。
Q 是整轮 Run 的业务查询参数,S 是 Queuebit 保存的分页状态。需要完整队列初始化与处理逻辑时,先运行快速开始。
const task = queue.define<Query, State>({
name: 'receipt-snapshot',
version: '1',
events: ['success'],
policy: {
attempts: 3,
timeoutMs: 30_000,
backoff: { baseMs: 500, maxMs: 5_000, jitter: 'full' },
},
}, {
async execute(ctx) {
// 读取业务页,写入回执;空页时 return ctx.end()
return ctx.next({ afterSourceId: 1002 });
},
async onSuccess(ctx) {
// Run 成功后交付完成通知
},
});name 和 version 组成任务定义身份。policy 是可选的任务级重试/超时覆盖;未声明字段从队列默认值补齐。不要把业务处理函数、局部变量或 worker 默认值当作定义身份的一部分。
定义必须在 await queue.ready() 前完成。ready() 之后再定义任务会破坏运行时成员的一致视图。
处理函数通过 ctx 读取当前 Run:
| 字段/方法 | 用途 | 边界 |
|---|---|---|
ctx.query |
整轮 Run 的业务选择参数 | 例如 snapshotId,不随翻页改变 |
ctx.state |
上一次已确认页保存的状态 | 初始为 null;不要原地修改 |
ctx.signal |
协作取消与超时信号 | 传给下游 repository/sink |
ctx.next(state) |
本页完成后提交下一状态 | 只在业务写入成功后返回 |
ctx.end() |
没有下一页时结束 Run | 不代表回调已经送达 |
至少一次执行意味着同一页可能重复执行。快速开始里的 sink.putOnce() 用稳定 receipt key 演示业务幂等;生产里要用数据库唯一键、事务或外部服务的幂等接口实现。
events: ['success'] 会在 Run 结束成功后创建完成事件。事件交付和 Run 执行相互独立:
run.status === 'success'不等于完成通知已经送达。- 回调失败会按回调策略重试,耗尽后进入死信。
ctx.eventId是事件身份;业务侧仍应使用稳定 completion key 做幂等。
| 方法 | 用途 | 需要特别处理 |
|---|---|---|
start(input) |
提交一轮运行 | created: false 不代表失败 |
get(runId) |
读取业务参数、状态和回调统计 | null 与“执行失败”不同 |
cancel(runId) |
请求将运行置为取消终态 | 不撤销已发生的外部写入 |
三种方法都依赖队列可用。runtime.mode: 'consumer' 不允许 start(),会报 MODE_OPERATION_NOT_ALLOWED;读取与取消仍通过各自契约校验。
start(input: {
query: Q;
idempotencyKey?: string;
}): Promise<{ runId: string; created: boolean }>;以上是方法签名,不是独立应用代码。
| 参数 | 必填 | 说明 |
|---|---|---|
query |
是 | 整轮 Run 的业务选择参数,必须符合 SDK 的 JSON 值约束 |
idempotencyKey |
否 | 业务方用于识别同一次提交的键;省略时不提供业务提交去重 |
成功返回 runId 和 created。首次创建时为 true;同一任务定义与幂等键命中相同输入时,返回已有 Run,created: false。相同键却提交不同输入会报 IDEMPOTENCY_CONFLICT,不能用同一个键表达两次不同业务操作。
提交去重不是“固定 TTL 到点立刻消失”。它随 Run 的实际回收清理;Run 达到保留期后也可能因为事件引用尚未释放而延后回收。因此不要把短期重复提交的行为推导为永久的业务唯一性保证。
收到 OUTCOME_UNKNOWN 时,不能断定提交失败。保存错误携带的 runId、操作和业务幂等键,先查询核对;查询暂时不可用或尚未查到时,不要立即换新键另投一轮。retryable 和 outcomeKnown 分别表达可重试性与结果是否已知,不能互相替代。
get(runId: string): Promise<RunInfo<Q, S> | null>;未找到运行时返回 null,例如 ID 不存在或运行已被回收。任务定义身份不匹配是错误,不应统一吞成“没找到”。
| 常用字段 | 如何理解 |
|---|---|
status、reason |
当前运行状态及原因 |
query、state |
只读业务参数与已保存状态;初始 state 为 null |
page、batchId |
页推进和批次标识,不应当作业务数据条数 |
businessFailures、scheduledRetries、recoveries |
分别观察业务失败、计划重试与恢复,不合并为一个“失败次数” |
error |
错误摘要,可能为空;并非永久完整原始堆栈 |
callbacks |
pending、delivered、deadLetters,与 Run 状态独立 |
可能的状态为 pending、running、retrying、blocked、pausing、paused、success、failed、cancelled。不要只处理 running/success/failed 三种情况,也不要把 blocked 自动当成终态。
查询结果是一个时点的快照。success 不保证所有回调已经交付;如果业务要求完成通知送达,需要进一步观察 callbacks。
cancel(runId: string): Promise<
| { found: false; runId: string }
| {
found: true;
runId: string;
status: 'cancelled' | 'success' | 'failed';
revision: number;
changed: boolean;
}
>;| 返回 | 调用方应如何处理 |
|---|---|
found: false |
没找到 Run,不包含 status 或 revision |
found: true, changed: true |
状态已改变;仍需理解业务执行的协作退出边界 |
found: true, changed: false |
没有改变已有终态,以返回的 status 为准 |
:::warning 取消不是业务回滚
Queuebit 会通知本地执行使用的 AbortSignal。处理函数应主动检查 signal,下游调用也应支持取消。已经发生的回执写入不会被撤销,已经启动且不响应取消的 Promise 不会被强制终止。
:::
带乐观并发、操作原因与命令 ID 的管理操作属于 queue.operator.runs,不是给这里的 cancel() 额外添加未定义参数。
以下函数接收已定义且已 ready 的 BatchTask;它们不是队列初始化代码。完整初始化见快速开始。
| 错误码 | 检查方向 |
|---|---|
QUEUE_NOT_READY |
等待 ready(),不要靠固定 sleep 猜测连接完成 |
QUEUE_CLOSING / QUEUE_CLOSED |
停止向该实例提交操作,检查生命周期管理 |
MODE_OPERATION_NOT_ALLOWED |
不要在仅 consumer 模式下提交 Run |
JSON_INVALID / PAYLOAD_TOO_LARGE |
检查 query 的数据类型及体积,不通过字符串强转掩盖无效输入 |
IDEMPOTENCY_CONFLICT |
同一提交键的业务输入不同,核对上游请求身份 |
TASK_IDENTITY_MISMATCH |
确认 runId 与当前定义匹配,不默认删除旧 Run |
CAPACITY_EXCEEDED |
检查 namespace 容量与回收情况,不用随机换 namespace 绕过生产限制 |
OUTCOME_UNKNOWN |
查询核对可能已生效的操作,不盲目重投 |
本页是任务对象参考,不是完整错误码目录。配置、运行控制和保留策略应与相同版本的队列协议共同使用。