Convex Workflow:接口、事件与恢复 API 实战

This article is extracted from the chat log with AI. Please identify it with caution.

说明:本文由 Codex 根据作者提供的主题与 Convex 官方资料辅助生成,属于 AI 生成/整理内容,非作者原创。请读者自行甄别并交叉验证。

Convex Workflow 是一个代码式 durable execution Component:把长任务拆成可持久化的步骤,记录每一步的结果,并在服务重启、等待事件或暂时失败后从下一未完成步骤继续。它适合订单履约、人工审批、第三方 API 编排、长时间 Agent 任务和批量处理;不适合只需一次数据库事务或一次短 Action 的工作。

最容易误解的一点是:workflow handler 不是一个一直挂在内存里的异步函数。每次准备执行下一步骤时,组件会从头确定性地重放 handler,并从 journal 中取回完成步骤的结果。因此,普通的 if、循环、try…catchPromise.all() 都能使用,但 handler 本身必须是确定性的;网络调用、环境变量和加密等副作用要放进 Action step。

本文按 @convex-dev/workflow npm 当前稳定版 0.4.4(2026-08-07 核对)说明。GitHub main 可能先出现尚未发布的 API;例如 step.withOptions() 不属于该稳定版的公开接口,不应直接复制到生产代码。npm 包信息 官方 README

先看执行与重放时间线#

sequenceDiagram
  participant M as Mutation / Action
  participant W as Workflow Component
  participant H as workflow handler
  participant S as durable step
  participant X as 外部服务 / 人工审批

  M->>W: start(ctx, workflow, args)
  W-->>M: WorkflowId
  W->>H: 从 handler 开头执行
  H->>S: runAction / runMutation / sleep
  S->>X: 执行或等待
  X-->>S: 结果或事件
  S->>W: 持久化步骤结果
  W->>H: 从头重放,跳过已完成步骤
  H-->>W: 返回完成结果

这带来四条设计规则:

  1. handler 的跨步骤状态只能来自输入、已完成步骤的返回值或事件值,不依赖模块变量和进程内内存。
  2. 外部副作用应放在 step.runAction(),并使用业务幂等键;网络超时不等于对方没有成功。
  3. 步骤名称、类型、参数与控制流都要在活动实例的生命周期内保持稳定。
  4. 只把小而可序列化的值放进步骤结果;大对象写入 Convex 数据库或文件存储,在步骤间传 ID。

1. 安装、注册与 WorkflowManager#

前提是已有一个 Convex 项目。安装 Component、在 convex/convex.config.ts 注册,然后运行 npx convex dev 或部署命令,让 Convex 生成 components.workflow 类型。

npm install @convex-dev/workflow@0.4.4
// convex/convex.config.ts
import workflowComponent from "@convex-dev/workflow/convex.config.js";
import { defineApp } from "convex/server";

const app = defineApp();
app.use(workflowComponent);

export default app;

配置保存后运行开发命令,让 Component 的类型出现在生成 API 中:

npx convex dev

在项目内创建一个 Manager。下面的并行度和重试策略只是示例,不是必须照抄的全局默认值。

// convex/workflow.ts
import { WorkflowManager } from "@convex-dev/workflow";
import { components } from "./_generated/api";

export const workflow = new WorkflowManager(components.workflow, {
  workpoolOptions: {
    maxParallelism: 10,
    defaultRetryBehavior: {
      maxAttempts: 3,
      initialBackoffMs: 1_000,
      base: 2,
    },
    retryActionsByDefault: false,
  },
});

components.workflow 对应上例的默认 Component 实例。maxParallelism 限制的是这个 Workpool 同时执行的步骤数,不是运行中 workflow 实例总数;同一个 Component 实例应只有一套一致的并行度配置。只有当所有 Action 都已实现幂等时,才考虑把 retryActionsByDefault 设为 true

2. 定义 Workflow:define、handler 与返回值#

核心形状是 workflow.define(config).handler(handler)

API作用关键配置
new WorkflowManager(components.workflow, options)绑定 Component,并设置共享 Workpool 行为workpoolOptions:并行度、默认重试、Action 是否默认重试
workflow.define(config)声明一个 workflow 函数argsreturns、可选的 workflow 级 workpoolOptions
.handler(async (step, args) => …)编写可确定性重放的控制流显式标注 Promise 返回类型,避免 internal.* 引用造成 TypeScript 类型循环
defineWorkflow(components.workflow, config)不创建 Manager 时的函数式替代写法define() 使用相同的 args、returns、handler 语义

