Backstage Scheduler Service 详解:在插件级任务调度器中实现分布式定时任务的调度、协调与运维检查
Backstage 的 Scheduler Service(coreServices.scheduler)为后端插件提供了一套“插件作用域”的分布式定时任务调度能力:你可以在插件初始化时注册类似 cron 的周期任务,由框架负责跨实例的执行协调、超时控制、取消信号与状态持久化,并通过每个插件的 REST API 暴露任务清单与手动触发/取消接口。读完全文后,你将掌握 scheduleTask 的完整参数语义、global/local 两种作用域背后的数据库协调机制、三个调度器 REST 端点的用法,以及如何在测试中用 MockSchedulerService 控制任务执行。
一、Scheduler Service 的定位
当编写 Backstage 后端插件时,经常需要“定时做某事”——轮询外部系统、刷新缓存、周期性入库等。原生 Node.js 的 setInterval 在多实例部署下会重复执行,而 Backstage 的每个插件实例都可能运行在多个 worker 主机上。为此,官方在 服务文档 中提供了任务调度器:
- 作用域按插件隔离:调度器以 plugin ID 为单位创建,任务 ID 只需在插件内部唯一;
- 两种执行作用域:
global(全局任务,同一时刻只在一台 worker 上运行,通过数据库协调)与local(每台 worker 各自按周期运行,类似setInterval,但同一实例内不重叠); - 状态持久化:全局任务的状态写入数据库,重启、多实例协作、超时回收都依赖这份共享状态;
- REST API 内建于插件路由:调度器自动把
/.backstage/scheduler/v1/...路由挂载到插件的 HTTP Router 上,方便运维检查与手动干预。
服务接口定义在 SchedulerService.ts 中,核心方法为:
| 方法 | 说明 |
|---|---|
scheduleTask(task) |
一次性完成“调度定义 + 任务函数”的注册,是最常用的入口 |
triggerTask(id) |
手动触发某个任务立即执行;任务不存在抛 NotFoundError,正在运行抛 ConflictError |
cancelTask(id) |
取消正在运行的任务,将其标记回 idle;任务不存在或不在运行中分别抛 NotFoundError/ConflictError |
createScheduledTaskRunner(schedule) |
预先生成一个“休眠”的调度运行器,供外层代码控制调度、内层代码只提供任务实现(依赖注入式解耦) |
getScheduledTasks() |
返回当前实例已注册的全部任务描述符(ID、scope、序列化后的 settings) |
任务函数类型为 SchedulerServiceTaskFunction,可以接收一个 AbortSignal 参数(也可以不提供);信号触发时应尽快中止处理并返回。这是取消语义能生效的前提——调度器只能发出中止信号,能否及时退出取决于任务实现是否消费了该信号。
二、在插件中使用 Scheduler Service
下面是在 example 插件中注册一个每 10 分钟运行一次、跨实例协调执行的任务的完整示例(与官方文档一致,可直接复制):
import {
coreServices,
createBackendPlugin,
} from '@backstage/backend-plugin-api';
createBackendPlugin({
pluginId: 'example',
register(env) {
env.registerInit({
deps: {
scheduler: coreServices.scheduler,
},
async init({ scheduler }) {
await scheduler.scheduleTask({
frequency: { minutes: 10 },
timeout: { seconds: 30 },
id: 'ping-google',
fn: async () => {
await fetch('http://google.com/ping');
},
});
},
});
},
});
参数语义详解
scheduleTask 的参数是 SchedulerServiceTaskScheduleDefinition & SchedulerServiceTaskInvocationDefinition 的联合,字段定义与 JSDoc 均可在 SchedulerService.ts 中查到:
| 字段 | 必填 | 取值 | 说明 |
|---|---|---|---|
frequency |
是 | { cron: string } / HumanDuration / Duration / { trigger: 'manual' } |
任务执行频率。支持 crontab 风格字符串(可带可选的秒位:* * * * * *,从秒、分、时、日、月到星期)、人类可读时长对象(如 { minutes: 10 })、Luxon Duration,或 manual(仅手动触发时运行,适合需要全局互斥锁但不应并发运行的任务)。这是尽力而为(best effort)的频率:当一次执行耗时超过频率且未超时时,下一次执行会顺延到上一次结束之后 |
timeout |
是 | HumanDuration / Duration |
单次调用的最长允许时长。超过后任务被视为超时并被“释放”,允许新的调用发生(可能在另一台 worker 上) |
initialDelay |
否 | HumanDuration / Duration |
首次执行前的等待时间,适合冷启动场景下让服务先稳定再执行重批处理。注意:从源码结构看这是按 worker 生效的延迟——在多实例集群中,其他长生命周期 worker 仍可能在单个新 worker 处于初始延迟期间继续处理该任务,因此它不能用于“全局暂停”任务 |
scope |
否 | 'global'(默认) / 'local' |
并发控制/加锁的作用域。global:调度器尽量保证同一时刻只有一台 worker 机器运行该任务,worker 数量增加不会提高任务频率,负载被随机分摊到各主机,适合访问共享资源的任务(如 Catalog 入库,避免多机重复导入互相踩踏);local:没有跨主机协调,每台主机各自按周期运行,类似 setInterval,但运行时保证单机内不重叠 |
id |
是 | string | 插件内唯一的任务 ID |
fn |
是 | 任务函数 | 周期性调用的实际逻辑,可接收 AbortSignal |
signal |
否 | AbortSignal |
传入后,该信号触发会停止任务的重复执行(实现中会与根生命周期关闭信号做委托合并,见 PluginTaskSchedulerImpl.ts) |
配置驱动的调度定义
除代码传参外,同一套字段也可以从 app-config 中读取。SchedulerService.ts 导出的 readSchedulerServiceTaskScheduleDefinitionFromConfig(config) 接受“该定义的子配置”(不是根配置),支持如下写法:
backstage:
example:
schedule:
frequency:
cron: '0 0 * * *' # 或人类可读时长字符串,或 trigger: manual
timeout: 10m
initialDelay: 2m # 可选
scope: global # 可选,仅允许 global / local,其他值会抛错
其中 frequency 的解析规则(见 readFrequency 实现):对象且含 cron 键 → cron;对象且 trigger === 'manual' → 手动触发;其余情况按时长字符串解析(readDurationFromConfig)。scope 若不是 global/local 会直接抛出 Only "global" or "local" are allowed for TaskScheduleDefinition.scope 错误。
三、实现原理:global 任务如何跨实例“只跑一次”
服务装配链路
调度器在 schedulerServiceFactory.ts 中注册为 coreServices.scheduler 的服务工厂,依赖 database、logger、rootLifecycle、httpRouter、pluginMetadata 与 metrics 六个服务。DefaultSchedulerService.create()(见 DefaultSchedulerService.ts)做了三件事:
- 通过
database.getClient()获取 Knex 连接,并执行migrateBackendTasks建表(可通过database.migrations.skip跳过); - 启动一个
PluginTaskSchedulerJanitor清理器(非测试环境),每分钟运行一次,负责清理数据库中残留的失效运行票据(例如 worker 崩溃后遗留的current_run_ticket); - 创建
PluginTaskSchedulerImpl实例,并把它的 Express Router 挂载到插件的httpRouter上——这就是 REST API 的来源。
两种 worker
PluginTaskSchedulerImpl.scheduleTask()(见 PluginTaskSchedulerImpl.ts)先把调度参数序列化为 version 2 的 settings 对象:
const settings: TaskSettingsV2 = {
version: 2,
cadence: parseDuration(task.frequency), // ISO 时长 / cron 串 / "manual"
initialDelayDuration: task.initialDelay && parseDuration(task.initialDelay),
timeoutAfterDuration: parseDuration(task.timeout),
};
parseDuration 会把 HumanDuration 经 Luxon 转换为 ISO 时长字符串(如 PT10M),cron 与 manual 则原样保留——这正是 REST API 中 settings.cadence 三种取值的来源。随后按 scope 分发:
global→ 创建TaskWorker(跨主机协作加锁)与一个共享的TaskStatePoller;local→ 创建LocalTaskWorker(见 LocalTaskWorker.ts),纯内存状态,完全不访问数据库。
数据库表结构
全局任务的状态全部落在 backstage_backend_tasks__tasks 表中,定义见 tables.ts:
| 列 | 用途 |
|---|---|
id |
任务 ID(插件内唯一) |
settings_json |
version 2 设置对象的 JSON 序列化(有 zod schema taskSettingsV2Schema 校验,见 types.ts) |
next_run_start_at |
下次计划开始时间;manual 任务该值为 NULL |
current_run_ticket |
当前运行的 UUID 票据,非空即“有人正在跑” |
current_run_started_at / current_run_expires_at |
本次运行开始时间与超时时刻 |
last_run_ended_at / last_run_error_json |
上一次运行结束时间与错误(JSON 序列化) |
认领(claim)、心跳与超时回收
TaskWorker 的核心循环在 TaskWorker.ts 中:
- 登记:
persistTask先按cadence计算首次/下次运行时间(cron 用CronTime.sendAt(),时长型用now + cadence,manual 置 NULL),然后INSERT ... ON CONFLICT (id) DO UPDATE。若任务已存在则只替换 settings,且不覆盖更晚的next_run_start_at——这保证了滚动部署时新旧 worker 不会把计划时间往回拨; - 轮询:worker 默认每 5 秒(
DEFAULT_WORK_CHECK_FREQUENCY)经TaskStatePoller检查一次是否有到期工作;若 cadence 小于 5 秒,轮询频率会自动提升到 cadence 本身; - 认领:
tryClaimTask用一条WHERE current_run_ticket IS NULL的条件 UPDATE 原子写入新票据(见 tryClaimTask),返回受影响行数是否为 1 来判定认领成败——多台 worker 同时到期时,只有数据库行锁胜出的那台会真正执行,其余得到claim-lost; - 执行与保活:执行期间设置超时定时器(到
timeoutAfterDuration即 abort 任务),并按轮询间隔周期性做checkLiveness心跳——若数据库中票据已被清除(被别的 host 取消,或被 janitor 清理),立即中止本次执行; - 释放:正常结束或抛错后调用
tryReleaseTask,按票据匹配清除运行状态、写入last_run_*,并把next_run_start_at推进到greatest(next_run_start_at + interval, now())——取 max/ greatest 是为了避免宕机追赶时一次性补跑历史所有错过的周期; - 容错:worker 主循环捕获到意外异常时会打 warn 日志、睡 1 秒后重新进入循环,任务调度本身不会因单次失败而退出。
每次任务执行还被 instrumentedFunction(见 PluginTaskSchedulerImpl.ts)包了一层 OpenTelemetry span 与 metrics 打点:backend_tasks.task.runs.count(按 started/completed/failed 计数的 counter)、backend_tasks.task.runs.duration(秒为单位的直方图)、backend_tasks.task.runs.started / backend_tasks.task.runs.completed(gauge),标签含 taskId 与 scope,可直接接入你的可观测体系。
local 任务
LocalTaskWorker 不写数据库、不参与跨主机协调:每台 worker 各自按 cadence 运行任务,用内存中的状态机记录 idle/running,同样支持初始延迟、超时与手动 trigger()/cancel()。适合“本机缓存刷新”这类与实例生命周期绑定的工作。
四、REST API
调度器在每个插件的 base URL 下暴露一组 REST 端点,用于检查和影响该插件所有任务的当前状态(路由实现见 PluginTaskSchedulerImpl.getRouter)。
GET <pluginBaseURL>/.backstage/scheduler/v1/tasks
列出该插件在启动时注册的所有任务及其当前状态。例如查询 Catalog 插件的全部调度任务:
curl 'https://<instance-name>/api/catalog/.backstage/scheduler/v1/tasks'
响应结构如下:
{
"tasks": [
{
"taskId": "InternalOpenApiDocumentationProvider:refresh",
"pluginId": "catalog",
"scope": "global",
"settings": {
"version": 2,
"cadence": "PT10S",
"initialDelayDuration": "PT10S",
"timeoutAfterDuration": "PT1M"
},
"taskState": {
"status": "idle",
"startsAt": "2025-04-11T20:35:13.418+02:00",
"lastRunEndedAt": "2025-04-11T20:35:03.453+02:00"
},
"workerState": {
"status": "initial-wait"
}
}
]
}
每个任务包含以下属性:
| 字段 | 格式 | 说明 |
|---|---|---|
taskId |
string | 任务在插件内的唯一 ID |
pluginId |
string | 任务所归属的插件 |
scope |
string | local(在每个 worker 节点上运行,可能重叠,类似 setInterval)或 global(同一时刻只在一个 worker 节点上运行,不重叠) |
settings |
object | 调度时传入的初始设置的序列化形式。唯一完全固定的已知字段是 version,其余字段依赖所用版本 |
settings.version |
string | settings 对象格式的内部标识符,格式可能随版本完全变化;本文档描述的是 version 2 |
settings.cadence |
string;ISO 时长 | 任务运行频率。要么是字符串 manual(仅手动触发时运行),要么是以字母 P 开头的 ISO 时长串,要么是 cron 格式串 |
settings.initialDelayDuration |
string;ISO 时长 | 服务启动后 worker 在开始寻找工作前等待多久,给服务留出稳定时间(如已配置),为 ISO 时长串 |
settings.timeoutAfterDuration |
string;ISO 时长 | 任务开始后多久视为超时、可被重试接管 |
taskState |
object | 任务当前状态(见下) |
workerState |
object | 负责任务的 worker 状态(见下) |
taskState 的形状取决于任务是否正在运行。运行中时:
| 字段 | 格式 | 可选 | 说明 |
|---|---|---|---|
taskState.status |
string | running |
|
taskState.startedAt |
string;ISO 时间戳 | 本次运行开始的时间 | |
taskState.timesOutAt |
string;ISO 时间戳 | 若本次运行在该时刻前未结束则超时 | |
taskState.lastRunError |
string;JSON 序列化的错误 | 可选 | 上次运行若抛错,此字段包含该错误 |
taskState.lastRunEndedAt |
string;ISO 时间戳 | 可选 | 上次运行的结束时间 |
**空闲(idle)**时:
| 字段 | 格式 | 可选 | 说明 |
|---|---|---|---|
taskState.status |
string | idle |
|
taskState.startsAt |
string;ISO 时间戳 | 可选 | 任务下一次计划运行的时间;手动调度的任务不会有该字段 |
taskState.lastRunError |
string;JSON 序列化的错误 | 可选 | 上次运行若抛错,此字段包含该错误 |
taskState.lastRunEndedAt |
string;ISO 时间戳 | 可选 | 上次运行的结束时间 |
workerState 的形状如下:
| 字段 | 说明 |
|---|---|
workerState.status |
负责任务的 worker 的状态:initial-wait(服务刚启动时)、running(任务正在运行)或 idle(任务当前未运行) |
POST <pluginBaseURL>/.backstage/scheduler/v1/tasks/<taskId>/trigger
将指定任务 ID 的任务调度为立即执行,而无需等待下一个计划时间槽。例如手动触发 Catalog 的某个任务:
curl -X POST "https://<instance-name>/api/catalog/.backstage/scheduler/v1/tasks/InternalOpenApiDocumentationProvider:refresh/trigger"
注意:worker 发现任务到期并真正接手之前可能还有短暂的额外延迟,通常不到 1 秒,但会有波动。请求没有请求体。响应:
200 OK成功404 Not Found该插件下没有这个已注册任务409 Conflict任务已处于运行状态
从源码看(TaskWorker.trigger,见 TaskWorker.ts),实现上先确认任务存在,再用一条 WHERE current_run_ticket IS NULL 的条件 UPDATE 把 next_run_start_at 置为当前时间——即“可认领”即成功,否则返回冲突。
POST <pluginBaseURL>/.backstage/scheduler/v1/tasks/<taskId>/cancel
取消指定任务 ID 正在运行的任务。注意 <taskId> 必须做 URL 编码以保持在 URL 中为单个路径段(例如 JavaScript 中用 encodeURIComponent,或标准百分号编码)。例如取消 Catalog 的某个任务(: 编码为 %3A):
curl -X POST "https://<instance-name>/api/catalog/.backstage/scheduler/v1/tasks/InternalOpenApiDocumentationProvider%3Arefresh/cancel"
注意:worker 发现任务被取消可能还有几秒以内的延迟;同时,任务能否真正停下来取决于任务实现是否正确响应了传入的 abort 信号。请求没有请求体。响应:
200 OK成功404 Not Found该插件下没有这个已注册任务409 Conflict任务当前不处于运行状态
源码中 TaskWorker.cancel 的做法是:校验任务存在且确有票据在跑,然后清票据、推进 next_run_start_at,并把 last_run_error_json 写为 Task was cancelled;由于执行侧有前述的票据心跳检查,远端 worker 上的任务随后会被中止。
五、测试:使用 MockSchedulerService
@backstage/backend-test-utils 包提供 mockServices.scheduler,它是调度器服务的 mock 实现,可用于单元测试。在 startTestBackend 中它默认被使用:只要注册的任务不是 manual 调度、也没有配置 initial delay,就会在启动时立即执行。测试中可以用独立实例获得更多控制(示例与官方文档一致):
it('should trigger a task', async () => {
const scheduler = mockServices.scheduler();
const { server } = await startTestBackend({
features: [scheduler.factory()],
});
await scheduler.triggerTask('some-task-id');
// Next verify that the plugin state is updated accordingly
// e.g. by calling the API or verifying database state
});
MockSchedulerService 的行为可用 MockSchedulerService.test.ts 中的用例对照验证:直接对 mock 实例 scheduleTask 后 triggerTask 即可驱动任务执行;而通过 startTestBackend 注册的插件任务也会因默认 mock 调度器在启动时被立即运行。典型用法是:注入 mock 工厂 → 调用 triggerTask 或等待自动执行 → 断言 API 响应或数据库状态。
六、关键要点回顾
- 调度器是插件级核心服务:每个插件拿到自己的 scheduler 实例,任务 ID 在插件内唯一;依赖 database/httpRouter 等服务自动装配(见 schedulerServiceFactory.ts);
- 选
global(默认)保证跨实例同一时刻只跑一份,代价是约 5 秒级的轮询发现延迟;选local则每台 worker 各跑一份,适合本机状态维护; timeout决定超时后任务被释放换 worker 接管的时机,initialDelay只是单 worker 视角的冷启动缓冲,不能当作全局暂停开关;- REST 三端点(
GET .../tasks、POST .../trigger、POST .../cancel)覆盖了“查看状态、手动触发、紧急取消”的日常运维闭环,注意 trigger/cancel 的404/409语义以及 cancel 的生效延迟; - 状态真相在
backstage_backend_tasks__tasks表:票据(current_run_ticket)是跨主机互斥的关键,janitor 每分钟清理失效票据,worker 执行期的心跳检查保证取消与清理能被远端执行及时感知。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00