容器版的 async/await:Trigger.dev 如何挂起并恢复正在运行的任务

我一直在找一个能运行 AI 工作流的平台——LLM 链、数据管道、智能体循环——又不想自己把队列、重试和调度那一整套基础设施都搭一遍。就这样我遇到了 Trigger.dev,一个用 TypeScript 运行后台任务的开源平台。这个项目做得不错,但其中有一个特性尤其吸引我。

如果你的任务调用了 wait.for(),或者用 triggerAndWait() 等待一个子任务:

export const myTask = task({
  id: "my-task",
  run: async () => {
    // 容器在这里被挂起——这一小时你不用付任何费用
    await wait.for({ hours: 1 });

    // 子任务在它自己的容器里运行期间,本容器处于挂起状态
    const result = await childTask.triggerAndWait({ data: "some data" });
  },
});

平台可以为任务容器创建检查点、释放计算资源,并在等待结束后恢复。开发者无需手动序列化局部变量。这适用于支持的等待操作,而非每个 JavaScript await;短暂的定时等待可能不会创建检查点。

Trigger.dev 将其称为检查点恢复系统。快照保留内存和 CPU 寄存器等进程状态。计算暂停期间仍需存储与调度,外部服务也不会随进程一起回退。

异步函数的类比很有帮助:等待时暂停,准备好后继续。区别在于释放资源的层级。普通 await 仍让进程保持运行,而受支持的 Trigger.dev 等待可以让平台暂停任务容器,并在兼容的工作节点上恢复。

这听上去好得有点不真实,于是我请 Claude 去读了 源代码,弄清它究竟是怎么工作的。答案涉及 CRIU(Checkpoint/Restore In Userspace)、Docker 的实验性 checkpoint API、用于创建 OCI 镜像的 Buildah,以及一套精心设计、把这一切协调起来的状态机。

这篇文章我想分享我的发现。

问题:无服务器任务中的漫长等待

先弄清楚检查点可能在哪些地方派得上用场。设想一个任务:先处理付款,再等待确认,然后发送收据:

import { task, wait } from "@trigger.dev/sdk";

export const processPayment = task({
  id: "process-payment",
  run: async (payload) => {
    const charge = await chargeCustomer(payload);

    // getConfirmation 在它自己的容器里运行期间,父容器处于挂起状态。
    // 这可能耗时数小时甚至数天——等待期间你不付费
    const confirmation = await getConfirmation.triggerAndWait({
      chargeId: charge.id,
    });

    await sendReceipt(charge, confirmation);

    return { success: true };
  },
});

等待数小时的工作流可能超过单次函数调用的执行时限,整个等待期间保留工作节点也会产生成本。

常见的等待处理方式包括:

  1. 让容器一直跑着等确认到来——这段时间你一直在为计算资源付费,尽管父任务什么也没做
  2. 拆成多个任务——把工作流拆成 chargeCustomer、一个定时触发器和 sendReceipt,失去单个函数的简洁。AWS Step Functions 或 Google Cloud Workflows 这类工作流编排器能帮上忙,但你现在调试的是状态机,而不是异步函数
  3. 把状态序列化到数据库——把 charge 存到某处,安排一个后续作业,恢复时再反序列化——这时你已经在造工作流引擎了

Trigger.dev 的答案是第 4 个选项:把容器的内存冻结到磁盘,关掉它,稍后再恢复。

当你调用 triggerAndWait 时,子任务会在一个独立容器中启动,随后父任务被创建检查点并挂起——释放其计算资源与并发额度——直到子任务完成。父任务带着子任务的返回值恢复,就像一次普通的 await。同样的机制也适用于 await wait.for({ hours: 24 }) 这类定时等待。

与传统做法的对比

工作流引擎保存进度的方式不同,可以显式持久化应用状态,也可以重放事件历史。二者都不同于保存进程镜像。

Temporal使用的事件历史重放通过确定性的工作流代码重放已记录事件,重建工作流状态。已完成活动使用记录的结果,而非再次执行普通副作用。这不需要序列化任意闭包或已打开的套接字。

容器检查点保存受支持的进程状态,包括局部变量和闭包所在的内存。它可以避免重放工作流代码,但需要兼容的操作系统和运行时支持。文件、套接字以及其他外部资源仍需单独处理。