下面是一个订单审批骨架。internal.orders.getForWorkflowinternal.shipping.createLabelinternal.orders.markFulfilled 是项目自己的内部 Convex 函数,需要按业务数据模型实现参数和返回值 validator;示例重点是 Workflow API 的可直接套用结构,不包含真实订单鉴权或物流服务。

// convex/orderWorkflow.ts
import { defineEvent } from "@convex-dev/workflow";
import { v } from "convex/values";
import { internal } from "./_generated/api";
import { workflow } from "./workflow";

export const approvalEvent = defineEvent({
  name: "order-approval-v1",
  validator: v.object({ approved: v.boolean() }),
});

export const fulfillOrder = workflow
  .define({
    args: { orderId: v.string() },
    returns: v.union(v.literal("fulfilled"), v.literal("rejected")),
  })
  .handler(async (step, { orderId }): Promise<"fulfilled" | "rejected"> => {
    const order = await step.runQuery(
      internal.orders.getForWorkflow,
      { orderId },
      { name: "orders/get-for-workflow-v1" },
    );

    const approval = await step.awaitEvent(approvalEvent);
    if (!approval.approved) {
      return "rejected";
    }

    await step.runAction(
      internal.shipping.createLabel,
      { orderId, address: order.shippingAddress },
      {
        name: "shipping/create-label-v1",
        retry: true,
      },
    );

    await step.runMutation(
      internal.orders.markFulfilled,
      { orderId },
      { name: "orders/mark-fulfilled-v1" },
    );

    return "fulfilled";
  });

returns 不只是 TypeScript 提示:它会在运行时验证 workflow 的最终返回值。对 handler 写出 Promise<…> 返回类型也很重要,因为 workflow 内部引用的 internal.* 函数会参与生成类型,省略类型时容易形成循环推导。

3. Step API:在 handler 内可以做什么#

每个 step 都是持久化边界。前四个执行 API 接受一个函数引用和参数;常用选项包括稳定的 name、延迟调度的 runAfterrunAt(二选一,后者是 epoch 毫秒)。延迟调度是尽力而为的时间点,实际执行可能稍晚。

Step API用途与恢复语义关键选项 / 注意点
step.runQuery(ref, args, options)运行一个 Convex query 并持久化结果;默认经 Workpool 作为独立事务执行可使用 inline: true 与 workflow 当前 mutation 共享事务;不能与 runAt/runAfter 一起使用
step.runMutation(ref, args, options)运行一个 Convex mutation 并持久化结果默认独立事务;inline: true 适合小读写,但共享事务预算
step.runAction(ref, args, options)执行网络请求、LLM、邮件、支付等外部副作用不能 inline;默认不重试;retry 可为 truefalse 或自定义退避策略
step.runWorkflow(ref, args, options)把子 workflow 作为父 workflow 的一个 durable step父流程等待子流程结束并接收返回值;父状态会反映活动子 workflow ID
step.sleep(milliseconds, options)持久化地等待一段时间,等待期间不占用执行资源可用 name 提高日志、排查和重启的可读性
step.awaitEvent(event)等待人工审批、支付 webhook 或其他外部信号可按 name 或动态创建的 id 等待;可附 validator 做类型与运行时验证
step.workflowId当前实例的 WorkflowId用于把父子流程、领域记录或审计日志关联起来

并行和延迟可以直接由普通 TypeScript 表达:

const [fraudCheck, inventory] = await Promise.all([
  step.runAction(internal.risk.checkOrder, { orderId }, { retry: true }),
  step.runMutation(internal.inventory.reserve, { orderId }),
]);

await step.sleep(30 * 60 * 1_000, { name: "payment-grace-period-v1" });

await step.runAction(
  internal.notifications.sendReminder,
  { orderId },
  { runAfter: 60 * 60 * 1_000, retry: true },
);

同一轮 Promise.all() 中的步骤会并行调度,handler 要等它们全部完成才继续。为避免一个 workflow 独占资源,Component 会受 Workpool 并行度约束。

inline 何时适用#

默认的 runQuery()runMutation() 各自使用独立事务,因此获得独立的读写预算与恢复边界。传入 { inline: true } 后,它们会共享 workflow body 的 mutation 事务,适合非常小、确实需要同一事务的读写:

const order = await step.runQuery(
  internal.orders.getForWorkflow,
  { orderId },
  { inline: true },
);

await step.runMutation(
  internal.orders.recordAudit,
  { orderId, status: order.status },
  { inline: true },
);

不要为“少一次调度”而滥用 inline:大读写会与 workflow body 共用 transaction limits,反而缩小可用预算;Action 也不能 inline。

