- 任务调度
- 后端
【免费下载链接】river
The polyglot queue: Fast and reliable background jobs in Go, Ruby, Rust, and JS/TS on Postgres or SQLite.
本文是 River TypeScript 可选扩展包@riverqueue/worker-threads的完整技术指南。该包通过一组有界的原生 worker thread(Node.js worker threads)执行 CPU 密集型的 River 任务处理器,使这类任务不再阻塞 River 用于认领、完成、取消和优雅停机的事件循环。读完本文,你将掌握如何定义线程处理器、以类型安全的方式注册处理器模块、控制线程池容量与资源上限、理解参数/结果/错误跨越线程边界的编码规则,以及如何配合取消、超时、崩溃恢复与关闭流程,将 CPU 密集任务安全地并入现有的 River 客户端体系。
一、为什么需要独立的线程执行器
River 的普通任务处理器(in-process handler)与运行时共享同一条事件循环。对以 I/O 为主的任务而言这很高效:promise 本身就允许它们在事件循环上交错执行,因此常规 I/O 处理器应继续留在进程内。但当处理器执行的是无法让出事件循环的 CPU 密集计算(例如大图缩放、素数搜索、重计算)时,它会长时间占用主线程,阻塞 River 的认领、完成、取消与停机逻辑。
@riverqueue/worker-threads提供了一条出路:把这类处理器放进一个有界的、由 River 管理的原生线程池中,每个线程并行运行独立的处理器,主线程的事件循环得以保持响应。从源码结构看,该包以可选执行器(executor)的形式接入 River 的运行时:WorkerThreads实现了 js/src/worker.ts 中定义的WorkExecutor接口(name、start(context, handler)、可选diagnostics()),再通过Workers.addExecutor注册到处理器注册表,与普通处理器完全兼容地参与任务调度。
二、核心类型全景:API 骨架速览
包的公共 API 全部集中在一个自动生成的声明文件中:js/worker-threads/etc/worker-threads.api.md,运行时实现与声明一一对应地存在于 js/worker-threads/src/index.ts。整包公开的类型与类如下:
| 类型/类 | 角色 |
|---|---|
WorkerThreads | 执行器主体,实现AsyncDisposable与WorkExecutor |
WorkerThreadsOptions | 执行器配置:maxThreads与可选的resourceLimits |
WorkerThreadWorkHandler | 线程处理器的类型签名 |
WorkerThreadWorkContext | 传给线程处理器的上下文(WorkContext的子集) |
WorkerThreadModule | 用模块导出类型注释过的 ESM 模块 URL |
WorkerThreadExportName | 模块中能处理指定定义的那些导出名的类型级计算 |
WorkerThreadHandlerTarget | exportName+module的组合,即"处理器在哪里" |
WorkerThreadHandlerError | 处理器抛出的错误经边界传递后的载体 |
WorkerThreadsDiagnostics | 线程池运行状况快照 |
下面逐一深入。
三、定义线程处理器:ESM 导出 + 类型安全的上下文
3.1 为什么必须是模块导出
闭包(closure)无法跨线程边界传递,因此线程处理器必须是ESM 模块的具名导出,并通过"模块 URL + 导出名"来引用。实际加载发生在线程内部:线程入口 js/worker-threads/src/thread.ts 收到run消息后await import(message.moduleUrl),再取出module[message.exportName],确认它是函数后才调用。
一个值得注意的工程细节是:任务定义(job definition)应放在独立模块中,这样生产者只需导入定义、无需连带加载处理器代码。官方示例 js/examples/worker-thread-cpu/src/jobs.ts 把定义放在jobs.ts:
import { defineJob } from "riverqueue"; import { z } from "zod"; export const findPrime = defineJob({ kind: "example.find_prime", schema: z.object({ ordinal: z.number().int().positive() }), });3.2 处理器签名
WorkerThreadWorkHandler的完整签名与普通处理器几乎一致,只是上下文类型被替换为线程专用版本:
export type WorkerThreadWorkHandler<Definition extends JobDefinition = JobDefinition> = ( context: WorkerThreadWorkContext<Definition> ) => PromiseLike<WorkOutcome | void> | WorkOutcome | void;它成功的方式与进程内处理器相同:不返回任何值,或返回一个 River outcome(如complete、snooze);失败则直接抛出异常。outcome 以 River JSON 形式跨回主线程——特别地,snooze的Temporal.Duration先被序列化为 ISO 8601 文本跨过边界,再由主线程侧 decodeOutcome 读回为Temporal.Duration。
处理器示例(对应 README 与示例中的写法):
// handler.ts —— 仅使用 type-only 导入,避免线程加载多余代码 import type { WorkerThreadWorkHandler } from "@riverqueue/worker-threads"; import { complete } from "riverqueue"; import type { findPrime } from "./jobs.ts"; export const findPrimeHandler: WorkerThreadWorkHandler<typeof findPrime> = ({ job, signal, }) => complete({ output: { prime: nthPrime(job.args.ordinal, signal) } });用定义来注释导出,可以让job.args获得完整类型推断;同时WorkerThreads.handler也据此在类型层面校验"该导出确实是为该定义写的处理器"。
3.3 线程上下文里有什么
WorkerThreadWorkContext是主线程WorkContext的Pick子集,只包含六项:
Pick<WorkContext<Definition>, "execution" | "job" | "logger" | "recordOutput" | "setMetadata" | "signal">job.args:主线程中由任务定义解码并校验过的参数(与进程内处理器完全一致);job.rawArgs:数据库中持久化的原始 JSON;execution:执行元数据(attemptedBy、startedAt),其中startedAt在线程侧被重建为原生Temporal.Instant;logger、recordOutput、setMetadata:接收 River JSON 值,异步转发回主线程执行;signal:取消信号。
没有client、completeTx和resumable——它们依赖主线程的连接与状态,无法跨线程。
四、模块引用与类型安全:URL 带类型
4.1 WorkerThreadModule:给 URL 附加模块导出类型
WorkerThreadModule在类型层面把"某个 URL 指向哪个模块"绑定到 URL 上:
export type WorkerThreadModule<Module = unknown> = URL & { readonly [workerThreadModuleExports]?: Module; };运行时它就是一个普通URL(任何URL都可赋值给它),唯一的区别是携带了类型信息。使用import type * as imageHandlers from "./image-handler.js"这种type-only 批量导入来获得模块类型,不会把模块代码拉进主线程:
const handlerModule: WorkerThreadModule<typeof imageHandlers> = new URL( "./handler.js", import.meta.url );4.2 类型级导出名校验
WorkerThreadExportName是一个条件映射类型:当模块类型未知(即普通URL)时退化为string;当模块类型已知时,它只保留那些值类型满足WorkerThreadWorkHandler<Definition>的导出名,其余全部映射为never。
WorkerThreadHandlerTarget把二者组合起来:
export interface WorkerThreadHandlerTarget<Definition, Module> { readonly exportName: WorkerThreadExportName<Module, Definition>; readonly module: WorkerThreadModule<Module>; }效果体现在编译期:写错的导出名、指向非处理器的导出、或参数类型不匹配的处理器都会在executor.handler(...)处直接报类型错误。测试 js/worker-threads/src/index.test.ts 用三处@ts-expect-error逐一验证了"模块没有该导出""nthSquare不是线程处理器""describeRich处理的是另一个定义的参数"三种被拒绝的情形。若使用普通URL,任何导出名都能通过编译,错误要等运行时(线程内)才暴露——此时加载到的不是函数会以ESM export "..." is not a function失败该次尝试。
4.3 注册到 Workers 注册表
把三者串起来,通过Workers.addExecutor注册(完整示例见 js/examples/worker-thread-cpu/src/index.ts):
await using executor = new WorkerThreads({ maxThreads: 4 }); const workers = new Workers().addExecutor( resizeImage, executor.handler(resizeImage, { exportName: "resizeImageHandler", module: imageModule, }) );executor.handler在运行时还会做防御性校验(js/worker-threads/src/index.ts):definition必须是带kind的对象、module必须是绝对 URL、exportName必须是非空字符串,否则抛ConfigurationError。注册的处理器同时覆盖定义的kind与kindAliases(与 Workers 注册表行为一致)。传入的 target 若来自其他 executor、或尝试处理别的 job kind,start会分别以"worker thread handler was not created by this executor"和 kind 不匹配错误拒绝。
五、参数、结果与错误如何跨过线程边界
这是本包最精细的部分,核心逻辑集中在 js/worker-threads/src/args.ts 与 js/worker-threads/src/protocol.ts。
5.1 参数编码:River JSON 文本优先,结构化克隆兜底
参数由主线程按定义解码并验证——任何语言插入的、参数非法的任务都会在线程接手前失败。解码后的参数穿越边界有两种编码(args.ts 与 protocol.ts):
- River JSON 文本(
encoding: "json"):凡是能被stringifyJson序列化的参数,一律以文本跨线程。这是首选路径,因为它能精确保留 JSON 数字——例如 Go 端产出的 int64 ID(ExactJsonNumber),结构化克隆会破坏其精度。测试 js/worker-threads/src/index.test.ts 验证了9007199254740993这样的精确整数在 outcome、日志属性、输出与元数据四个方向上都能原样保留。 - 受检结构化克隆(
encoding: "clone"):非 River JSON 的已解码参数,只要能在结构化克隆后原样到达,就按值传递。允许的类型包括bigint、Date、Map、Set、Uint8Array以及由这些构成的普通对象与数组。
encodeArgs会递归检查(args.ts):函数、symbol、symbol 键、访问器属性、非普通类实例(如自定义Money类,其原型会在克隆时丢失)以及混在非 JSON 参数里的精确 JSON 数字,都会让该次尝试直接以ConfigurationError失败,并给出具体路径(如$.price is a Money instance、$.nested[0] is a function、$.id is an exact JSON number),绝不会让处理器拿到类型失真后的参数。
5.2 结果与遥测的边界规则
outcome、输出、元数据与日志属性均以 River JSON 文本回传(thread.ts):
- 输出与元数据:
recordOutput和每个setMetadata值若超过32 MiB,在线程侧即抛ValidationError("job output must not exceed 32 MiB"),与主线程上的行为一致; - 日志:单条消息被截断到32 KiB;日志属性的 JSON 若超过 32 KiB,属性被丢弃并在消息末尾附注
[log attributes omitted: N characters],而不会让处理器失败; - 这些转发在池侧(pool.ts)若因主线程回调抛出异常,会通过"abandon"流程使当前尝试失败并弃用该线程,而不是把异常泄漏到主进程。
5.3 错误的边界封装
处理器抛出的错误经 serializeError 序列化:尽力读取message、name、stack(即使 getter 会抛异常也能安全退化),并截断到 River 持久化尝试错误所用的上限——错误名 256 字符、消息与栈各 32 768 字符。主线程侧收到后包装为WorkerThreadHandlerError,尽量还原原始name、message、stack,并以其使该次尝试失败。线程自身崩溃或意外退出(未捕获异常、未处理的 rejection、process.exit()、资源限制等)也以该错误使正在运行的尝试失败。测试用 10 万字符的巨型错误验证了消息恰被截到 32 768 字符(index.test.ts)。
六、取消、超时与崩溃:线程不会卡死任务
6.1 协作式取消 + 强制终止
当尝试被取消、超时或所属客户端停机时,处理器的signal会被 abort。若处理器在客户端jobStuckThreshold(默认 10 秒)之后仍未自行结束,River 会直接 terminate 该线程——因此一个永不 yield 的 CPU 死循环也无法无限占用 River。关键保证(pool.ts):只有在线程已自行结束或已被终止之后,River 才持久化尝试的结局。
结局分类很讲究:
- 若 abort 来自客户端停机(
run.stop的取消模式),被强制终止的处理器以JobAbortedError失败——该尝试计入次数、适用重试策略; - 若 abort 来自任务取消或超时,则尝试按处理器自行停止处理。
测试 js/worker-threads/src/client.test.ts 精确验证了这两个方向:spin处理器无视 abort 时,stop 会等待满jobStuckThreshold(200ms 测试阈值)才终止线程,任务以attempt: 1+ "job aborted after ignoring cancellation" 错误回到retryable/available;而排队中从未开始的等待任务则attempt: 0、无错误地回到available。
6.2 等待线程不计入任务超时
当全部maxThreads线程忙碌时,后续尝试在池的 FIFO 队列中排队(pool.ts)。River 通过WorkExecutorHandle.started这一能力(worker.ts 中的started?)感知任务真正开始执行的时间点,只在线程接手后才启动任务超时与卡死检测。因此客户端可以配置比线程数更多的maxWorkers,而排队中的任务不会白白消耗超时时间。
6.3 崩溃隔离与懒替换
- 线程崩溃只失败它正在运行的那一次尝试;闲置时崩溃(典型来源:处理器返回后遗留的后台工作,如残留 timer 抛出异常)只是丢弃该线程;
- 新线程在后续尝试需要时才懒启动替换;
- 被复用线程在接到下一个任务前就死掉时,该任务会被重新入队至多一次(pool.ts 的
#settleOrphan路径),避免一次后台崩溃吞掉排队任务; - 线程的
error/exit/message监听器常驻(pool.ts),确保线程内任何失败都不会以主进程未捕获异常的形式出现——测试 index.test.ts 专门监听主进程的uncaughtException/unhandledRejection,断言线程失败从不泄漏到宿主进程。
七、线程池配置与运行诊断
7.1 WorkerThreadsOptions
WorkerThreadsOptions只有两个字段:
| 字段 | 说明 | 校验 |
|---|---|---|
maxThreads | 活跃原生线程数量上限;超出部分排队等待,不消耗任务超时 | 必须是 ≥ 1 的安全整数,否则ConfigurationError(requireInteger) |
resourceLimits? | 应用于每个线程的 V8 堆与栈限制(ResourceLimits类型) | 仅接受codeRangeSizeMb、maxOldGenerationSizeMb、maxYoungGenerationSizeMb、stackSizeMb四个键,且均为正有限数(index.ts) |
配置示例:
const limited = new WorkerThreads({ maxThreads: 2, resourceLimits: { maxOldGenerationSizeMb: 256 }, });线程超过堆限制会被 Node 终止,其尝试失败,后续尝试由替换线程接管。测试用 16 MiB 老生代限制验证了分配死循环的处理器会以内存限制失败、随后健康线程仍能正常工作(index.test.ts)。
7.2 诊断快照
WorkerThreadsDiagnostics提供线程池实时快照,River 会将其挂到客户端诊断的executors.worker_threads下:
| 字段 | 含义 |
|---|---|
activeThreads | 正在运行尝试的线程数 |
idleThreads | 等待尝试的存活线程数 |
pendingTasks | 等待线程的尝试数 |
totalThreads | 存活原生线程总数(含正在被终止的) |
crashedThreads | 自构造以来因未捕获异常、未处理 rejection、process.exit()或资源限制而丢失的线程数(abort 或 close 导致的不计) |
其中crashedThreads的持续上升通常意味着处理器在返回后遗留了失败的后台工作。
八、所有权与关闭:executor 属于应用
与 River 的运行时不同,executor 归应用所有(index.ts):
- 多个客户端可以共享同一个 executor;
- 停止任一客户端不会关闭 executor;
- 正确流程是:先停止所有使用它的客户端,再
await executor.close(),或让await using作用域结束时自动关闭([Symbol.asyncDispose]委托给close)。
close()会:使排队的尝试与正在运行的尝试以LifecycleError失败("worker thread executor is closed"/"worker thread executor closed while the attempt was running")、终止所有线程、等待线程退出;关闭是幂等的,之后新的start立即失败。测试 index.test.ts 与 client.test.ts 分别验证了共享 executor 在其中一个客户端停止后继续服务另一个客户端、以及await using自动关闭的行为。
另一个贴心设计:空闲线程不会让进程保持存活(worker.unref(),见 pool.ts),因此一个从未被关闭的 executor 也不会阻碍进程退出。
九、开发模式下的 TypeScript 源码加载
这是一个面向开发体验的特性。当宿主进程启用了 Node 的类型剥离能力(process.features.typescript)时,线程入口 thread.ts 会通过registerHooks注册解析钩子:当某个.js文件(处理器 URL 或其相对导入)解析失败、而同目录存在对应的.ts源码时,加载该源码。要点:
- 始终按编译后的
.js名称引用处理器模块,同一 URL 三处通用:构建产物(.js真实存在)、node src/main.ts源码运行(主线程相对导入需用.ts扩展名,例如配合 TypeScript 的rewriteRelativeImportExtensions)、以及 Vitest 测试; - 该回退仅在解析失败后生效,因此永远不会改变构建实际加载的模块;
- 通过
--import注册的 loader 也会作用于线程内(Node 会把进程的execArgv传给 worker thread); - 处理器模块必须使用 Node 可剥离的 TypeScript 语法,即排除 enum、参数属性(parameter properties)等需要运行时转换的特性。
测试数据中的 handlers.ts 故意只存在.ts而缺失.js,使每个测试都顺带覆盖了这条回退路径;而 format.ts 则验证了线程内相对导入经.js名称解析到.ts源码的能力。
十、安全边界与运行要求
- 可用性隔离,而非安全沙箱:worker thread 与主进程共享进程、环境变量与文件系统访问权限。只能运行受信任的处理器模块。API 报告的原文表述为 "Worker threads isolate availability, not security. Only run trusted handler modules."
- Node.js 版本:要求 Node.js 26 或更新,且内置原生
Temporal(node -p "typeof Temporal"应输出object)。官方二进制包含Temporal,但部分从源码编译的发行版(包括某些发行版与 Homebrew 包)没有。 - 版本对齐:
riverqueue是本包的 peer dependency,必须安装与本包完全相同的版本,npm 会在版本不匹配时直接拒绝,而不是加载两份副本。 - TypeScript 用户:需要 TypeScript 6.0 及以上与
@types/node,并在compilerOptions.types中列出"node"。
十一、何时该用、何时不该用
最后给出一个清晰的适用判断(综合 README 与 API 报告):
- 该用:CPU 密集、无法让出事件循环的处理器,且不希望它们阻塞 River 的认领/完成/取消/停机逻辑;
- 不该用:普通 I/O 处理器——promise 已经让它们高效共享事件循环,引入线程只会增加序列化、通信与内存开销;
- 记住三个边界:处理器必须是 ESM 导出(不能是闭包);参数必须能完整跨线程(River JSON 或受检的结构化克隆类型);线程上下文里没有
client/completeTx/resumable。
本文所有结论均有仓库证据可查:API 声明见 js/worker-threads/etc/worker-threads.api.md,完整使用指南见 js/worker-threads/README.md,运行时可执行示例见 js/examples/worker-thread-cpu/src/index.ts。理解这些机制后,你就能把最重的计算安全地挪出 River 的事件循环,同时仍然享受完整的任务生命周期管理。
- 任务调度
- 后端
【免费下载链接】river
The polyglot queue: Fast and reliable background jobs in Go, Ruby, Rust, and JS/TS on Postgres or SQLite.
相关推荐
River JS 官方示例 worker-thread-cpu 全解:用有界 Worker 线程池托管 CPU 密集型任务
River JS 官方示例 worker thread cpu 全解:用有界 Worker 线程池托管 CPU 密集型任务 River 是构建于 Postgre
任务调度后端River JS 的 worker-threads 执行器:把 CPU 密集型任务搬出事件循环的完整实践指南
River JS 的 worker threads 执行器:把 CPU 密集型任务搬出事件循环的完整实践指南 @riverqueue/worker thread
任务调度后端Midway 线程池组件 @midwayjs/piscina 实战指南:Worker 线程池中的 CPU 密集任务与 Midway 容器化执行
Midway 线程池组件 @midwayjs/piscina 实战指南:Worker 线程池中的 CPU 密集任务与 Midway 容器化执行 导读 @midw
后端微服务云原生
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考