这里的取舍在于重放约束与快照基础设施:检查点减少应用层的状态处理,却增加存储、传输和兼容性要求。

CRIU:底层的技术

CRIU 是保存和恢复受支持进程状态的 Linux 工具,可处理内存页、寄存器、文件描述符及相关内核状态,也可作用于进程树。

CRIU 位于语言运行时之下,因此可以保存多种语言编写的程序。但并非所有进程都能创建检查点:设备、内核特性、命名空间和运行时支持都会限制恢复能力。

Docker 通过 docker checkpoint create 提供了对 CRIU 的实验性支持,Kubernetes 则通过 CRI(Container Runtime Interface)用 crictl checkpoint 支持它。Trigger.dev 视部署模式两者都用。当 CRIU 完全不可用时——二进制文件缺失、内核不支持,或者 Docker 的实验特性没开——Trigger.dev 会退回到 docker pause,它只挂起容器,并不捕获状态。工作流照样继续,但如果容器挂了,这次运行就丢了。这个退路是给那些搭建 CRIU 不现实的开发环境准备的。

自己动手试试

本示例在特权 Docker 容器中运行 CRIU,需要兼容的 Linux 内核与 CRIU 配置;仅添加 --privileged 并不保证可用。

docker run -d --name criu-demo --privileged python:3.12-slim bash -c 'apt-get update -qq && apt-get install -y -qq criu > /dev/null 2>&1 && sleep infinity'

拷入计数器脚本——它只是把一个数字加一,每秒写进文件一次:

docker exec criu-demo bash -c 'cat > /counter.py << "EOF"
import time
count = 0
while True:
    count += 1
    with open("/output.txt", "a") as f:
        f.write(f"count = {count}\n")
    time.sleep(1)
EOF'

拷入演示脚本——它启动计数器、做检查点,然后再恢复它:

docker exec criu-demo bash -c 'cat > /demo.sh << "EOF"
#!/bin/bash
python3 /counter.py &
PID=$!
disown
sleep 5

echo "--- before checkpoint ---"
cat /output.txt

mkdir -p /checkpoint
criu dump -t $PID -D /checkpoint --shell-job -v0

echo "--- checkpointed, process killed ---"
> /output.txt

criu restore -d -D /checkpoint --shell-job -v0
sleep 5

echo "--- after restore ---"
cat /output.txt
EOF
chmod +x /demo.sh'

运行这个演示:

docker exec criu-demo /demo.sh

输出:

--- before checkpoint ---
count = 1
count = 2
count = 3
count = 4
count = 5
--- checkpointed, process killed ---
--- after restore ---
count = 7
count = 8
count = 9
count = 10
count = 11

计数器在写完 5 之后被做了检查点,进程随即被杀死,然后 CRIU 从检查点把它还原——它像什么都没发生过一样继续计数。变量 count 原本躺在 Python 的堆内存里,而 CRIU 捕获并还原了整个内存状态。这正是 Trigger.dev 所用的机制,只是外面裹了多得多的编排逻辑。

示例从容器内部为进程创建检查点。后文所述实现则由协调器从外部请求运行时为任务创建检查点。

检查点流程是怎么走的

下面是父任务触发子任务时完整的 checkpoint-resume 流程(依据 Trigger.dev 文档中的示意图):

图中展示的是 triggerAndWait 的流程,但同样的机制也适用于 wait.for()——唯一的区别在于由什么来解除等待点(定时器,还是子任务完成)。

下面的源码链接指向本文参考的具体版本,描述的是该实现,并不保证当前部署仍采用相同的组件布局。

图中参与方源码类
Trigger.devrun-engine/engine/index.tsRunEngine
父任务/子任务managed/controller.tsManagedRunController
CR 系统coordinator/checkpointer.tsCheckpointer
存储coordinator/exec.tsBuildah

图中的 CR 系统 对应的是 Coordinator 容器——这些组件在工作节点虚拟机上的布局如下:

工作节点 VM
├── Supervisor 容器 (apps/supervisor)
│   ├── 从平台队列中取出运行任务
│   ├── 按需创建任务容器
│   └── 与 Coordinator 协同完成检查点
│
├── Coordinator 容器 (apps/coordinator)  ← 图中的「CR 系统」
│   ├── 运行 Checkpointer(CRIU、Buildah)
│   ├── 拥有访问 Docker 守护进程的权限
│   └── 从外部冻结任务容器
│
└── 任务容器(临时的,每次运行一个)
    ├── Controller (ManagedRunController)  [入口点]
    │   └── 在任务可挂起时发出信号
    │
    └── Worker  [子进程,通过 IPC 派生]
        └── 你的任务代码 (task.run())