4. 启动、完成回调与状态查询#

Workflow 不能像普通同步函数一样把最终结果直接返回给启动者。应从 mutation 或 action 调用 start()(也可使用 Manager 上的同名方法),立刻得到 WorkflowId 并将它保存到自己的业务记录。

// convex/orderWorkflow.ts
import { cleanup, start, vWorkflowId } from "@convex-dev/workflow";
import { vResultValidator } from "@convex-dev/workpool";
import { v } from "convex/values";
import { components, internal } from "./_generated/api";
import { internalMutation } from "./_generated/server";

export const startFulfillment = internalMutation({
  args: { orderId: v.string() },
  returns: vWorkflowId,
  handler: async (ctx, { orderId }) => {
    return await start(
      ctx,
      internal.orderWorkflow.fulfillOrder,
      { orderId },
      {
        onComplete: internal.orderWorkflow.recordCompletion,
        context: { orderId },
      },
    );
  },
});

export const recordCompletion = internalMutation({
  args: {
    workflowId: vWorkflowId,
    context: v.object({ orderId: v.string() }),
    result: vResultValidator,
  },
  returns: v.null(),
  handler: async (ctx, args): Promise<null> => {
    if (args.result.kind === "success") {
      console.log("workflow completed", args.context.orderId);
    } else if (args.result.kind === "failed") {
      console.error("workflow failed", args.result.error);
    } else {
      console.log("workflow canceled", args.context.orderId);
    }

    // 先持久化业务结果或审计信息,再决定是否删除 workflow 历史。
    await cleanup(ctx, components.workflow, args.workflowId);
    return null;
  },
});

上例使用内部 mutation,意味着公开入口应在自己的 public mutation/action 中先完成身份认证、订单归属和租户检查后再调用它。不要把“知道 WorkflowId 就能查看状态、取消或发送审批”暴露给未授权客户端。

onComplete 是 workflow 的完成回调,适合更新领域状态、写审计记录和清理资源。稳定版 0.4.4result.kindsuccessfailedcanceled;不要照抄旧示例中出现的 error 分支。完成后的 workflow 不会自动清理,若前端仍需读取最终状态,就不要在回调中立刻调用 cleanup()

start() 可接受 startAsync。默认值为 false,创建时会同步评估 handler 的首个阶段,便于尽早暴露参数或引用错误;设为 true 则通过 Workpool 异步入队,适合大量创建且调用方必须快速返回的场景。

生命周期与运维 API#

API应在何处调用作用与边界
getStatus(ctx, components.workflow, workflowId)query、mutation 或 action返回 inProgresscompletedfailedcanceled 状态;包一层 query 后,前端可得到响应式更新
cancel(ctx, components.workflow, workflowId)mutation 或 action停止后续编排;已经开始的 runAction() 不会被强制中断,仍可能完成
restart(ctx, components.workflow, workflowId, options)mutation 或 action用既有 journal 重放;from 可指定从 0 起的步骤序号、步骤/事件名或函数引用开始
cleanup(ctx, components.workflow, workflowId)mutation 或 action清除已完成 workflow 的 Component 存储;清理前先决定保留多久的状态和审计信息
listlistByNamelistSteps通常在受限的 query、mutation、action 或 cron 中分页查看实例和步骤,用于运维控制台、告警排查与批量清理

restart() 默认在当前事务内执行;若 handler 仍报错,restart 操作也会回滚。传入 { startAsync: true } 可把重启交给 Workpool。from 会删除该步骤及之后的记录,所以它应是明确的运维动作,而不是面向普通用户的“再试一次”按钮。

5. 人工审批、Webhook 与外部事件#

事件是 Convex Workflow 的“暂停后恢复”接口。前面的 approvalEvent 同时定义名称和 validator,workflow 与发送方复用它,避免名称或 payload 漂移。

// convex/orderApproval.ts
import { sendEvent, vWorkflowId } from "@convex-dev/workflow";
import { v } from "convex/values";
import { components } from "./_generated/api";
import { internalMutation } from "./_generated/server";
import { approvalEvent } from "./orderWorkflow";

export const submitApproval = internalMutation({
  args: {
    workflowId: vWorkflowId,
    approved: v.boolean(),
  },
  returns: v.null(),
  handler: async (ctx, args): Promise<null> => {
    await sendEvent(ctx, components.workflow, {
      ...approvalEvent,
      workflowId: args.workflowId,
      value: { approved: args.approved },
    });
    return null;
  },
});

