Skip to content

Latest commit

 

History

History
168 lines (124 loc) · 7.72 KB

File metadata and controls

168 lines (124 loc) · 7.72 KB
title 任务定义与接口
description queue.define、execute、onSuccess 与 BatchTask 的 start、get、cancel:如何定义任务、推进分页、提交和查询 Run。

任务定义与接口

Queuebit 的用户侧入口分两层:

  1. queue.define<Q, S>() 定义任务名称、版本、事件、执行函数和回调。
  2. 使用返回的 BatchTask<Q, S> 提交、查询和取消某一任务定义下的 Run。

Q 是整轮 Run 的业务查询参数,S 是 Queuebit 保存的分页状态。需要完整队列初始化与处理逻辑时,先运行快速开始

define:定义任务

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 成功后交付完成通知
  },
});

nameversion 组成任务定义身份。policy 是可选的任务级重试/超时覆盖;未声明字段从队列默认值补齐。不要把业务处理函数、局部变量或 worker 默认值当作定义身份的一部分。

定义必须在 await queue.ready() 前完成。ready() 之后再定义任务会破坏运行时成员的一致视图。

execute:推进分页

处理函数通过 ctx 读取当前 Run:

字段/方法 用途 边界
ctx.query 整轮 Run 的业务选择参数 例如 snapshotId,不随翻页改变
ctx.state 上一次已确认页保存的状态 初始为 null;不要原地修改
ctx.signal 协作取消与超时信号 传给下游 repository/sink
ctx.next(state) 本页完成后提交下一状态 只在业务写入成功后返回
ctx.end() 没有下一页时结束 Run 不代表回调已经送达

至少一次执行意味着同一页可能重复执行。快速开始里的 sink.putOnce() 用稳定 receipt key 演示业务幂等;生产里要用数据库唯一键、事务或外部服务的幂等接口实现。

onSuccess:完成回调

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:提交运行

start(input: {
  query: Q;
  idempotencyKey?: string;
}): Promise<{ runId: string; created: boolean }>;

以上是方法签名,不是独立应用代码。

参数 必填 说明
query 整轮 Run 的业务选择参数,必须符合 SDK 的 JSON 值约束
idempotencyKey 业务方用于识别同一次提交的键;省略时不提供业务提交去重

成功返回 runIdcreated。首次创建时为 true;同一任务定义与幂等键命中相同输入时,返回已有 Run,created: false。相同键却提交不同输入会报 IDEMPOTENCY_CONFLICT,不能用同一个键表达两次不同业务操作。

去重持续多久

提交去重不是“固定 TTL 到点立刻消失”。它随 Run 的实际回收清理;Run 达到保留期后也可能因为事件引用尚未释放而延后回收。因此不要把短期重复提交的行为推导为永久的业务唯一性保证。

提交结果不确定

收到 OUTCOME_UNKNOWN 时,不能断定提交失败。保存错误携带的 runId、操作和业务幂等键,先查询核对;查询暂时不可用或尚未查到时,不要立即换新键另投一轮。retryableoutcomeKnown 分别表达可重试性与结果是否已知,不能互相替代。

get:读取运行

get(runId: string): Promise<RunInfo<Q, S> | null>;

未找到运行时返回 null,例如 ID 不存在或运行已被回收。任务定义身份不匹配是错误,不应统一吞成“没找到”。

常用字段 如何理解
statusreason 当前运行状态及原因
querystate 只读业务参数与已保存状态;初始 statenull
pagebatchId 页推进和批次标识,不应当作业务数据条数
businessFailuresscheduledRetriesrecoveries 分别观察业务失败、计划重试与恢复,不合并为一个“失败次数”
error 错误摘要,可能为空;并非永久完整原始堆栈
callbacks pendingdelivereddeadLetters,与 Run 状态独立

可能的状态为 pendingrunningretryingblockedpausingpausedsuccessfailedcancelled。不要只处理 running/success/failed 三种情况,也不要把 blocked 自动当成终态。

查询结果是一个时点的快照。success 不保证所有回调已经交付;如果业务要求完成通知送达,需要进一步观察 callbacks

cancel:取消运行

cancel(runId: string): Promise<
  | { found: false; runId: string }
  | {
      found: true;
      runId: string;
      status: 'cancelled' | 'success' | 'failed';
      revision: number;
      changed: boolean;
    }
>;
返回 调用方应如何处理
found: false 没找到 Run,不包含 statusrevision
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 查询核对可能已生效的操作,不盲目重投

本页是任务对象参考,不是完整错误码目录。配置、运行控制和保留策略应与相同版本的队列协议共同使用。