在这个设计中,控制器发出就绪信号,协调器从任务容器外请求检查点。检查点包含控制器与工作进程状态。要在其他节点恢复,还需要兼容的运行时及必要的文件系统状态。

我们逐步走一遍。

第 1 步:开始执行

运行在工作节点 VM 上的 supervisor 从平台队列取出这次运行,并通过 workloadManager.create() 创建任务容器,同时传入环境变量(例如 TRIGGER_SUPERVISOR_API_DOMAIN),好让任务容器内的控制器进程知道如何通过 HTTP 访问 supervisor 的 Workload API。控制器 ManagedRunController 是任务容器内的主 Node.js 进程——它用 Node 的 fork() 派生出一个 worker 子进程来运行你的任务代码。两个进程通过 Node.js 的 IPC 通信。

第 2 步:触发子任务

当你的代码调用 await childTask.triggerAndWait(...) 时,会发生两件事:

  • 运行在 worker 进程中的 SDK 直接向 Trigger.dev 平台(图中的「Trigger.dev」参与方)发起 API 调用,绕过控制器——worker 有自己通往平台的 HTTP 客户端。这会把子任务排入执行队列,并创建一个等待点(waitpoint)——平台数据库中的一条记录,表示「这次运行正在等待这个子任务完成」。
  • worker 通过 IPC 告知控制器:它可以被挂起了——即可以在不丢数据的前提下冻结。

父任务并不等子任务启动;它只是告诉平台「运行这个」,然后发信号说「现在可以给我做检查点了」。注意这里的分工:控制器负责执行的生命周期(可挂起信号、快照管理),但 SDK 用来触发任务和创建等待点的 API 调用完全绕过它——它们直接由 worker 经 HTTP 发往平台。对 wait.for() 而言,等待点是一个时间点而不是子任务,但流程的其余部分完全一致。

第 3 步:请求快照

任务容器内的控制器在 supervisor 的 Workload API 上调用 suspendRun()(使用第 1 步建立的 HTTP 连接)。supervisor 通过 CheckpointClient 把任务委派给 coordinator。coordinator 从外部调用 CRIU 来冻结任务容器。检查点逻辑位于 checkpointAndPush(),有两种模式:

Docker 模式(本地/开发):

docker checkpoint create --leave-running <container-name> <checkpoint-name>

Kubernetes 模式(生产):

crictl checkpoint --export=/checkpoints/<identifier>.tar <container-id>

两条命令都要求运行时创建进程检查点。Docker 示例使用 --leave-running,因此快照本身不会永久停止原容器;暂停和清理由外围调度流程负责。

第 4 步:存储快照

在生产环境(Kubernetes 模式)中,检查点会被导出为 tar 归档。coordinator 用 Buildah 类把它封装成一个 OCI 容器镜像,并推送到镜像仓库:

buildah from scratch
buildah add <container> /checkpoints/<identifier>.tar /
buildah config --annotation=io.kubernetes.cri-o.annotations.checkpoint.name=<shortCode> <container>
buildah commit <container> <registry>/<namespace>/<project>:<version>.prod-<shortCode>
buildah push --tls-verify <imageRef>

OCI 仓库负责传输检查点制品。恢复仍需要兼容且支持检查点的运行时;将检查点打包成 OCI 镜像,并不意味着任何节点都能像普通镜像一样运行它。

第 5 步:释放资源

检查点镜像存好之后,平台会:

  1. 在平台数据库中把运行状态更新为 WAITING_TO_RESUME(就是第 2 步创建等待点的那个数据库)
  2. 存入一条 TaskRunCheckpoint 记录(类型、位置、镜像引用)
  3. 释放这次运行占用的全部并发额度

如果你的队列设了 concurrencyLimit: 5,而有三个任务处于挂起状态,那三个槽位就会被腾出来。挂起的任务既不消耗计算资源,也不占用并发额度。

第 6 步:子任务完成