将它接到真实前端或 webhook 时,外围 public mutation/action 必须先验证用户身份、签名、订单与 workflow 的归属。事件按名称等待第一个未消费的匹配值;如果 webhook 先到、workflow 后开始等待,组件会保留该事件并在等待步骤到达时立即恢复。发送 { error: “…” } 会让 awaitEvent() 抛异常,可在 handler 的 try…catch 中实现拒绝、超时或补偿路径。

当事件名称本身也要按运行时生成时,用 createEvent() 创建并保存 event ID,再在 workflow 中 awaitEvent({ id })、在 mutation/action 中用同一 ID 调用 sendEvent()。动态 ID 适合工具调用回调或每个任务独立的人工输入;稳定、可复用的业务事件优先使用 defineEvent()

6. 重试、幂等、并发与错误处理#

Convex 的重试语义要按步骤类型区分:

工作单元系统错误时的基础语义设计要求
workflow handler、query、mutation使用 Convex 的事务与系统错误重试保证mutation 不会部分提交或重复提交;仍应控制单个事务的读写规模
runAction()默认不重试对可安全重试的外部调用显式传 retry: true 或自定义策略
已启用重试的 Action指数退避和抖动由 Workpool 策略处理第三方 API、支付、发信、创建资源都必须使用业务幂等键或去重记录

为单个 Action 定义更严格的策略:

await step.runAction(
  internal.billing.capturePayment,
  { orderId, idempotencyKey: "payment:" + orderId },
  {
    name: "billing/capture-payment-v1",
    retry: {
      maxAttempts: 3,
      initialBackoffMs: 1_000,
      base: 2,
    },
  },
);

retry: true 使用 Manager 或 workflow 级的默认策略,retry: false 明确关闭重试。不要因为“失败了就重跑”而给所有 Action 开启全局重试:如果支付服务已经扣款、但网络响应在返回前中断,下一次尝试可能产生重复扣款。正确的补救是在 Action 内传递稳定 idempotency key,或把外部请求 ID 和业务状态持久化后再决定是否调用。

7. 确定性、步骤命名与版本演进#

组件在重放时会比对 journal 中的步骤名称、种类和参数。以下规则能避免活动实例在一次部署后出现 determinism violation:

  1. handler 中不直接做副作用。 不要直接 fetch、读取环境变量或调用加密 API;放进 runAction()Math.random() 在 handler 中是确定性伪随机值,不能作安全用途。
  2. 条件与循环只依赖稳定数据。 可以依据 workflow 输入、前一个步骤的返回值和事件值分支;不要依赖进程内状态、未持久化的当前时间结果或随部署变化的配置。
  3. 为长期步骤指定稳定 name。 默认名称与函数路径有关;重命名或移动函数会影响活动实例。显式使用类似 shipping/create-label-v1 的名称能让迁移意图更清楚。
  4. 不要在活动实例中增删或重排步骤。 新流程应部署为新版本 workflow,或等旧实例完成后再改原实现。
  5. 谨慎使用 unstableArgs: true 它只跳过步骤参数的 journal 比对,不跳过名称和种类比对;适合非常有限的兼容场景,不能把不确定性“修好”。

可恢复的流程不是静态 DAG:正常 TypeScript 的分支、循环和异常处理都被支持。约束在于同一个 WorkflowId 的未来重放必须看到与过去一致的步骤历史。

8. 容量、清理与上线检查#

当前组件文档给出的重要边界包括:单次 workflow execution 中步骤参数与返回值总量不超过 1 MB,journal 上限为 8 MiB,workflow body 作为 mutation 还受到 Convex mutation 的事务限制。文件、模型输出、抓取页面或大批量数据应写入数据库/文件存储,步骤之间只传 ID、对象键、摘要或 checksum。

上线前逐项检查:

  1. 所有公开启动、状态、取消、重启和事件入口都验证身份、租户与领域对象归属。
  2. 每个 Action 的副作用都有幂等键或可查询的外部请求记录,再决定是否启用 retry。
  3. 每个长等待和关键副作用都有稳定的 name;规划好旧 workflow 完成前的版本策略。
  4. WorkflowId 保存到领域记录;UI 通过受权限保护的 query 读取 getStatus(),不要把 Component 管理 API 直接交给浏览器。
  5. 定义完成后的留存策略:需要展示最终进度就延迟 cleanup;只需领域结果就先保存结果再清理。
  6. npx convex dev 验证 Component codegen、函数 validator 和 TypeScript 类型;对真实外部 Action 做一次带幂等键的端到端测试。

继续阅读#

本文共 4964 字,创建于 Aug 7, 2026

相关标签: Serverless, DevOps, TypeScript, Convex, ByAI