ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

Backstage Scheduler Service 详解:在插件级任务调度器中实现分布式定时任务的调度、协调与运维检查 Backstage Scheduler Service 详解在插件级任务调度器中实现分布式定时任务的调度、协调与运维检查【免费下载链接】backstageBackstage is an open framework for building developer portals项目地址: https://gitcode.com/GitHub_Trending/ba/backstageBackstage 的 Scheduler ServicecoreServices.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正在运行抛ConflictErrorcancelTask(id)取消正在运行的任务将其标记回 idle任务不存在或不在运行中分别抛NotFoundError/ConflictErrorcreateScheduledTaskRunner(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 }、LuxonDuration或manual仅手动触发时运行适合需要全局互斥锁但不应并发运行的任务。这是尽力而为best effort的频率当一次执行耗时超过频率且未超时时下一次执行会顺延到上一次结束之后timeout是HumanDuration/Duration单次调用的最长允许时长。超过后任务被视为超时并被“释放”允许新的调用发生可能在另一台 worker 上initialDelay否HumanDuration/Duration首次执行前的等待时间适合冷启动场景下让服务先稳定再执行重批处理。注意从源码结构看这是按 worker 生效的延迟——在多实例集群中其他长生命周期 worker 仍可能在单个新 worker 处于初始延迟期间继续处理该任务因此它不能用于“全局暂停”任务scope否global默认 /local并发控制/加锁的作用域。global调度器尽量保证同一时刻只有一台 worker 机器运行该任务worker 数量增加不会提高任务频率负载被随机分摊到各主机适合访问共享资源的任务如 Catalog 入库避免多机重复导入互相踩踏local没有跨主机协调每台主机各自按周期运行类似setInterval但运行时保证单机内不重叠id是string插件内唯一的任务 IDfn是任务函数周期性调用的实际逻辑可接收AbortSignalsignal否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 的来源。两种 workerPluginTaskSchedulerImpl.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 时长字符串如PT10Mcron 与manual则原样保留——这正是 REST API 中settings.cadence三种取值的来源。随后按 scope 分发global→ 创建TaskWorker跨主机协作加锁与一个共享的TaskStatePollerlocal→ 创建LocalTaskWorker见 LocalTaskWorker.ts纯内存状态完全不访问数据库。数据库表结构全局任务的状态全部落在backstage_backend_tasks__tasks表中定义见 tables.ts| 列 | 用途 | | -- | ---- | |id| 任务 ID插件内唯一 | |settings_json| version 2 设置对象的 JSON 序列化有 zod schemataskSettingsV2Schema校验见 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 cadencemanual 置 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.completedgauge标签含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.41802:00, lastRunEndedAt: 2025-04-11T20:35:03.45302:00 }, workerState: { status: initial-wait } } ] }每个任务包含以下属性字段格式说明taskIdstring任务在插件内的唯一 IDpluginIdstring任务所归属的插件scopestringlocal在每个 worker 节点上运行可能重叠类似setInterval或global同一时刻只在一个 worker 节点上运行不重叠settingsobject调度时传入的初始设置的序列化形式。唯一完全固定的已知字段是version其余字段依赖所用版本settings.versionstringsettings 对象格式的内部标识符格式可能随版本完全变化本文档描述的是 version 2settings.cadencestringISO 时长任务运行频率。要么是字符串manual仅手动触发时运行要么是以字母 P 开头的 ISO 时长串要么是cron格式串settings.initialDelayDurationstringISO 时长服务启动后 worker 在开始寻找工作前等待多久给服务留出稳定时间如已配置为 ISO 时长串settings.timeoutAfterDurationstringISO 时长任务开始后多久视为超时、可被重试接管taskStateobject任务当前状态见下workerStateobject负责任务的 worker 状态见下taskState的形状取决于任务是否正在运行。运行中时字段格式可选说明taskState.statusstringrunningtaskState.startedAtstringISO 时间戳本次运行开始的时间taskState.timesOutAtstringISO 时间戳若本次运行在该时刻前未结束则超时taskState.lastRunErrorstringJSON 序列化的错误可选上次运行若抛错此字段包含该错误taskState.lastRunEndedAtstringISO 时间戳可选上次运行的结束时间**空闲idle**时字段格式可选说明taskState.statusstringidletaskState.startsAtstringISO 时间戳可选任务下一次计划运行的时间手动调度的任务不会有该字段taskState.lastRunErrorstringJSON 序列化的错误可选上次运行若抛错此字段包含该错误taskState.lastRunEndedAtstringISO 时间戳可选上次运行的结束时间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 的某个任务:编码为%3Acurl -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 上的任务随后会被中止。五、测试使用 MockSchedulerServicebackstage/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 执行期的心跳检查保证取消与清理能被远端执行及时感知。【免费下载链接】backstageBackstage is an open framework for building developer portals项目地址: https://gitcode.com/GitHub_Trending/ba/backstage创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表