- 任务调度
- 后端
【免费下载链接】river
The polyglot queue: Fast and reliable background jobs in Go, Ruby, Rust, and JS/TS on Postgres or SQLite.
本文是 River JavaScript 运行时(@riverqueue核心运行时)的技术指南,围绕 js/docs/runtime.md 展开。River 是一个在 Postgres 或 SQLite 上运行的 polyglot 后台任务队列,JavaScript 客户端提供异步、受监督(supervised)的运行时:client.start()返回RunHandle,应用持有该句柄直到停机。读完本文,你将掌握 River.js 的完整生命周期管理、maxWorkers与事件循环的正确用法、重试/错误/取消语义、数据库故障下的韧性行为,以及"永不领导"客户端与fetchOnlyKnownKinds多语言共享队列的部署模式。
运行时生命周期:异步且受监督
River.js 的运行时是异步的、受监督的。client.start()返回一个RunHandle,该句柄的所有权属于应用,直到停机。核心约定是:必须持续观察run.completed,因为一旦运行时故障,这个 Promise 会 reject——如果应用忽略它,故障将不可见(Node.js 中未处理的 rejection 也可能导致进程异常)。
启动与优雅停机的最小示例
await using run = await client.start(); process.once("SIGTERM", () => { void run.stop({ mode: "graceful", timeout: { seconds: 30 } }); }); await run.completed;要点解读:
await using run借助RunHandle[Symbol.asyncDispose]()(实现于 runtime.ts)在作用域退出时自动执行stop({ mode: "graceful" })。- 优雅停机(graceful):停止领取新任务,让已激活的 handler 全部完成;若给了
timeout,超过该时限后升级为取消。 - River 从不自行安装信号处理器。
SIGTERM等信号必须像上面这样由应用自己接线;这保持了库的进程语义中立。 cancel模式还会中止所有活跃 handler 的AbortSignal。
停机与故障的语义
- 与 River for Go 一致:停机一开始,正在领导的客户端立即辞去领导并结束维护工作,同时其队列继续排空;在此期间,另一个客户端可以接管维护。
- 多次调用 stop 共享同一个停机流程(
RuntimeController.stop中#shutdownPromise只创建一次,见 runtime.ts)。 - 当运行时内部故障时,它以同样的方式停机:
#fail触发完整的#finishStop,run.completed在清理完成后以故障原因 reject(见 runtime.ts)。
从源码看,#finishStop的停机顺序非常明确(runtime.ts):先等待ready,中止维护与领导服务 → 排空所有队列生产者 → 中止取消轮询 → 等待全部后台任务 → 中止运行期服务 → 关闭完成管道(completion pipeline)并等待后期任务 → 排空事件投递 → 停止事件循环延迟监控 → 释放 liveness 定时器 → 才 resolvecompleted。因此run.completed的 settle 一定发生在所有清理完成之后。
协作式取消与尝试身份守卫
作业超时(job timeout)、远程取消或停机中止对进程内 handler 都是协作式的:它们通过AbortSignal请求 handler 自行收尾。迟到的结果被**尝试身份(attempt identity)**保护:一次完成操作只有在作业仍属于产生它的那次尝试(相同的attempt编号、由同一 client 认领)时才会改写作业(参见 completion-pipeline.ts 的committedBy校验逻辑),因此迟到结果不可能覆盖更新的一次尝试。
并发模型与事件循环
maxWorkers 是 I/O 并发,不是线程池
maxWorkers限制每个队列的活跃 handler 数量。从源码看,每个配置队列对应一个独立的 claim 循环(QueueProducer的#queueLoop,queue-producer.ts),一次 claim 最多填满当前空闲的 worker 槽位。
- 异步数据库/网络 handler 天然适合进程内:Node 可以高效地同时运行大量此类 handler。
- 同步 CPU 密集型工作会阻塞该进程内的 job 领取(claim)、完成(completion)、取消以及应用代码——即使
maxWorkers设得很高也无济于事,因为它们共享同一个事件循环。 maxWorkers的合法范围是 1~10,000(normalizeQueueConfig,settings.ts)。
领取节奏:fetchCooldown、pollInterval 与抖动
每个队列还拥有两个领取节奏参数:
| 参数 | 含义 | 默认值 |
|---|---|---|
fetchCooldown(客户端级fetchCooldown或队列级fetchCooldownMs) | 两次 claim 查询之间的最小间隔 | 未设置时为客户端fetchCooldown,即 100 毫秒 |
pollInterval(队列级pollIntervalMs) | 没有插入通知到达时的轮询兜底间隔 | 1 秒 |
源码中的默认值定义在 settings.ts:DEFAULT_FETCH_COOLDOWN_MS = 100,pollIntervalMs默认1_000,且校验要求pollInterval不得短于fetchCooldown。
与 River 其他运行时一致,每次轮询会额外等待一个随机增量:最多pollInterval的十分之一、且至少 10 毫秒(jitteredPollInterval,queue-producer.ts),这样多个生产者不会步调一致地同时发起查询。
通知与冷却的关系:插入通知(insert notification)会唤醒轮询等待,但绝不会绕过冷却期。因此突发流量会被聚合成一批,而不是为每个插入的作业都发起一次数据库查询。实践建议:高吞吐队列可以有意识地调低fetchCooldown,但保持pollInterval至少与fetchCooldown一样大。
插入通知限流
与 River for Go 一致,一个 client 在每个队列、每个fetchCooldown周期内最多发送一条插入通知,而且适用于所有 driver——无论作业是直接插入、由 periodic jobs 插入,还是到点调度产生。理由来自客户端选项文档(options.ts):生产者本来每个冷却周期最多领取一次,被抑制通知的作业会通过轮询拿到;重试(retry)作业不发送插入通知。
事件循环延迟监控与 worker-threads
River 会测量事件循环延迟并暴露在运行时诊断(run.diagnostics.eventLoopDelay)中。监控基于node:perf_hooks的monitorEventLoopDelay(event-loop-delay-monitor.ts),默认每 1 秒采样一次(reportIntervalMs),直方图分辨率为 20ms,超过 100ms(warningThresholdMs)记为exceededThreshold并触发runtime_event_loop_delay事件。
把持续的事件循环延迟当作运维故障处理:正确的做法是把 CPU 密集工作迁移到@riverqueue/worker-threads、独立进程或专门的服务,而不是调高队列并发。
Worker 线程有两个关键性质:
- 有界:线程数量受控,不会随队列膨胀无限增长。
- 可强停:如果某个 handler 在超过
jobStuckThreshold(客户端选项,默认 10 秒,见 settings.ts)后仍无视 abort 信号,worker-threads 执行器会强制终止它。
Worker 线程通过structured clone边界传输的是持久化的 JSON 参数,而不是解码函数或任意转换后的类实例。如果 handler 需要更丰富的本地表示,应在 worker 模块内部自行解码。
重试与错误语义
handler 的四种显式结果
- handler resolve → 作业完成(completed)。
- handler 抛错或 reject → 记录一次尝试错误(attempt error),并应用配置的重试策略。
return snooze(...)→ 本次尝试不计入尝试次数,按给定时长重新调度。return discard(...)→不再重试,作业被丢弃。return cancel({ reason })→ 等价于 Go 端的river.JobCancel,永久取消作业,并把JobCancelError: <reason>记录为尝试错误。
上述类型定义在 worker.ts:CancelOutcome、CompleteOutcome、DiscardOutcome、SnoozeOutcome共同构成WorkOutcome。
安全提醒:River 会限制记录的错误与堆栈大小(源码中canonicalError将错误消息与堆栈截断到 32,768 字符,failures.ts)。因此不要把凭证或敏感载荷放进抛出的错误消息里。
错误时间戳与"panic"类比
每条记录的尝试错误都带有所属尝试的开始时间(at)。
Go 端只为 panic 记录堆栈;River 的 JavaScript 端把panic 类比为运行时故障,即以原生TypeError、RangeError、ReferenceError、SyntaxError、EvalError、URIError形式浮出的错误。只有这些错误记录堆栈;而 handler 主动抛出的错误——包括这些类的子类——只记录消息,堆栈为空(isRuntimeFault通过原型链精确判定,failures.ts),这也让数据行保持紧凑。
近未来重试/暂停的快速通道
与其他运行时一致:距离当前不足一个调度器周期(scheduler interval,默认 5 秒)的重试或 snooze,会直接以available状态存储并携带未来的scheduled_at,从而准点执行,而不是等领导者的下一次调度器扫描。源码中DEFAULT_SCHEDULER_INTERVAL_MS = 5_000(settings.ts),完成管道在生成完成命令时把这个区间考虑在内(CompletionPipeline的schedulerIntervalMs参数)。
AbortSignal 与三种中止场景
每个 handler 收到一个AbortSignal。请把它传入数据库、HTTP 等可取消感知的操作中;如需 catch abort,只用于清理资源。中止后的记录方式与其他运行时完全一致:
远程取消之后:
- handler 仍正常 resolve → 作业完成;
- 抛错(例如通过
signal.throwIfAborted())、或返回snooze(...)/discard(...)→ 作业被取消。
作业超时之后:
- 成功 → 作业完成;
- handler 因 abort 停止 → 记录超时作为其错误;
- 其他任何错误 → 按抛出记录。
停机期间:
- 只有因停机 abort 而停止的 handler 会被不消耗尝试次数地重新置为
available; - 成功 → 作业完成;
- 真实错误 → 正常记录并按策略重试。
这些规则在AttemptRunner.#persistAttempt中精确实现(attempt-runner.ts):远程取消优先于除"干净成功完成"之外的一切结果;超时后成功照常完成;停机中止下只有被该中止停止的 handler 才走 interrupt 路径。
job_cancelled与job_interrupted事件在数据库状态迁移提交之后恰好发出一次(完成管道在持久化成功后投递事件,completion-pipeline.ts)。来自已不再拥有自己数据行的尝试的迟到结果会被直接忽略。
数据库故障:可运维,而非致命
后台工作中的数据库错误属于运维性故障,不是致命的。被其他会话持有的锁、statement_timeout、故障转移(failover)、连接池饱和——都会通过客户端的logger记录并重试,绝不会让运行时停下。run.completed只在配置错误和内部不变量被破坏时 reject。
指数退避重试
- claim、队列控制轮询、通知流在失败后按指数退避 + 抖动重试:从 250 毫秒起、封顶 30 秒,成功后重置。源码常量
BACKGROUND_BACKOFF = { baseMs: 250, maxMs: 30_000 },且每次重试叠加 ±10% 的抖动(backoff.ts),让同一事故中恢复的多个进程错峰散开。 - 通知流断线重订阅:一旦重订阅成功,会立刻轮询所有队列及其控制信息(pause/resume 等),确保监听器离线期间插入的作业不会干等到下一个轮询周期;同时读取该 client 正在运行的作业,取消那些在断线期间被取消的(
NotificationPump.#recoverMissedNotifications,notification-pump.ts)。 - 与 River for Go 一致:通知流在启动时无法连接并监听,会在领取任何作业之前让
client.start()以该错误 reject(见 runtime.ts 的初始化顺序:先启动通知泵,再启动队列)。 - 领导者选举在失败后会比正常间隔更早地重试;维护服务把失败记录为
maintenance_failed事件,并在下一个间隔重试。 - 日志级别:瞬时故障记为警告(warning),其他故障记为错误(error)。判别逻辑在 context.ts:永久性错误(配置、能力缺失、后端不匹配)直接重抛,因为它重试也修不好。
完成持久化的批处理与重排队
完成操作(completion)以批次持久化:
- 每次尝试(attempt)上限 10 秒(
COMPLETION_TIMEOUT_MS); - 一个批次最多 3 次尝试,退避间隔为 1、2、4 秒(
COMPLETION_ATTEMPTS = 3、COMPLETION_BACKOFF = { baseMs: 1_000, maxMs: 4_000 },completion-pipeline.ts)。
之后的分流:
- 瞬时失败(
DatabaseOperationError且retryable为 true)→ 批次被重新入队(requeue),worker 等待完成容量,直到数据库恢复; - 任何其他失败→ 批次被丢弃(drop),其中的作业保持
running状态,由 rescuer 在maintenance.rescueAfter之后重试; - 两种结局都会被记录,并分别上报为
job_completion_requeued与job_completion_dropped指标。
停机时:第一个持久性失败会放弃剩余的完成操作、交给 rescuer,因此即使数据库不可用,stop()也能结束(CompletionPipeline.drain()配合#persistCompletions中"lastAttempt && draining 时不延迟"的逻辑)。
完成操作的归属与竞态
一次完成操作只有在作业仍属于产生它的尝试时才会改写作业:尝试编号相同、且由同一 client 认领。若该尝试的作业已离开running状态(例如 rescuer 已重试或它被取消),该尝试记录的 output 与 metadata 仍会合并进作业(与 River for Go 一致),但其状态保持不变。
与 Go 不同:一个迟到的完成操作永远不会改写另一个 client 后来重新认领的作业——这种情况会作为job_race事件上报。
进程扩展:跨语言协作
单个进程可以高效运行 I/O 密集型工作。当需要应对 CPU/控制开销、可用性或部署隔离时,增加普通进程副本即可。
关键在于:River 的数据库协议在 JavaScript、Go 和 Rust 进程之间协调认领与领导权——库中没有隐藏任何 JavaScript 专属的进程管理器。这意味着你可以放心地在多个 Node.js 进程之间水平扩展,也可以与 Go/Rust 客户端混布在同一个数据库和 schema 上共享队列。
永不领导的客户端(leaderElectionDisabled)
leaderElectionDisabled: true让一个 client 不参与领导者选举,等价于 River for Go 的Config.LeaderElectionDisabled。它会照常从配置的队列领取并执行作业,但永远不会运行调度器(scheduler)、rescuer、cleaner、reindexer 或 periodic jobs。
适用场景:专门处理特定作业类型、不应承担任何其他职责的进程。
const videoClient = new Client(new PgDriver(new Pool()), { leaderElectionDisabled: true, queues: { video: { maxWorkers: 4 } }, workers: videoWorkers, });三个必须遵守的约束(均有源码校验支撑):
- 必须存在至少一个可领导的 client:同一数据库和 schema 上(任意语言皆可)必须有另一个已启动且仍可参与选举的 client。否则定时作业、重试、periodic jobs、卡住作业救援和清理都会停止推进——禁用选举的 client 即使没有任何其他 client 在运行,也永远不会领导。
- 不能配置
periodicJobs,也不能修改client.periodicJobs:两者都会抛出ConfigurationError(settings.ts 中"periodicJobs must be empty when leaderElectionDisabled is true")。但它仍会执行领导者插入到其队列中的 periodic 作业。 - 它的
maintenance设置不起任何作用。
与只认识部分作业类型的 client 共享队列
默认情况下,client 会领取其队列中的每一个作业;对于没有对应 worker 的 kind,作业会以"未知作业类型"错误失败。
fetchOnlyKnownKinds: true会把认领范围限制在启动时该 client 已注册 worker 的 kinds(等价于 River for Go 的Config.FetchOnlyKnownKinds)。其他 kind 的作业保持可用且不消耗尝试次数,因此拥有不同 worker 的 client 可以共享一个队列——典型场景是作业类型在语言之间迁移时:
const client = new Client(new PgDriver(new Pool()), { fetchOnlyKnownKinds: true, leaderElectionDisabled: true, queues: { default: { maxWorkers: 10 } }, workers: migratedWorkers, });底层实现:resolveRuntimeSettings在启动时把fetchKinds固定为workers.kinds()的排序快照(settings.ts),QueueProducer认领时把它作为kinds过滤参数传入。
两个重要边界:
- 该选项只影响认领。领导者的 rescuer 仍会处理每个队列中的卡住作业,并丢弃它不认识的 kind。因此一个只认识部分 kind 的 client 应当同时设置
leaderElectionDisabled,并确保另有(任意语言的)具备全部 kind worker 的可领导 client 在运行。 - 若不加此选项,未知 kind 的作业会被认领并以未知作业类型错误失败(源码注释明确说明这一取舍:未知 kind 会消耗一次尝试并持久化兼容的执行错误,而不是无限搁置,见 queue-producer.ts)。
常用配置速查
以下是运行时相关客户端选项(全部接受Temporal.Duration形式,例如{ seconds: 5 };不要使用带Ms后缀的选项名,会直接触发校验错误,见 options.ts):
| 选项 | 作用 | 默认值 |
|---|---|---|
queues.<name>.maxWorkers | 每队列并发 handler 上限 | 必填,1~10000 |
queues.<name>.fetchCooldown/ 客户端fetchCooldown | claim 最小间隔 | 100ms |
queues.<name>.pollInterval | 无通知时的轮询兜底 | 1s |
jobTimeout | 默认协作式作业超时 | 1 分钟;null禁用 |
jobStuckThreshold | 超时后判定卡住、执行器强停的等待 | 10s |
leaderElectionDisabled | 不参与领导选举 | false |
fetchOnlyKnownKinds | 只认领已注册 kind | false |
queueControlPollInterval | 队列控制(pause/resume)轮询、无通知时的取消轮询 | 2s |
queueHeartbeatInterval | 已配置队列的心跳上报 | 30s |
eventLoopDelay | 事件循环延迟监控,false关闭 | 开启:report 1s / resolution 20ms / 阈值 100ms |
maintenance | 领导端维护服务(保留期、选举间隔、各服务间隔、rescueAfter等) | 见 options.ts |
以上运行时骨架——RuntimeController统筹QueueProducer(认领)、AttemptRunner(执行)、CompletionPipeline(持久化结果)、NotificationPump(后端通知)四个协作者——均可在 js/src/runtime.ts 与 js/src/runtime/ 目录下对照阅读。结合 js/docs/runtime.md 原文与 js/src/client.ts 的公开 API,即可完整掌握从启动到停机的每一步行为。
- 任务调度
- 后端
【免费下载链接】river
The polyglot queue: Fast and reliable background jobs in Go, Ruby, Rust, and JS/TS on Postgres or SQLite.
相关推荐
BEX Engine 深度解析:BAML 语言异步运行时的事件循环、并发模型与基于 Epoch 的垃圾回收协调
BEX Engine 深度解析:BAML 语言异步运行时的事件循环、并发模型与基于 Epoch 的垃圾回收协调 导读 本文以开源仓库 BAML(The prog
编程语言AI Agent编译器CLI人工智能BFKit本地化方案:轻松支持11种语言的iOS应用国际化指南
BFKit本地化方案:轻松支持11种语言的iOS应用国际化指南 想要让你的iOS应用走向全球市场?BFKit为你提供了一套完整的本地化解决方案!BFKit是一个
移动开发开发工具FoundationDB Flow 运行时深度解析:Actor 协程、Future/Promise 异步原语与 Net2 事件循环
FoundationDB Flow 运行时深度解析:Actor 协程、Future/Promise 异步原语与 Net2 事件循环 导读 Flow 是 Foun
分布式数据库KV存储数据库后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考