子任务在自己的容器里运行。它结束时,其控制器会在 supervisor 上调用 completeRunAttempt(),由 supervisor 把结果上报给平台。平台解除第 2 步创建的父任务等待点,从而触发恢复流程。对 wait.for() 来说,这一步换成定时器到期,它以同样的方式解除等待点。

第 7 步:取回快照并恢复状态

平台向 CR 系统请求检查点,后者从存储中取回快照镜像。以该检查点镜像启动一个新容器——CRIU 把所有进程还原到它们确切的内存状态。

被还原的控制器会察觉自己是恢复而来,于是调用 continueRunExecution():

POST /api/runs/{runId}/continue
Body: { snapshotId: "...", workerId: "...", runnerId: "..." }

恢复后的容器可以运行在另一台兼容工作节点上。共享仓库让制品可被获取,主机与运行时兼容性则决定恢复能否成功。

第 8 步:恢复并完成执行

后端校验快照,把运行状态更新为 EXECUTING,任务从 await 之后的那一行继续。从你代码的视角看,什么都没发生过——await 以子任务的返回值完成,执行照常往下走。

可能出什么问题

检查点不是魔法。也有一些边界情况:

连接可能失效。保存套接字的本地状态并不能让远端一直保持连接。长时间等待后,数据库、HTTP 或 WebSocket 连接可能需要重建,应用代码应处理这种情况。

进程快照与文件系统快照不同。文件描述符指向的文件,在恢复时还必须具备所需内容和元数据。容器可写层或临时文件是否保留,取决于运行时和部署方式,不能仅凭内存快照推断。

内存占用越大,检查点通常也越大。快照大小不一定等于分配的 RAM,驻留页、压缩、稀疏文件和运行时选项都会影响它。大型检查点可能增加存储、传输和恢复时间。

CRIU 需要内核支持。 CRIU 需要特定的内核特性(命名空间、cgroups)以及 Docker 的实验模式。在 Kubernetes 中,容器运行时(CRI-O 或 containerd)必须配置为支持检查点。这并非到处都有——这也正是 Trigger.dev 准备了模拟退路的原因。

同样的套路,用在虚拟机层面

CRIU 工作在进程层面——它捕获的是容器内的单个进程树。但同样的 checkpoint/restore 套路在虚拟机层面也成立。Firecracker——AWS Lambda 和 Fly.io 背后的微虚拟机管理器——能够暂停整台虚拟机,把它完整的内存和设备状态转储成文件,之后再从这些文件恢复,而且是在一个全新的 Firecracker 进程里。

我在开启了 KVM 的 WSL2 上试了一下。配置是:Firecracker v1.12.0、一个带有每秒自增 shell 计数器的 Alpine Linux rootfs,以及一个预编译的 Linux 内核。启动虚拟机、让计数器数到 20 之后,我暂停虚拟机并通过 Firecracker 的 REST API 创建了快照:

# 暂停
curl --unix-socket /tmp/firecracker.socket -X PATCH \
  http://localhost/vm -H 'Content-Type: application/json' \
  -d '{"state": "Paused"}'

# 快照
curl --unix-socket /tmp/firecracker.socket -X PUT \
  http://localhost/snapshot/create -H 'Content-Type: application/json' \
  -d '{"snapshot_type": "Full", "snapshot_path": "./snapshot_file", "mem_file_path": "./mem_file"}'

然后我把 Firecracker 进程彻底杀掉,启动一个新的,并加载这份快照:

curl --unix-socket /tmp/firecracker.socket -X PUT \
  http://localhost/snapshot/load -H 'Content-Type: application/json' \
  -d '{"snapshot_path": "./snapshot_file", "mem_file_path": "./mem_file", "enable_diff_snapshots": false, "resume_vm": true}'

计数器从 21 继续。恢复耗时约 29 毫秒。

Firecracker 快照包含客户机内存与设备状态,但恢复仍受 CPU、版本和主机兼容性要求限制。AWS Lambda SnapStart利用已初始化环境的快照减少启动延迟。这与恢复执行到一半的工作流不同,也不保证每个 Lambda 函数都能在 100 毫秒内启动。

这套设计的优雅之处

这个设计中最打动我的是它的抽象边界。从开发者的视角看:

await wait.for({ hours: 24 });

在这个等待调用背后,平台协调快照创建、存储、资源释放和恢复。任务保留普通异步控制流程,但仍需处理外部连接、副作用和重试。