首页
/ Backstage Scheduler Service 详解:在插件级任务调度器中实现分布式定时任务的调度、协调与运维检查

Backstage Scheduler Service 详解:在插件级任务调度器中实现分布式定时任务的调度、协调与运维检查

2026-09-09 17:48:26作者:袁立春Spencer

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 的服务工厂,依赖 databaseloggerrootLifecyclehttpRouterpluginMetadatametrics 六个服务。DefaultSchedulerService.create()(见 DefaultSchedulerService.ts)做了三件事:

  1. 通过 database.getClient() 获取 Knex 连接,并执行 migrateBackendTasks 建表(可通过 database.migrations.skip 跳过);
  2. 启动一个 PluginTaskSchedulerJanitor 清理器(非测试环境),每分钟运行一次,负责清理数据库中残留的失效运行票据(例如 worker 崩溃后遗留的 current_run_ticket);
  3. 创建 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 中:

  1. 登记persistTask 先按 cadence 计算首次/下次运行时间(cron 用 CronTime.sendAt(),时长型用 now + cadence,manual 置 NULL),然后 INSERT ... ON CONFLICT (id) DO UPDATE。若任务已存在则只替换 settings,且不覆盖更晚的 next_run_start_at——这保证了滚动部署时新旧 worker 不会把计划时间往回拨;
  2. 轮询:worker 默认每 5 秒DEFAULT_WORK_CHECK_FREQUENCY)经 TaskStatePoller 检查一次是否有到期工作;若 cadence 小于 5 秒,轮询频率会自动提升到 cadence 本身;
  3. 认领tryClaimTask 用一条 WHERE current_run_ticket IS NULL 的条件 UPDATE 原子写入新票据(见 tryClaimTask),返回受影响行数是否为 1 来判定认领成败——多台 worker 同时到期时,只有数据库行锁胜出的那台会真正执行,其余得到 claim-lost
  4. 执行与保活:执行期间设置超时定时器(到 timeoutAfterDuration 即 abort 任务),并按轮询间隔周期性做 checkLiveness 心跳——若数据库中票据已被清除(被别的 host 取消,或被 janitor 清理),立即中止本次执行;
  5. 释放:正常结束或抛错后调用 tryReleaseTask,按票据匹配清除运行状态、写入 last_run_*,并把 next_run_start_at 推进到 greatest(next_run_start_at + interval, now())——取 max/ greatest 是为了避免宕机追赶时一次性补跑历史所有错过的周期;
  6. 容错: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),标签含 taskIdscope,可直接接入你的可观测体系。

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 实例 scheduleTasktriggerTask 即可驱动任务执行;而通过 startTestBackend 注册的插件任务也会因默认 mock 调度器在启动时被立即运行。典型用法是:注入 mock 工厂 → 调用 triggerTask 或等待自动执行 → 断言 API 响应或数据库状态。

六、关键要点回顾

  • 调度器是插件级核心服务:每个插件拿到自己的 scheduler 实例,任务 ID 在插件内唯一;依赖 database/httpRouter 等服务自动装配(见 schedulerServiceFactory.ts);
  • global(默认)保证跨实例同一时刻只跑一份,代价是约 5 秒级的轮询发现延迟;选 local 则每台 worker 各跑一份,适合本机状态维护;
  • timeout 决定超时后任务被释放换 worker 接管的时机,initialDelay 只是单 worker 视角的冷启动缓冲,不能当作全局暂停开关;
  • REST 三端点(GET .../tasksPOST .../triggerPOST .../cancel)覆盖了“查看状态、手动触发、紧急取消”的日常运维闭环,注意 trigger/cancel 的 404/409 语义以及 cancel 的生效延迟;
  • 状态真相在 backstage_backend_tasks__tasks 表:票据(current_run_ticket)是跨主机互斥的关键,janitor 每分钟清理失效票据,worker 执行期的心跳检查保证取消与清理能被远端执行及时感知。
登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395